Tempo 数据入口枢纽:Distributor 组件的接收、校验、限流与路由机制深度解析
Tempo 数据入口枢纽Distributor 组件的接收、校验、限流与路由机制深度解析【免费下载链接】tempoGrafana Tempo is a high volume, minimal dependency distributed tracing backend.项目地址: https://gitcode.com/GitHub_Trending/tempo1/tempoDistributor 是 Grafana Tempo 分布式追踪后端写入路径的第一站负责接收来自被插桩应用的所有 span 数据执行同步的租户级限流校验并依据部署模式将数据路由到 Kafka微服务模式或进程内的 live-store 与 metrics-generator单体模式。本文以 Tempo 仓库中 distributor 架构文档 为主线结合 distributor 模块源码 与配置文档系统讲解其接收协议、限流策略、分片路由与观测指标帮助你理解并运维 Tempo 的写入路径。Distributor 在 Tempo 架构中的定位Distributor 是所有 trace 数据进入 Tempo 的唯一入口。它从被插桩的应用接收 span并按配置的租户限制对数据进行校验。它的转发方式完全取决于 部署模式微服务模式microservicesDistributor 按 trace ID 对 trace 分片shard写入 Kafka。下游的 block-builder、live-store 和 metrics-generator 各自独立地从 Kafka 消费数据。这是生产环境推荐的模式。单体模式monolithicDistributor 在进程内直接把数据推送给 live-store 和 metrics-generator不需要 Kafka。此模式适合开发环境与低到中等流量场景。在 components/_index.md 描述的组件拓扑中Distributor 是写入路径write path的第一环关于其他组件live-store、block-builder、metrics-generator的职责可进一步参考对应文档。接收链路基于 OpenTelemetry Collector 的多协议接入Distributor 复用了 OpenTelemetry Collectorjaegerreceiver、zipkinreceiver、otlpreceiver与kafkareceiver。支持的 span 协议格式OpenTelemetry Protocol (OTLP)支持 gRPC 与 HTTP这是官方推荐格式JaegerThrift 与 gRPCZipkinKafka官方建议尽可能使用 OTLP over gRPC无论是 Grafana Alloy 还是 OpenTelemetry Collector 都原生支持 OTLP 导出。默认接收器配置如果在配置中未显式指定receiversDistributor 会启用 config.go 中定义的默认接收器var defaultReceivers map[string]interface{}{ jaeger: map[string]interface{}{ protocols: map[string]interface{}{ grpc: nil, thrift_http: nil, }, }, otlp: map[string]interface{}{ protocols: map[string]interface{}{ grpc: nil, }, }, }即默认开启JaegergRPC thrift_http与OTLP gRPC。完整的distributor配置块结构见 configuration/_index.md可以按需启停协议distributor: receivers: otlp: protocols: grpc: # 默认监听 localhost:4317 http: # 默认监听 localhost:4318 jaeger: protocols: thrift_http: grpc: thrift_binary: thrift_compact: zipkin: kafka:需要注意接收器默认监听 localhost生产环境要配置为可监听外部接口的 IP且只应启用真正需要的接收器。shim 层的关键细节Tempo 通过receiversShimshim.go把 OTel Collector 的 receiver 与 Tempo 自己的PushTraces对接起来。有两个值得注意的实现细节元数据注入对otlpHTTP、zipkin、jaegerthrift_http接收器Tempo 会开启IncludeMetadata将 HTTP 头注入 context这是后续基于头信息做租户认证Authentication的前提shim.go。可重试错误包装当 Distributor 返回ResourceExhausted错误且集群启用了重试信息时wrapErrorIfRetryable会为 gRPC 错误附加RetryInfo详情shim.go让 OTLP exporter 等客户端知道应在RetryDelay之后重试而不是立即放弃或反复重试。同步校验与限流写入路径上唯一同步执行的限制在转发数据之前Distributor 会依据配置的摄取限制校验进入的数据。这些是唯一在摄取时同步强制执行的限制其余限制如max_live_traces_bytes、max_bytes_per_trace都在下游异步执行。请求处理主流程从 distributor.go 的PushTraces可以看到完整的处理流水线执行可选的TracePushMiddleware钩子失败仅记录日志不阻断推送即 fail-open从 context 提取租户 IDvalidation.ExtractValidTenantID无效则返回InvalidArgument统计限制前字节数tempo_distributor_ingress_bytes_total空批次直接返回执行checkForRateLimits同步限流将 OTel 格式转换为 Tempo 内部tempopb.Trace两者 wire 兼容可选地记录接收到的 span 日志、生成调试指标按 trace ID 重组批次requestsByTraceID同时校验 trace ID / span ID 合法性、截断超长属性依据pushSpansToKafka开关走 Kafka 写入或进程内直推。速率与突发限流摄取速率限制ingestion rate limit定义了每个租户每秒允许的最大字节数超限时客户端会收到RATE_LIMITED错误gRPC 码为ResourceExhausted错误文本以overrides.ErrorPrefixRateLimited即RATE_LIMITED开头。摄取突发大小ingestion burst size则控制允许超出持续速率的最大突发量。对应的限流实现见 checkForRateLimits它调用ingestionRateLimiter.AllowN判断本次批次字节数是否放行拒绝时调用overrides.RecordDiscardedSpans(spanCount, overrides.ReasonRateLimited, userID)记录丢弃指标并返回包含本地/全局速率、突发大小与租户信息的详细错误便于排查。限流器的构造在 Newlimiter.NewRateLimiter(ingestionRateStrategy, 10*time.Second)其中的策略实现位于 ingestion_rate_strategy.go。local 与 global 两种速率策略Tempo 提供两种摄取速率策略源码中对应localStrategy与globalStrategy两个实现策略适用场景工作方式local默认希望每个 distributor 各自独立处理一个固定速率接受集群总速率随实例数增长每个 distributor 强制完整的rate_limit_bytes。例如 5 个 distributor、每个30 MB/s集群总吞吐最高150 MB/sglobal需要一个不随 distributor 数量变化、可预期的集群级摄取预算将rate_limit_bytes除以健康 distributor 数量。例如 5 个 distributor、每个6 MB/s合计仍是30 MB/s在globalStrategy的实现中ingestion_rate_strategy.go健康实例数取自 distributor ring 中处于ACTIVE状态的实例数量若为 0 则退化为直接使用配置值。burst 大小burst_size_bytes不受rate_strategy影响始终按实例独立生效。配置位于 overrides 的defaults.ingestion块configuration/_index.mdoverrides: defaults: ingestion: rate_strategy: local # global | local默认 local rate_limit_bytes: 30000000 # 默认 30MB/s burst_size_bytes: 20000000 # 默认 20MB始终按实例生效注意rate_strategy是集群级设置无法按租户覆盖这在 runtime_config_overrides.go 的注释中有明确说明而rate_limit_bytes和burst_size_bytes均可在租户级 overrides 中配置。各摄取限制的执行位置与策略影响配置项执行组件受rate_strategy影响rate_limit_bytesDistributor是burst_size_bytesDistributor否始终按实例max_traces_per_userLive-store否始终按实例max_global_traces_per_user运行时未强制执行否仅指标默认关闭max_bytes_per_traceLive-store、block-builder否文档中还特别指出max_live_traces_bytes这类限制由下游 live-store异步强制执行而max_bytes_per_trace同样在下游执行——在微服务模式下还包括 block-builder。也就是说Distributor 只负责最前端的同步限流更复杂的单 trace 大小、活跃 trace 数量限制发生在数据消费端。Distributor ringglobal 策略的基石启用global策略时Distributor 需要借助一个专用的 distributor ring 来统计健康实例数。其配置结构见 distributor_ring.go关键配置项distributor: ring: kvstore: store: memberlist # 默认 memberlist prefix: collectors/ heartbeat_period: 5s # 心跳周期 heartbeat_timeout: 5m # 心跳超时视为不健康 instance_id: hostname # ring 中注册的实例 ID instance_interface_names: [eth0, en0] instance_addr: ip # 默认取 instance_interface_names instance_port: int enable_inet6: false与常规 hash ring 不同distributor ring每个实例只注册 1 个 tokenringNumTokens 1其用途并非数据分片而纯粹是统计HealthyInstancesCount作为 global 速率的除数distributor.go 注释明确说明。健康实例数的计算只统计ACTIVE状态、不扩展到LEAVINGdistributor.go。此外ring 启用了 auto-forget 机制实例在2 * heartbeat_timeout时间内失联会被自动从 ring 移除distributor.go。记录被丢弃的 span当 Distributor 因限流等原因拒绝 span 时会递增tempo_discarded_spans_total指标并携带reason标签标明丢弃原因。该指标的实现在 modules/overrides/discarded_spans.go标签还包括tenant。当前定义的可导出丢弃原因包括rate_limited租户速率超限trace_too_large单 trace 的 span 数过多live_traces_exceeded该租户的活跃 trace 数超限invalid_trace_idtrace ID 不是 128 bitrequestsByTraceID中校验distributor.goinvalid_span_idspan ID 不是 64 bit 或全零distributor.go为调试而记录被丢弃的 span若要逐条记录被丢弃的 span 以便调试可在distributor配置块开启日志distributor: log_discarded_spans: enabled: true include_all_attributes: falseinclude_all_attributes: true会产生更冗长的日志包含 span 属性有助于定位行为异常misbehaving的客户端。filter_by_status_error: true时只记录状态为错误的 span见 logSpans 的实现。此外还有同构的log_received_spans记录每个收到的 span官方不建议生产环境开启与metric_received_spans按租户/span 名/service 记录调试指标可设root_only: true只看根 span。对应配置参数均定义在 config.go。微服务模式按 trace ID 分片写入 Kafka在微服务模式下校验通过后 Distributor 会对 trace ID 做哈希util.TokenFor(userID, traceID)distributor.go进行分片查询partition ring确定哪些 Kafka 分区处于 active 状态将记录写入对应的分区并等待 Kafka 确认后才向客户端返回响应。只有 Kafka 返回成功后写入才算成功这保证了一旦客户端收到成功响应数据就已经被持久化存储。写入实现见 sendToKafka先按租户对 partition ring 做ShuffleShard再通过ring.DoBatchWithOptions把同一分片的 trace 批量编码为 Kafka 记录并ProduceSync同步等待结果随后统计tempo_distributor_kafka_records_per_request、tempo_distributor_kafka_write_latency_seconds与tempo_distributor_kafka_write_bytes_total指标。分片Partitioning设计Distributor 按trace ID分片意味着同一 trace 的所有 span 都会进入同一个 Kafka 分区。这一设计带来两个关键收益block-builder可以在单次消费周期内构建出同一 trace 的所有 span 共置的 blocklive-store可以从单个分区服务完整的 trace 查询无需跨分区协调。需要特别强调Distributor 使用Tempo 自身的 partition ring而非 Kafka 原生的分区路由来决定目标分区这让 Tempo 可以独立于 Kafka 控制分区的生命周期。分区 ring 维护了 pending、active、inactive 三种分区状态Distributor 只向 active 分区写入数据细节可参考 partition-ring 文档。单体模式进程内直推在单体模式-targetall下Distributor 不再初始化 Kafka producer也不使用 partition ring 做路由而是通过LocalPushTargets回调在进程内直接把数据推送给 live-store 与 metrics-generatorpushLocal。客户端响应在 live-store 接受数据后返回。单体模式的实现要点distributor.goPushBytesFunc把预序列化的 trace 字节推送给进程内 live-storePushSpansFunc把 span 数据推送给进程内 metrics-generator对 metrics-generator 的推送会先进入一个按租户隔离的异步队列generatorForwarderforwarder.go队列大小与 worker 数可通过 overrides 的metrics_generator_forwarder_queue_size/metrics_generator_forwarder_workers配置默认分别为 100 与 2且会周期性监听 overrides 变化并平滑迁移队列。单体模式下所有组件共享同一进程资源查询高峰可能影响写入吞吐无法独立扩缩容因此官方建议生产环境使用微服务模式deployment-modes.md。关键指标Distributor 暴露了丰富的 Prometheus 指标可用于监控摄取状态。以下是指南中列出的核心指标及其源码定义位置指标说明源码定义tempo_distributor_spans_received_totalDistributor 接收到的 span 总数按租户distributor.gotempo_discarded_spans_total被丢弃的 span带reason与tenant标签discarded_spans.gotempo_distributor_bytes_received_total接收到的总字节数限制通过后统计distributor.gorate(tempo_distributor_spans_received_total[5m])当前摄取速率span/s由接收计数经 PromQL 派生—以下相关指标同样值得关注均定义于 distributor.gotempo_distributor_ingress_bytes_total限制执行之前接收的字节数与bytes_received_total对比可评估限流拦截量tempo_distributor_push_duration_seconds处理并路由一个批次的总耗时receiver/shim.gotempo_distributor_traces_per_batch每批次中的 trace 数直方图tempo_distributor_kafka_records_per_request/tempo_distributor_kafka_write_latency_seconds/tempo_distributor_kafka_write_bytes_totalKafka 写入观测tempo_distributor_metrics_generator_pushes_total/_failures_total推送给 metrics-generator 的次数与失败次数tempo_distributor_attributes_truncated_total被截断的属性 key/value 数按租户与 scoperesource/scope/span/event/linktempo_distributor_debug_spans_received_totalmetric_received_spans开启时的调试计数。其他实用配置与进阶机制属性大小截断distributor.max_attribute_bytes默认2048即 2KBconfig.go用于控制单个属性 key 或 value 的最大字节数超限的字符串属性会在入库前被截断processAttributes将该参数设为0可关闭此检查。每租户可通过 overrides 的ingestion.max_attribute_bytes覆盖distributor.go。截断事件会以限速日志每秒 1 条记录首条示例便于诊断。资源耗尽后的重试指引distributor.retry_after_on_resource_exhausted默认5s控制 Tempo 在返回 gRPCResourceExhausted时携带的重试等待时长设为0表示集群级禁用重试信息。租户级可通过 overrides 的ingestion.retry_info_enabled单独控制distributor.go。Forwarders异步复制distributor.forwarders支持把摄入的 trace 异步复制best-effort到外部端点目前仅支持otlpgrpc后端且需在租户 overrides 中启用失败不会影响主链路写入forwarder 配置。租户分片与全局速率联动Kafka 写入前会按ingestion_tenant_shard_size对 partition ring 做 shuffle sharddistributor.go将租户流量收敛到其分片内的分区子集这是微服务模式下控制分区扇出与数据局部性的重要旋钮。相关资源Distributor 架构文档本文依据部署模式说明Partition ring 机制Tempo 写入路径架构图Distributor 完整配置参考 与 Ingestion rate strategy 章节源码入口modules/distributor/distributor.go、modules/distributor/config.go、modules/distributor/ingestion_rate_strategy.go、modules/distributor/receiver/shim.go【免费下载链接】tempoGrafana Tempo is a high volume, minimal dependency distributed tracing backend.项目地址: https://gitcode.com/GitHub_Trending/tempo1/tempo创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
上一篇/下一篇内容由系统自动关联
返回资讯列表 →