OpenTelemetry Collector 是 OpenTelemetry 的核心组件,但在底层基础设施(如 Kubernetes 节点)故障时,可能暴露出阻塞或延迟问题。本文通过一次因 Sampling 服务节点宕机引发的故障,结合代码分析其原因,并提供临时和长期解决方案。
一天,收到告警,OpenTelemetry 出现 Exporter Trace 异常的情况,具体表现为:


故障发生在一个多层遥测数据处理架构中,具体流程如下:
客户端 -> OpenTelemetry-Collector -> OpenTelemetry-Collector-Sampling -> Trace 后端服务初步排查确认,此次故障影响了第一层 Collector 的正常运行。
触发本次问题的 OpenTelemetry Collector 版本为 0.73.0(发布于 2023 年初)。此版本的负载均衡器实现存在已知问题,尤其是在处理端点下线时的阻塞行为尚未优化(后续版本通过 PR #31602 修复)。
通过日志、指标和代码分析,我们锁定了问题根源: Sampling 服务的一个 Pod 因所在节点宕机下线,节点宕机后,导致第一层 Collector 的 gRPC 连接未正常关闭,Collector 的数据发送和 DNS 解析组件在处理下线端点时发生阻塞。以下是详细故障链条,结合代码剖析:
Sampling 服务的一个 Pod 因节点宕机下线。由于是意外故障,gRPC 连接未执行正常关闭流程,客户端(第一层 Collector)gRPC 未立即感知服务器不可用。
dnsResolver组件每 5 秒解析 DNS,监控 Sampling 服务端点:
func (r *dnsResolver) resolve(ctx context.Context) ([]string, error) {
r.shutdownWg.Add(1)
defer r.shutdownWg.Done()
addrs, err := r.resolver.LookupIPAddr(ctx, r.hostname)
if err != nil {
_ = stats.RecordWithTags(ctx, resolverSuccessFalseMutators, mNumResolutions.M(1))
return nil, err
}
_ = stats.RecordWithTags(ctx, resolverSuccessTrueMutators, mNumResolutions.M(1))
var backends []string
for _, ip := range addrs {
var backend string
if ip.IP.To4() != nil {
backend = ip.String()
} else {
backend = fmt.Sprintf("[%s]", ip.String())
}
if r.port != "" {
backend = fmt.Sprintf("%s:%s", backend, r.port)
}
backends = append(backends, backend)
}
sort.Strings(backends)
if equalStringSlice(r.endpoints, backends) {
return r.endpoints, nil
}
// **关键点 1:端点变化触发回调**
r.updateLock.Lock()
r.endpoints = backends
r.updateLock.Unlock()
_ = stats.RecordWithTags(ctx, resolverSuccessTrueMutators, mNumBackends.M(int64(len(backends))))
r.changeCallbackLock.RLock()
for _, callback := range r.onChangeCallbacks {
callback(r.endpoints) // 调用 onBackendChanges
}
r.changeCallbackLock.RUnlock()
return r.endpoints, nil
}onBackendChanges处理端点变化:
func (lb *loadBalancerImp) onBackendChanges(resolved []string) {
newRing := newHashRing(resolved)
if !newRing.equal(lb.ring) {
// **关键点 2:加锁更新**
lb.updateLock.Lock()
defer lb.updateLock.Unlock()
lb.ring = newRing
ctx := context.Background()
lb.addMissingExporters(ctx, resolved)
lb.removeExtraExporters(ctx, resolved)
}
}removeExtraExporters下线多余端点:
func (lb *loadBalancerImp) removeExtraExporters(ctx context.Context, endpoints []string) {
endpointsWithPort := make([]string, len(endpoints))
for i, e := range endpoints {
endpointsWithPort[i] = endpointWithPort(e)
}
for existing := range lb.exporters {
if !endpointFound(existing, endpointsWithPort) {
// **关键点 3:Shutdown 清理积压数据**
_ = lb.exporters[existing].Shutdown(ctx)
delete(lb.exporters, existing)
}
}
}数据路由到 Sampling 服务的逻辑在traceExporterImp.consumeTrace中实现:
func (e *traceExporterImp) consumeTrace(ctx context.Context, td ptrace.Traces) error {
var exp component.Component
routingIds, err := routingIdentifiersFromTraces(td, e.routingKey)
if err != nil {
return err
}
for rid := range routingIds {
endpoint := e.loadBalancer.Endpoint([]byte(rid)) // 获取端点
// **关键点 4:根据 endpoint 获取 Exporter**
exp, err = e.loadBalancer.Exporter(endpoint) // 获取 Exporter
if err != nil {
return err
}
te, ok := exp.(exporter.Traces)
if !ok {
return fmt.Errorf("unable to export traces, unexpected exporter type: expected exporter.Traces but got %T", exp)
}
start := time.Now()
// **关键点 5:发送数据到 Sampling 服务**
err = te.ConsumeTraces(ctx, td)
duration := time.Since(start)
if err == nil {
_ = stats.RecordWithTags(
ctx,
[]tag.Mutator{tag.Upsert(endpointTagKey, endpoint), successTrueMutator},
mBackendLatency.M(duration.Milliseconds()))
} else {
_ = stats.RecordWithTags(
ctx,
[]tag.Mutator{tag.Upsert(endpointTagKey, endpoint), successFalseMutator},
mBackendLatency.M(duration.Milliseconds()))
}
}
return err
}loadBalancer.Exporter方法负责返回指定端点的 Exporter:
func (lb *loadBalancerImp) Exporter(endpoint string) (component.Component, error) {
// NOTE: make rolling updates of next tier of collectors work. currently, this may cause
// data loss because the latest batches sent to outdated backend will never find their way out.
// for details: https://github.com/open-telemetry/opentelemetry-collector-contrib/issues/1690
// **关键点 6:根据 endpoint 获取 Exporter 时加读锁**
lb.updateLock.RLock()
exp, found := lb.exporters[endpointWithPort(endpoint)]
lb.updateLock.RUnlock()
if !found {
return nil, fmt.Errorf("couldn't find the exporter for the endpoint %q", endpoint)
}
return exp, nil
}在 OpenTelemetry Collector v0.73.0 的故障场景中,Sampling 服务节点宕机导致 gRPC 连接未正常关闭,Shutdown 操作阻塞了负载均衡器的更新和数据发送,恢复时间长达 15 分钟。为缩短恢复时间,我们调整了 gRPC 的 Keepalive 参数,使 gRPC 客户端更快感知连接异常,从而加速端点下线流程。以下是具体的配置和分析。
loadbalancing:
protocol:
otlp:
compression: none
tls:
insecure: true
keepalive:
time: 10s
timeout: 3s
permit_without_stream: truereceivers:
otlp:
protocols:
grpc:
endpoint: 0.0.0.0:55681
keepalive:
enforcement_policy:
min_time: 8s
permit_without_stream: true社区通过 PR #31602 优化了 0.73.0 版本的问题:
在客户端 -> OpenTelemetry-Collector (v0.73.0) -> OpenTelemetry-Collector-Sampling -> Trace 后端服务架构中,Sampling 节点宕机导致 Collector 侧的 gRPC 连接 异常,Collector 的数据发送和 DNS 解析组件在处理下线端点时发生阻塞。恢复时间达 15 分钟。
通过版本背景和代码分析,我们理解了 0.73.0 的问题根源,希望这篇博文为类似场景提供参考!