- 可观测性
- 后端
【免费下载链接】highlight
highlight.io: The open source, full-stack monitoring platform. Error monitoring, session replay, logging, distributed tracing, and more.
本篇技术指南基于 highlight.io(开源全栈可观测性平台)仓库中的架构文档展开,梳理其代码库的高层组织结构:从前端 SDK 数据采集、Public Graph 数据摄入、Private Graph 前端查询,到 Worker 异步处理与 Kafka / InfluxDB / OpenTelemetry 集成。读完本文,你将掌握该仓库的目录定位方法、两大 GraphQL 端点的职责与本地调试端口,以及 Worker 后台任务的处理机制,可直接上手对源码进行探索与二次开发。
架构文档定位与阅读前提
仓库内的架构说明位于 docs-content/general/4_company/open-source/contributing/architecture.md,它以极简的方式给出了整个代码库的"顶层导览":SDK、Public Graph、Private Graph、Workers四大块。文档原文以列表与示意图为主,本文将在其基础上,逐一对每个组件深入源码,还原它们在实际仓库中的真实实现与调用关系。
- SDKs sdk/ - Firstload - Client - highlight-node / other SDKs - Public Graph backend/public-graph/graph/schema.resolvers.go (SDK 数据摄入 GraphQL 端点,本地地址 http://localhost:8082/public) - Private Graph backend/private-graph/graph/schema.resolvers.go (供前端使用的 GraphQL 端点,本地地址 http://localhost:8082/private) - Workers backend/worker.go - Public graph worker processPublicWorkerMessage - Async worker Start下文将分别从数据采集端(SDK)、数据摄入端(Public Graph)、查询服务端(Private Graph)与后台处理端(Workers)四个层次展开,最后结合 Kafka、InfluxDB、OpenTelemetry 三张集成示意图说明其周边基础设施。
一、数据采集层:sdk/ 目录与各类语言 SDK
架构文档指出,所有采集端代码位于根目录的sdk/之下。实际仓库中该目录包含三类核心形态(见 sdk/ 目录列表):
- Firstload(引导加载器):负责在浏览器页面早期注入并加载核心采集脚本,是会话录制(Session Replay)与错误采集的"起搏器",决定了前端页面首屏性能与采集启动时机。
- Client(浏览器客户端):实现事件采集、网络请求记录、控制台日志、WebSocket 事件等浏览器侧数据收集,并将数据批量上报到后端。
- highlight-node / 其他 SDK:覆盖服务端与框架侧,仓库中可以看到 highlight-go、highlight-java、highlight-py、highlight-ruby、highlight-rust、highlight-next、highlight-node、highlight-react、highlight-remix、highlight-dotnet、highlight-ex、highlight-hono、highlight-chrome 等十余种语言的实现,外加 highlight-wordpress 等平台适配。
SDK 采集到的原始数据(会话事件、日志、资源、WebSocket 消息、错误等)并不会直接写入存储,而是被打包后发送给后端的数据摄入端点——即下一层的 Public Graph。
二、数据摄入层:Public Graph 端点
架构文档给出的第二个关键入口是:
Public Graph
backend/public-graph/graph/schema.resolvers.go—— SDK 数据摄入 GraphQL 端点,本地调试地址为 http://localhost:8082/public
该文件确实位于 backend/public-graph/graph/schema.resolvers.go。从源码看,它是基于99designs/gqlgen(v0.17.70)生成的 GraphQL 解析器实现文件,定义了一系列 Mutation 解析器用于接收 SDK 上报的数据,例如:
InitializeSession:接收会话初始化参数(sessionSecureID、organizationVerboseID、隐私设置enableStrictPrivacy、网络录制开关enableRecordingNetworkContents、客户端与 Firstload 版本号、environment、appVersion、serviceName、fingerprint、clientID、networkRecordingDomains、disableSessionRecording、privacySetting),解析后写入会话模型。见 schema.resolvers.go#L29-L31。IdentifySession、AddTrackProperties等:承载用户识别与业务属性上报。
在数据摄入链路中,Public Graph 并不立即完成全部处理,而是将任务投递到 Kafka 队列,由后台 Worker 异步消费处理(详见下文 Workers 部分)。这种"SDK → Public Graph(摄入)→ Kafka → Worker(处理)"的异步解耦,是 highlight.io 支撑高吞吐采集的核心设计。
三、查询服务层:Private Graph 端点
架构文档中的第二个 GraphQL 入口:
Private Graph
backend/private-graph/graph/schema.resolvers.go—— 供前端使用的 GraphQL 端点,本地调试地址为 http://localhost:8082/private
对应的源码位于 backend/private-graph/graph/schema.resolvers.go。与 Public Graph 面向 SDK 写入数据不同,Private Graph 面向前端控制台(Frontend)提供查询与业务操作能力,例如会话列表、错误分组、日志检索、告警配置等。两个端点按职责拆分为独立的 GraphQL Schema(各自的gqlgen.yml分别定义在 backend/public-graph/graph/gqlgen.yml 与 backend/private-graph/graph/gqlgen.yml),并通过localhost:8082/public与localhost:8082/private两个本地路径对外暴露。
四、后台处理层:Workers
架构文档将 Worker 拆成两条主线:
- Public graph worker ——
processPublicWorkerMessage - Async worker ——
Start
需要说明的是,文档中提到的backend/worker.go在实际仓库中已演进为独立的backend/worker/包,两个核心函数均实现在 backend/worker/worker.go。
4.1 processPublicWorkerMessage:Kafka 消息的消费处理器
processPublicWorkerMessage位于 backend/worker/worker.go#L359。从源码可以看到它承担了两个职责:
其一,长任务监控。函数启动后即创建一个后台 goroutine 观察任务执行时长:每秒输出一次进度日志,超过10 分钟仍未完成则主动cancel()取消任务(见 backend/worker/worker.go#L370-L403),避免单个消息拖垮整个消费链路。
其二,按消息类型分发处理。函数体是一个switch task.Type,覆盖了 Public Graph 摄入后投递到 Kafka 的各类型消息,例如:
kafkaqueue.PushPayload:调用PublicResolver.ProcessPayload处理会话数据(事件、消息、资源、WebSocket 事件、错误、日志等);kafkaqueue.PushCompressedPayload:调用ProcessCompressedPayload处理压缩上报的负载;kafkaqueue.InitializeSession:调用InitializeSessionImpl完成会话初始化,并写入worker.session.initialize.count指标;kafkaqueue.IdentifySession:调用IdentifySessionImpl完成用户身份关联。
见 backend/worker/worker.go#L405-L458。由此可以清晰还原链路:SDK 上报 → Public Graph 摄入 → Kafka 队列 →processPublicWorkerMessage分发 →PublicResolver.*Impl落库。
此外,PublicWorker(ctx, topic)(backend/worker/worker.go#L518)是这一处理器的运行入口,负责从指定 Kafka topic 持续拉取消息并交给processPublicWorkerMessage执行。
4.2 Start:异步会话处理 Worker 主循环
Start位于 backend/worker/worker.go#L1179,它实现了一个带资源约束的轮询式异步任务池:
- 初始化最大 10 个并发的 workerpool(
maxWorkerCount := 10),并挂载 panic 恢复处理器; - 每 1 秒从数据库拉取一批待处理会话(数量上限约
processSessionLimit(200)加随机抖动),见 backend/worker/worker.go#L1183-L1189; - 提交任务前检查系统内存:若启用了
WorkerMaxMemoryThreshold配置且当前内存使用率超过阈值,则每 5 秒重试直到内存回落,见 backend/worker/worker.go#L1215-L1226; - 每个会话交由
processSession处理(backend/worker/worker.go#L653),失败时递增RetryCount,超过MAX_RETRIES(5 次)后将会话标记为Excluded并改投递SessionDataSync同步任务,见 backend/worker/worker.go#L1229-L1256; - 等待队列饱和时(
WaitingQueueSize() >= processSessionLimit)轮询休眠,实现背压控制。
Start之外,同一文件中还挂载了多个后台任务入口,包括StartLogAlertWatcher(日志告警监听,backend/worker/worker.go#L1288)、StartMetricAlertWatcher(指标告警监听,backend/worker/worker.go#L1292)与StartSessionDeleteJob(会话删除任务,backend/worker/worker.go#L1296)。这些任务在Worker结构体(持有Resolver、PublicResolver与StorageClient,见 backend/worker/worker.go#L75-L79)之上运行,与部署配置 deploy/worker-task.json、docker/compose.yml 中的 worker 服务对应。
五、通用架构图:全链路数据流
架构文档以一张全局图总结端到端链路:浏览器/服务端 SDK 采集数据 → Public Graph 摄入 → Kafka 缓冲 → Worker 处理(写入 PostgreSQL / ClickHouse / 对象存储)→ Private Graph 供前端查询。
代码结构图
第二张图按"组件职责"拆分代码库:SDK(浏览器与后端)、Public Graph、Private Graph、Workers,以及支撑它们的 ClickHouse、Kafka、Redis、存储等服务,与上文四个层次的划分一一对应。
六、Kafka:异步消息缓冲与解耦
highlight.io 将 Kafka 作为 SDK 数据摄入与后台处理之间的异步消息总线。上一节中processPublicWorkerMessage消费的kafkaqueue.Message即定义于 backend/kafka-queue(其中 types.go 定义消息类型,kafkaqueue.go 实现生产者/消费者,aws.go 提供 AWS MSK 兼容支持)。
消息类型大致包括PushPayload、PushCompressedPayload、InitializeSession、IdentifySession、AddTrackProperties、SessionDataSync等,覆盖"摄入 - 处理 - 同步"的全链路。文档中的 Kafka 示意图如下:
从源码结构看(结合 backend/kafka-queue/types.go 中的类型定义),Kafka 队列不仅用于事件摄入,还承担了数据同步(DataSyncQueue,见 backend/worker/worker.go#L105)等异步职责,使采集链路具备削峰与重试能力。
七、InfluxDB:指标与时序数据
架构文档中的 InfluxDB 图对应 highlight.io 的指标(Metrics)与告警评估场景。仓库中与 InfluxDB 直接相关的是后端指标告警任务 backend/jobs/metric-alerts(由StartMetricAlertWatcher驱动),以及 backend/clickhouse/metric_history.go 中的指标历史记录。
需要说明的是,随着仓库演进,时序与日志数据已逐步迁移至 ClickHouse(见 backend/clickhouse 与迁移脚本 backend/clickhouse/migrations),InfluxDB 相关图示反映的是架构文档编写时期的设计形态;当前源码中的指标存储与查询路径以 ClickHouse 为主,阅读旧文档时需注意这一演进差异。
八、OpenTelemetry:标准化遥测接入
最后一张示意图说明 highlight.io 对OpenTelemetry 标准的接入:SDK 与后端均可通过 OTLP 协议上报 traces、metrics、logs,纳入统一的采集与查询链路。仓库中的支撑证据包括:
- 采集侧样例与转换逻辑:backend/otel/otel.go、backend/otel/extract.go(及其测试 backend/otel/otel_test.go、extract_test.go),负责将 OTLP 数据提取为标准事件;
- 数据同步任务:deploy/otel-collector.yaml、deploy/opentelemetry-collector.Dockerfile,以及 docker/collector.yml 中的 Collector 配置;
- 各语言 SDK 的 OTel 集成,如 sdk/highlight-go、sdk/highlight-py 等。
总结:从架构文档到源码的对照速查表
| 架构文档中的组件 | 文档中的路径/函数 | 仓库实际位置(以根目录为起点) |
|---|---|---|
| SDKs | sdk/ | sdk/,含 Firstload、Client 与各语言 SDK |
| Public Graph | backend/public-graph/graph/schema.resolvers.go | backend/public-graph/graph/schema.resolvers.go(本地/public) |
| Private Graph | backend/private-graph/graph/schema.resolvers.go | backend/private-graph/graph/schema.resolvers.go(本地/private) |
| Workers | backend/worker.go | backend/worker/worker.go(processPublicWorkerMessage于 L359,Start于 L1179) |
对希望深入本仓库的开发者,建议的探索顺序是:先看 sdk/ 中语言对应 SDK 的上报逻辑,再到 Public Graph 端点确认摄入字段,最后阅读 backend/worker/worker.go 的消息分发与会话处理流程,即可在脑海中建立"采集 → 摄入 → 队列 → 处理 → 查询"的完整闭环。
- 可观测性
- 后端
【免费下载链接】highlight
highlight.io: The open source, full-stack monitoring platform. Error monitoring, session replay, logging, distributed tracing, and more.
相关推荐
highlight.io Cloudflare Worker SDK 全指南:错误监控、日志采集与分布式追踪实战
highlight.io Cloudflare Worker SDK 全指南:错误监控、日志采集与分布式追踪实战 本篇技术指南围绕 highlight.io 的
可观测性后端yfinance数据导出架构:从采集到应用的全链路优化
yfinance数据导出架构:从采集到应用的全链路优化 数据导出的现实困境与解决方案 你是否曾面临这样的困境:下载的金融数据格式混乱,无法直接导入分析工具?导出
数据分析金融科技BFL FLUX API 集成实战指南:端点选型、异步轮询与 Webhook 全链路解析
BFL FLUX API 集成实战指南:端点选型、异步轮询与 Webhook 全链路解析 本指南基于 OpenMontage 仓库中的 BFL API 集成 S
人工智能AI Agent音视频媒体生成工作流自动化
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考