Vector 中 Datadog APM Stats 的摄入设计:从 RFC 9862 到 datadog_traces Sink 的实现演进
【免费下载链接】vectorA high-performance observability data pipeline.项目地址: https://gitcode.com/GitHub_Trending/vect/vector
导读
本文以仓库 rfcs/2021-11-03-9862-ingest-apm-stats-along-traces-in-dd-agent-source.md 为核心主线,系统讲解 Vector 如何将 Datadogtrace-agent上报的 APM Stats(应用性能监控统计量)与 Traces 一并纳入datadog_agentsource,并最终由datadog_tracessink 按 Datadog 官方算法重新聚合、按 10 秒时间桶刷新的完整链路。读完本文,你将理解 APM Stats 的 Protobuf 数据形态、Vector 内部表示与 Sketch 转换、source 的多命名输出(named outputs)设计,以及该 RFC 在 src/sources/datadog_agent 与 src/sinks/datadog/traces 中落地后的真实实现细节与当前取舍。
背景:为什么 Traces 之外还要摄入 APM Stats
Datadog 的trace-agent向 Datadog 后端提交的数据并不仅仅是 Trace(链路),还包括APM Stats——即对每个被插桩资源(某段被监控的代码)运行时间的统计量,这些统计量基于收到的100% Traces按时间聚合而成。APM Stats 的核心价值在于:它能直接暴露代码热点、极大简化聚合分析,因此 Vector 不能丢弃这部分数据。
RFC 指出,APM Stats 可能由追踪库(tracing lib)在客户端直接计算,但更多时候由trace-agent中的concentrator组件计算;二者最终通过同一套发送代码(stats writer)上报。无论来自哪一方,trace-agent都会把 APM Stats 发送到与 Trace 完全相同的端点(endpoint),仅仅路径(path)不同,因此可以依据路径轻松区分 Trace 与 APM Stats 负载。
在 Vector 中,datadog_agentsource 对/api/v0.2/traces与/api/v0.2/stats两个路径分别建了 warp 过滤器,见 src/sources/datadog_agent/traces.rs。
APM Stats 的数据形态:ClientGroupedStats Protobuf
APM Stats 负载是若干ClientGroupedStats的聚合结果。RFC 中给出了该结构的 protobuf 字段定义,这也是理解后续转换逻辑的基础:
string service = 1; string name = 2; string resource = 3; uint32 HTTP_status_code = 4; string type = 5; string DB_type = 6; // db_type might be used in the future to help in the obfuscation step uint64 hits = 7; // count of all spans aggregated in the groupedstats uint64 errors = 8; // count of error spans aggregated in the groupedstats uint64 duration = 9; // total duration in nanoseconds of spans aggregated in the bucket bytes okSummary = 10; // ddsketch summary of ok spans latencies encoded in protobuf bytes errorSummary = 11; // ddsketch summary of error spans latencies encoded in protobuf bool synthetics = 12; // set to true on spans generated by synthetics traffic uint64 topLevelHits = 13; // count of top level spans aggregated in the groupedstats可以看到,这是一组可以映射为多种度量的字段组合:
- 计数类字段:
hits、errors、top_level_hits(分别统计被聚合的 span 总数、错误 span 数、顶层 span 数); - 时长字段:
duration(时间桶内 span 的总时长,单位为纳秒); - 分布类字段:
okSummary与errorSummary分别是成功/失败 span 延迟的ddsketch 摘要; - 标识类字段:
service、name、resource、type、HTTP_status_code、synthetics等用于后续聚合与过滤的维度。
RFC 特别指出:虽然在 proto 定义里okSummary/errorSummary只是无结构的bytes切片,但实际填充的是protobuf 编码的 ddsketch。Vector 内部的 Sketch 数据结构本身同样「重度基于 ddsketch」(参见 lib/vector-core/src/metrics/ddsketch.rs 中的AgentDDSketch实现),因此 APM Stats 的 sketch 与 Vector 内部表示之间存在相互转换的可行性,但这需要不小的实现工作量——这正是本文后面要展开的设计分歧点。
关键设计抉择:APM Stats 在 Vector 内部用什么事件表示
RFC 给出了两条截然不同的路径:
- 作为 Log 事件:每个
ClientGroupedStats映射为一条 log event; - 作为 Metric 事件:
ClientGroupedStats中的每个数值拆成一条 metric,上层所有信息存为 tags,再交由datadog_tracessink 做重新聚合——代价是 sink 侧需要相当复杂的重聚合逻辑。
两条路径还可以混合:例如允许单个 metric 事件携带多个 metric sample,或允许 log event 内嵌 metric(即给Value枚举增加Metric类型)。
此外,RFC 还抛出了datadog_agentsource 输出流组织的问题:假设 source 已接受 Datadog Agent metrics,则可能出现一条与普通 metrics 无关的第二条 metric 流。由于 APM Stats 事件必须与 Traces 一起路由(通常与 core-agent 上报的普通 metrics/logs 走不同路径),RFC 建议重组datadog_agentsource,具体有两种路线(详见下文「Source 重组」小节)。
Source 重组:单 source 多输出与多 source 方案
RFC 列出了两种重组思路:
方案 A:保留单个datadog_agentsource
- 添加一个
agent: <TYPE>开关,TYPE可为core(同时支持 logs 与 metrics,也可用logs/metrics单独限定)、trace,未来可扩展process、security等; - 或为 source 引入多输出能力,例如
<SRC_ID>.metrics、<SRC_ID>.logs、<SRC_ID>.traces、<SRC_ID>.apm_stats。
方案 B:每种 Datadog Agent 一个 Vector source
datadog_core(或沿用现有datadog_agent)接收收集 logs/metrics 的 core Agent 数据;datadog_trace支持trace-agent发送的全部数据;- 随支持列表增长继续派生
datadog_process、datadog_security等。
最终方案 B 的配置形态如下(RFC 中给出的「功能等价」示例,注意这里为每种数据类型分配了不同端口):
[sources.dd_in_logs] type = "datadog_logs" address = "[::]:8081" [sources.dd_in_metrics] type = "datadog_metrics" address = "[::]:8082" [sources.dd_in_traces] type = "datadog_traces" address = "[::]:8083" [sinks.dd_traces] type = "datadog_traces" inputs = ["dd_in_traces" ] [sinks.dd_out_logs] type = "datadog_logs" inputs = ["dd_in_logs"] [sinks.dd_out_metrics] type = "datadog_metrics" inputs = ["dd_in_metrics"] [sinks.debug] type = "console" inputs = ["dd_in_*"] encoding.codec = "json"方案 A 的用户体验:命名输出(named outputs)
为了免去在拓扑中用复杂且不可靠的routetransform 去区分「同样表示为 log 的 traces」「core Agent 来的普通 metrics」与「trace-agent 来的 APM stats metrics」,RFC 建议扩展 remap transform 已有的 failed-event-routing 行为,让 source 也能暴露多个命名输出。RFC 设想的配置如下:
[sources.dd_agents] type = "datadog_agent" address = "[::]:8081" [sinks.dd_traces] type = "datadog_traces" inputs = ["dd_agents.traces", "dd_agents.apm_stats" ] [sinks.dd_logs] type = "datadog_logs" inputs = ["dd_agents.logs"] [sinks.dd_metrics] type = "datadog_metrics" inputs = ["dd_agents.metrics"] [sinks.debug] type = "console" # Optionally the non-suffixed name could receive everything, this will be configurable inputs = ["dd_agents"] encoding.codec = "json"实现上,RFC 建议:
- 把原本只用于 transform 的
named_outputs能力扩展到 source,使其暴露<SRC_ID>.<OUTPUT_NAME>形式的命名输出; - 在
datadog_agent中增加<SRC_ID>.traces、<SRC_ID>.apm_stats、<SRC_ID>.metrics、<SRC_ID>.logs四个命名输出; - 无后缀输出要有可预期行为,可增加
top_level_output开关让用户选择无后缀输出拿到的数据类型。
当前源码中的落地情况
从当前仓库源码看,方案 A 的「多输出」已落地:datadog_agentsource 在 src/sources/datadog_agent/mod.rs 中定义了LOGS/METRICS/TRACES/LLMOBS四个端口常量,并提供了multiple_outputs配置项(默认false):
If this is set to
true, logs, metrics (beta), and traces (alpha) are sent to different outputs. For a source component namedagent, the received logs, metrics (beta), and traces (alpha) can then be configured as input to other components by specifyingagent.logs,agent.metrics, andagent.traces, respectively.
当multiple_outputs=true时,outputs()方法会通过.with_port(...)分别为 logs、metrics、traces、llmobs 生成独立的SourceOutput;否则所有数据类型走同一个无后缀输出(DataType::all_bits())。同时每个数据类型还有对应的禁用开关:disable_logs、disable_metrics、disable_traces、disable_llmobs。
值得注意的是:当前实现中APM stats 并没有作为独立命名输出暴露,traces.rs的build_stats_filter对/api/v0.2/stats路径直接返回 200 OK,源码注释明确说明:
APM stats are discarded on purpose, they will be computed in the
datadog_tracessink
也就是说,最终实现的走向是 RFC 中「备选方案」的变体:APM Stats 不在 source 侧解码,而是在datadog_tracessink 侧根据 100% Trace 重新计算。这与 RFC 早期设想(source 摄入并转换 APM stats)有所不同,但完全呼应了 RFC 在「Alternatives」与「Future Improvements」中提到的思路——即在 sink 中为任意 trace 格式计算 APM stats。
实现方案:把 APM Stats 转为 Vector Metrics 再重聚合
RFC 主提案中的实现计划分为三大块,互为独立:
- 重组
datadog_agentsource(详见上节); - 把 APM Stats 全部导入为标准的 Vector metric:
- 将每个
ClientGroupedStats拆成相关 metrics,携带全部上层元数据以保证无损透传(passthrough)场景与与 traces 同等的过滤/路由能力; - APM stats 的 sketch 转换为 Vector 内部 sketch,Vector 内部 sketch 获得参数化的
gamma与maxbin(默认仍取 agent 的 sketch 值);
- 将每个
datadog_tracessink 重新聚合:- 按相关维度使用
Partitionertrait 对 metric 重新聚合,重建 APM stats 负载; - 进入的 metric 被缓冲,填充与
ClientGroupedStats基础对象匹配的结构体,并按trace-agent使用的同类 key 存入 map; - 每 10 秒(
trace-agent的发送间隔)序列化并刷新到 Datadog;为吸收迟到的 metric,sink 需要保留过去 2~3 个时间桶并相应延迟刷新,依赖trace-agent在负载中保存的桶时间戳。
- 按相关维度使用
RFC 还规划了与 sketches-rs 生态对齐的后续演进:待 sketches-rs 生产可用后,将 Vector 的 sketch 实现切换为该 crate,并把转换逻辑从 traces 处理中剥离(届时datadog_agentsource 与datadog_metricssink 只需处理 Vector sketch 与 agent 变体之间的转换)。
当前仓库的 sink 侧实现:重聚合的完整脉络
虽然 source 侧选择了丢弃 APM stats,但datadog_tracessink 侧的 APM stats 重聚合实现是完整的,且代码注释明确表示其建模紧贴 Datadog trace-agent。整个模块位于 src/sinks/datadog/traces/apm_stats,分为aggregation.rs、bucket.rs、flusher.rs、weight.rs四个文件。
10 秒时间桶与聚合 key
src/sinks/datadog/traces/apm_stats/mod.rs 定义了核心常量与数据结构:
/// The duration of time in nanoseconds that a bucket covers. pub(crate) const BUCKET_DURATION_NANOSECONDS: u64 = 10_000_000_000;模块文档明确指出:该模块「基于进入 sink 的 trace 事件计算 APM 统计量,其建模紧密参照 Datadog Agent 的 trace-agent 组件,并以 10 秒间隔、独立于 sink 的 trace 负载,发送由同一算法格式化并聚合的 StatsPayload 包」。
同时,mod.rs中完整定义了与 Datadog 官方 stats.proto 一一对应的 msgpack 结构(StatsPayload、ClientStatsPayload、ClientStatsBucket、ClientGroupedStats),字段名与 go 生成代码完全一致(serde(rename_all = "PascalCase"),并对RuntimeID、ContainerID、HTTPStatusCode、DBType做了显式 rename),其中ClientGroupedStats恰好对应 RFC 展示的 protobuf 结构。
Aggregator:把 Trace 变成统计
src/sinks/datadog/traces/apm_stats/aggregation.rs 实现了核心聚合器:
AggregationKey由PayloadAggregationKey(env、hostname、version、container_id)与BucketAggregationKey(service、name、resource、type、status_code、synthetics)组成,聚合维度与 RFC 描述的「按 trace-agent 同类 key」一致;Aggregator::handle_trace遍历每条 trace 的 spans,依据 trace-agent 的 span 标记(_top_level、_dd.measured、_dd.partial_version)筛选需要计量的 span,再进行加权统计;- 桶时间通过
align_timestamp将 span 结束时间对齐到 10 秒边界,太旧的 span 会被钳制到允许的最老桶; BUCKET_WINDOW_LEN = 2:正常运行时只保留最近 2 个桶长度(即 20 秒)的缓存,超过即刷新——这正好实现了 RFC 中「保留过去 2~3 个桶、延迟刷新以吸收迟到 metric」的要求;flush(force)在正常模式按now - 10s * 2计算刷新截止时间,force=true(进程退出时)则刷新全部剩余桶。
Bucket:hits / errors / duration 与 ddsketch 编码
src/sinks/datadog/traces/apm_stats/bucket.rs 中的GroupedStats持有hits、top_level_hits、errors、duration四个浮点计数器,以及ok_distribution/err_distribution两个AgentDDSketch:
- 每个 span 进入时按权重累加计数,错误 span 与成功 span 的延迟分别插入各自的 sketch;
encode_sketch把 Vector 内部的AgentDDSketch转换为 proto/vector/ddsketch_full.proto 定义的DdSketchprotobuf(携带gamma、index_offset、Interpolation::None以及正/负 store 与零计数),用prost编码为字节——对应 RFC 中「把 Vector 内部 sketch 转换回 APM stats 的 sketch 字段」;convert_stores按符号位把 sketch 的 bin 拆分为正/负/零三部分。
AgentDDSketch的默认参数定义在 lib/vector-core/src/metrics/ddsketch.rs:AGENT_DEFAULT_BIN_LIMIT = 4096、AGENT_DEFAULT_EPS = 1.0/128.0、AGENT_DEFAULT_MIN_VALUE = 1.0e-9,这正对应 RFC 中「默认仍取 agent 的 sketch 值」的设想(gamma 由 EPS 推导)。
Flusher:10 秒独立线程刷新
src/sinks/datadog/traces/apm_stats/flusher.rs 运行一个独立异步线程flush_apm_stats_thread:
- 以
BUCKET_DURATION_NANOSECONDS(10 秒)为间隔周期刷新最老的桶; - sink 关闭时通过 oneshot tripwire 通知线程做
force刷新,全部剩余桶发完后再确认退出; - 负载经
rmp_serde::to_vec_named编码为msgpack(与 trace-agent 的 stats_gen.go 一致),压缩后发送到 stats 端点,content-type为application/msgpack。
端点与路由信息在 src/sinks/datadog/traces/config.rs 中定义:DatadogTracesEndpoint::{Traces, APMStats}两个枚举值,分别映射到https://trace.agent.<site>/api/v0.2/traces与/api/v0.2/stats。
Partition 与 sink 主流程
src/sinks/datadog/traces/sink.rs 中的EventPartitioner实现了Partitionertrait,用api_key、env、hostname、agent_version、target_tps、error_tps构造PartitionKey;TracesSink::run_inner把 trace 流按 key 分批,并额外启动 APM stats 刷新线程,在 sink 主流程结束后通过 oneshot 等待该线程完成最终刷新再退出,避免进程在 stats 尚未发完时被终止。
整个计算流程在 src/sinks/datadog/traces/apm_stats/mod.rs 的compute_apm_stats(key, aggregator, trace_events)中触发:先更新 agent 级属性,再逐条 trace 交给 aggregator 处理。其内部机制可在测试中进一步验证,例如 src/sinks/datadog/traces/tests.rs 与 apm_stats 模块的集成测试(datadog-traces-integration-testsfeature)。
备选方案与 RFC 的取舍理由
RFC 明确讨论了「不摄入 APM stats」的两种替代方案:
- 完全丢弃 APM stats:不可取,会直接导致用户体验劣化(用户会失去对执行时间的洞察);
- 在 trace-agent 侧禁用采样、由 datadog_traces sink 计算 APM stats:可行,但对首版实现而言工作量过大(需要在计算逻辑之上再支持 plain ddsketch),且为了匹配当前 APM stats 的精度,Vector 必须收到 100% 的 traces——这不一定总能做到。不过该方案能为「无论 traces 来自何方都能计算通用 APM stats」铺路。
关于内部表示,备选方案还包括:用带数字字段的 log event 表示 APM stats,或采用混合方案(log event 携带 metric 字段 / metric event 携带多值)。关于 sketch,由于 APM stats 的 sketch 与 Vector 内部表示并非完全一致,也可以选择不解码这些 sketch,在 Vector 内部把它们当作不透明字节切片保留,从而省去转换所需的 plumbing——这与当前实现中「sketch 由 sink 端重新生成」的取舍形成了有趣的对照。
演进路线:Plan of Attack 与未来改进
RFC 末尾给出了实施路线图:
- 实现 source 的多输出(multiple outputs per source)选项
- 实现 APM stats 解码为 Vector metrics
- 为
datadog_tracessink 增加 APM 支持 - (视 sketches-rs 时间线)把 Vector 内部 sketch 切换到 sketches-rs crate
未来改进方向包括:在datadog_tracessink 中为任意 trace 格式计算 APM stats;若 APM stats 以 log event 表示,则对 schema 施加约束会非常有用。
结合当前仓库代码可以看到:第 1、3 项已实质落地(multiple_outputs与 apm_stats 聚合/flusher),第 2 项则被「sink 端重计算」的路线所替代(datadog_agentsource 有意丢弃 stats 请求体并返回 200 OK),而 sketch 生态迁移仍在演进中。理解这条从 RFC 设计到源码实现的完整脉络,能帮助你在实际部署 Datadog 相关组件时,正确选择datadog_agentsource 的multiple_outputs配置、datadog_tracessink 的default_api_key等参数,并清楚 APM Stats 在 Vector 管道中的真实生命周期。
【免费下载链接】vectorA high-performance observability data pipeline.项目地址: https://gitcode.com/GitHub_Trending/vect/vector
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考