ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

SkyWalking 链路数据处理实战:Segment 反序列化入 ClickHouse、Refs 拓扑去重与分钟级多粒度聚合

SkyWalking 链路数据处理实战:Segment 反序列化入 ClickHouse、Refs 拓扑去重与分钟级多粒度聚合 前言SkyWalking 是目前国内使用最广的开源 APM 之一Java Agent 无侵入采集链路数据Trace/Segment默认写入 ES、H2 或 BanyanDB。但在真实的生产可观测体系里我们往往还有更进一步的需求链路明细要和日志、指标放在同一个分析引擎里方便一次 SQL 关联排查“这个时间段报错的服务对应的慢 SQL 和错误日志是什么”调用拓扑要回流到 CMDB让应用关系数据资产化而不是只活在 APM 的界面里明细量大、查询以聚合统计为主需要冷热分层控制存储成本。我们的方案是把 OAP 输出到 Kafka 的序列化 Segment 数据消费下来用 SkyWalking 源码反序列化成明文 JSON逐字段解析后写入 ClickHouse 明细表同时利用 Segment 里的 Refs 引用还原服务调用关系去重后推送 Kafka 供 CMDB 消费再按端点/实例/服务/应用四个粒度做分钟级聚合。本文完整分享这套链路数据处理的设计字段映射规则、建表 DDL、拓扑去重逻辑和踩过的坑。表结构设计思路MergeTree 选型、集群化改造在上一篇《ClickHouse 可观测数据建模实战》里已经详细讲过本文聚焦链路数据这条流水线本身。一、整体链路与数据流先看全貌SkyWalking Agent | v OAP Server --kafka exporter-- Kafka(sw-trace-segments, 序列化Segment) | Flink 消费任务 | SkyWalking 源码反序列化 SegmentObject | 明文 JSON 数据集 | ------------------------------------------------------------ v v v v dwd_trace_segment (source,target) 分钟级窗口聚合 Database Span (明细表,TTL 7天) 拓扑去重 | 拆解 | | -------- | v v v v v 明细检索/ Kafka(trace- dws_trace_ Kafka(trace- dwd_trace_sql 关联分析 topology-ref) *_min 4级 agg-minute) (慢SQL明细) | | | v v v CMDB 服务关系 大盘/告警/报表 dws_trace_sql_min关键设计点消费侧不解析协议细节直接复用 SkyWalking 源码里的反序列化能力拿到SegmentObject避免自己啃 protobuf 格式版本升级时跟随 OAP 的序列化协议走一份数据三条出口明细入 ClickHouse、拓扑去重后回 Kafka、聚合结果双写ClickHouse Kafka下游各自订阅互不牵制明细与聚合分层 TTL明细 7 天、常规聚合 7 天、应用级和 SQL 聚合 30 天用存储策略把冷数据下沉到 HDD。二、Segment 数据结构从序列化字节到明文 JSON反序列化后得到的是一条条 Segment链路段。一个跨服务调用的 Trace 由多个 Segment 组成每个 Segment 对应一个服务实例内的执行片段核心字段如下脱敏后的真实样例{traceId:9a233cefb9f2441188c30e5e5f76c5a2.59.17318268143338523,traceSegmentId:9a233cefb9f2441188c30e5e5f76c5a2.59.17318268143338522,spans:[{spanId:0,parentSpanId:-1,startTime:1731826814333,endTime:1731826814335,refs:[],operationName:Mysql/JDBC/PreparedStatement/executeBatch,peer:192.168.10.31:3306,spanType:Exit,spanLayer:Database,componentId:33,isError:false,tags:[{key:db.type,value:Mysql},{key:db.instance,value:demo_db},{key:db.statement,value:}],logs:[],skipAnalysis:false}],service:shop::dc1::order-service,serviceInstance:46712e94e11b4560a99a1b73be170e39192.168.10.69,isSizeLimited:false}先建立几个直觉traceId / traceSegmentId一次请求一条 Trace每个服务实例内的执行各占一个 Segment。跨服务时靠 Refs 关联父子spansSegment 内的 Span 数组。parentSpanId -1的 Span 是入口 Span其operationName就是这个接口的端点名endpointspanTypeExit且spanLayerDatabase的 Span 是数据库访问带 db.type / db.instance 等 tags——第六节的慢 SQL 分析就靠它refs跨进程/跨线程时指向父 Segment。一个 Segment 里至多一个 Span 带 refs也可能没有——这个至多一个的约定后面解析时要用service 与 serviceInstance我们的 Agent 配置把 service 命名成应用::数据中心::服务名三段式SkyWalking 支持 service name 覆写这是常见做法业务维度的应用编码、数据中心信息就藏在里面serviceInstance 是实例UUIDIP结构一个 拆两半实例 ID 和 IP 都有了。三、明细入仓Segment 到 ClickHouse 的字段映射明细表dwd_trace_segment的设计原则是能用原生字段就不加工、加工字段必须可追溯。完整映射关系目标字段类型数据来源trace_idStringSegment 自带 traceIdsegment_idStringSegment 自带 traceSegmentIdspansStringspans 数组原文JSON 字符串保存refsString从带 refs 属性的 span 中提取异步场景可能多个ref_typeString带 refs 的 span 的 spanLayer 值service_idStringserviceInstance 按 拆分的前半段service_instance_idStringserviceInstance 原值service_nameStringservice 按 :: 分组的第 3 段app_codeStringservice 按 :: 分组的第 1 段data_centerStringservice 按 :: 分组的第 2 段endpoint_nameStringparentSpanId -1 的 span 的 operationNamestart_timeDateTime64(3)spans 中最早的 startTimeend_timeDateTime64(3)spans 中最晚的 endTimelatencyInt32endTime - startTime毫秒ipStringserviceInstance 按 拆分的后半段is_startInt8入口标记见下方规则 5is_errorInt8任一 span 的 isError 为 true 即记错误几条解析规则值得单独说明都是实际踩出来的refs 的恢复refs 挂在 span 上而不是 Segment 顶层提取时要遍历 spans 找带 refs 的那个 span至多一个找不到就存空数组。早期版本我们丢了 refs导致拓扑还原不完整后来专门补上ref_type 的语义直接取该 span 的 spanLayer 字符串如 Database/Http表结构上经历过 int 编码到 String 的调整——枚举值扩列时 int 映射表要跟着改String 更抗变化没有 refs 时存空字符串endpoint_name 取入口 Span不是随便取一个 operationName而是 parentSpanId -1 那个这样 JDBC、HTTP Client 这类 Exit Span 的操作名不会污染接口维度时间与耗时start_time / end_time 分别取 spans 的最早/最晚值精度用 DateTime64(3) 保留毫秒latency 用两者差值聚合时不用再二次解析is_start / is_error 的 0/1 约定is_start 0 表示是链路首节点spans 中不存在带 refs 的 span1 表示不是is_error 0 表示有错误1 表示正常——两列都是0 才是特殊值的 errno 风格和字段名的直觉相反查询时极易写反第八节还会说。建表 DDLCREATETABLEobs_dw.dwd_trace_segment(trace_id StringCOMMENT链路IDAgent原生生成,segment_id StringCOMMENTSegment实例/线程IDAgent原生生成,spans StringCOMMENTJSON格式的spans数组原文,refs StringCOMMENT父Segment引用异步场景可能多个,service_id StringCOMMENT服务ID,service_name StringCOMMENT服务名service第3段,service_instance_id StringCOMMENT实例IDserviceInstance 前半段,ref_type StringCOMMENT父调用类型取refs span的spanLayer,endpoint_name StringCOMMENT接口名parentSpanId-1的span,start_time DateTime64(3)COMMENTspans最早startTime,end_time DateTime64(3)COMMENTspans最晚endTime,latency Int32COMMENT耗时时长ms,ip StringCOMMENT服务IPserviceInstance 后半段,app_code StringCOMMENT应用编码service第1段,data_center StringCOMMENT数据中心service第2段,is_start Int8COMMENT是否首节点 0是 1否,is_error Int8COMMENT0异常 1正常,start_dateDate)ENGINEMergeTreePARTITIONBYtoYYYYMMDD(start_time)ORDERBY(start_time,trace_id)TTL start_datetoIntervalDay(7)SETTINGS index_granularity8192;明细保留 7 天按天分区TTL 到期自动清理排序键用 (start_time, trace_id)时间范围 链路 ID 点查是最高频的两种查询路径。四节点以上集群部署时换成 ReplicatedMergeTree Distributed 两层结构做法在《ClickHouse 可观测数据建模实战》里写过不再展开。四、服务拓扑用 Refs 还原调用关系并去重APM 界面里的拓扑图很好看但这份调用关系数据如果只活在 APM 里就太可惜了——把它沉淀下来推给 CMDB应用调用关系就成了可治理的数据资产。我们的做法规则一有 Refs 的 Segment 产生一条调用边。当前 Segment 的 service 是 targetrefs 里的 parentService 是 source当前 Segment 的 startTime 作为边的时间戳source refs.parentService target 当前segment.service time 当前segment.startTime规则二没有 Refs 的入口 Segment 补一个虚拟根。判断依据是入口 span 的 parentSpanId -1 且 refs 为空数组说明这条链路从这里开始补 source “0”虚拟根节点source 0 target 当前segment.service这样下游按 source 聚合时能直接区分入口流量和内部调用。规则三按 (source, target) 去重。高流量服务之间每秒产生成百上千个 Segment全部推送没有意义——边的语义就是存在调用关系所以内存里维护一个 (source, target) 集合出现过的边不再重复推送。规则四周期性批量输出。每隔一段时间参数可配也可通过配置文件临时调整把新增的边打成一批推送到 Kafka 主题trace-topology-ref报文格式[{source:0,target:shop::dc1::gateway,time:1713320388},{source:shop::dc1::gateway,target:shop::dc1::order-service,time:1713320388},{source:shop::dc1::order-service,target:shop::dc1::stock-service,time:1713320388}]三个字段足够下游消费source/target 是应用::数据中心::服务名的完整标识time 是数字时间戳。注意这里有一个格式兼容坑不同环境的 Agent service 命名模板不一样我们遇到过应用::数据中心::服务名三段和应用::环境_数据中心::服务名环境与数据中心用下划线拼在同一段两种解析侧必须按实际分段数兜底不能写死永远三段、第 2 段就是数据中心。五、多粒度分钟级聚合明细有了大盘和告警要的是聚合数字。按使用场景拆成四个粒度全部以 startTime 向下取整到分钟作为时间条件聚合表分组条件输出指标dws_trace_endpoint_minapp_code data_center service_name service_ip endpoint_namerequest_count / error_count / timeout_count / avg_used_timedws_trace_instance_minapp_code data_center service_name service_ip同上dws_trace_service_minapp_code data_center service_name同上 satisfied_count / tolerant_countdws_trace_app_minapp_code同上指标口径request_count分钟 分组条件下组内所有 Segment 的数量error_count组内 is_error 标记为错误的 Segment 数量avg_used_time组内所有 latency 的平均值mstimeout_count超时数量——依赖超时的定义阈值从哪来我们暂时搁置了这列宁可空着也不要给一个口径不清的数字satisfied_count / tolerant_count服务级表增加了满意/容忍计数配合平均耗时可以算 Apdex满意度指数这是服务级别才需要的视角。端点级 DDL其余粒度只是分组列的增减CREATETABLEobs_dw.dws_trace_endpoint_min(app_code StringCOMMENT应用编码,data_center StringCOMMENT数据中心,service_name StringCOMMENT服务名称,service_ip StringCOMMENT服务IP,endpoint_name StringCOMMENT端点,timeDateTimeCOMMENT统计分钟,request_count Int32COMMENT请求量,error_count Int32COMMENT错误数,timeout_count Int32COMMENT超时数,avg_used_time Int32COMMENT平均耗时ms)ENGINEMergeTreePARTITIONBYtoYYYYMMDD(time)ORDERBYtimeTTLtimetoIntervalDay(7)SETTINGS index_granularity8192;应用级聚合是查询频次最高、保留周期最长的表单独放宽 TTL 到 30 天并挂存储策略让冷数据落到 HDD 卷CREATETABLEobs_dw.dws_trace_app_min(app_code StringCOMMENT应用编码,timeDateTimeCOMMENT统计分钟,request_count Int32COMMENT请求量,error_count Int32COMMENT错误数,timeout_count Int32COMMENT超时数,avg_used_time Int32COMMENT平均耗时ms)ENGINEMergeTreePARTITIONBYtoYYYYMMDD(time)ORDERBYtimeTTLtimetoIntervalDay(30)SETTINGS storage_policyssd_hdd_policy,index_granularity8192;聚合结果同时写一份到 Kafka 主题trace-agg-minute——报表、告警等下游系统不必都来查 ClickHouse各自订阅即可。聚合这一步用通用任务编排就能实现分组 汇总算子是成熟算子不像明细和拓扑必须自定义代码这个取舍下一节展开。六、SQL 慢语句分析Database Span 的二次拆解第三节的明细表把 spans 原文存成了 JSON 字符串通用分析没问题但慢 SQL 排查是高频专项场景从 JSON 字符串里反复抽取太浪费。于是对spanLayer Database的 Span 做二次拆解单独落两张表。慢 SQL 明细表来源Database Span 的 tags 与 peerCREATETABLEobs_dw.dwd_trace_sql(app_code StringCOMMENT应用系统编码,data_center StringCOMMENT数据中心,service_name StringCOMMENT服务名,ip StringCOMMENT服务IP,trace_id StringCOMMENT链路ID,db_instance StringCOMMENT数据库实例ipportdatabase,db_type StringCOMMENT数据库类型,sql_statement StringCOMMENTSQL语句详细信息,start_time DateTime64(3)COMMENT开始时间,end_time DateTime64(3)COMMENT结束时间,used_time Int32COMMENT耗时时长ms,is_error Int8COMMENT0异常 1正常,error_msg StringCOMMENT异常信息,start_dateDate)ENGINEMergeTreePARTITIONBYtoYYYYMMDD(start_time)ORDERBY(start_time,sql_statement)TTL start_datetoIntervalDay(7)SETTINGS index_granularity8192;慢 SQL 分钟聚合TTL 30 天和应用级聚合同样的冷热策略CREATETABLEobs_dw.dws_trace_sql_min(app_code StringCOMMENT应用系统编码,data_center StringCOMMENT数据中心,service_name StringCOMMENT服务名,ip StringCOMMENT服务IP,sql_statement StringCOMMENTSQL语句详细信息,timeDateTimeCOMMENT统计分钟,request_count Int32COMMENT执行数,error_count Int32COMMENT异常数,avg_used_time Int32COMMENT平均耗时ms,max_used_time Int32COMMENT最慢执行时间ms)ENGINEMergeTreePARTITIONBYtoYYYYMMDD(time)ORDERBYtimeTTLtimetoIntervalDay(30)SETTINGS index_granularity8192;聚合粒度里加了max_used_time平均耗时会被高峰毛刺摊平最慢一次执行了多久是慢 SQL 治理时更被关心的数字。sql_statement 建议在写入前做参数归一化把字面量替换成 ?否则同一句 SQL 会因为参数不同炸出无数分组。七、为什么明细和拓扑没有走通用任务编排我们的数据平台有可视化的任务编排能力拖算子、配流。评估过把整条流水线都搬上去结论是分层实现环节编排能否实现卡点明细解析入仓不能缺反序列化算子吃不了序列化数据缺解析算子只能处理已结构化数据某字段子集存在某属性时给当前字段赋枚举值这类条件加列表达不出来拓扑去重不能缺去重算子分钟级聚合可以分组 汇总是成熟算子直接编排配置这个结论的意义在于划清平台与代码的边界编排平台适合结构化数据的确定性变换而涉及协议反序列化、状态化去重、跨字段条件逻辑的部分老老实实写 Flink 任务更可控。全部硬塞进编排的结果是每个算子都在勉强拼装排查问题时既看不到代码又没有算子级日志。八、踩坑清单最后把散落在前面各节的坑集中列一遍is_error / is_start 的 0/1 反直觉0 才代表有错误/“是首节点”errno 风格。SQL 里WHERE is_error 1查出来的是正常请求写报表前先把口径注释贴在查询旁边refs 在 span 上且至多一个 span 有解析时遍历查找而不是取固定下标没有 refs 时 ref_type 存空字符串而不是 null下游好处理ref_type 从 int 改 String枚举值跟着 SkyWalking 版本可能扩列int 映射表每次都要同步改String 直接存 spanLayer 更皮实service 命名模板不统一三段式和环境_数据中心合并式并存:: 拆分必须按分段数兜底判断这在多环境接入时一定会遇到spans 原文存 JSON 字符串看似浪费空间实际是救命设计——任何明细字段漏解析了都能从原文回补不用重放 KafkaTTL 分层要和查询频次对齐明细 7 天 / 常规聚合 7 天 / 应用级与 SQL 聚合 30 天再配合 ssd_hdd_policy 把冷分区压到 HDD存储成本能差出一个量级表名上线前 review 拼写我们的明细表第一版把 trace 拼成了 trance集群建好后改名成本极高只能将错就错——DDL 上线前多看一眼timeout_count 宁缺毋滥超时阈值来源没有定论之前这列保持空避免下游基于一个口径不明的数字做告警。结语这套链路数据处理跑通后SkyWalking 对我们来说不再只是一个看调用链的工具明细进 ClickHouse 与日志、指标同引擎关联分析调用关系回流 CMDB 沉淀为数据资产四级聚合支撑大盘与告警慢 SQL 独立成表专项治理。整个方案里最值得复用的思想是一份数据、多条出口、各取所需以及明细宁大勿丢、聚合分层降本。相关阅读ClickHouse 可观测数据建模实战指标、日志、链路、告警四类数据的 MergeTree 表设计与集群化改造云原生日志采集与检索架构实战K8s 日志和传统日志统一接入日志平台的设计方案Kafka ElasticsearchClickHouse 物化视图实战AggregatingMergeTree 实现监控指标多粒度聚合你们的链路明细数据保留多久7 天明细 30 天聚合的分层是否够用还是会更激进地全部落对象存储欢迎评论区交流。
返回列表