FastAPI从OpenTelemetry Collector接收追踪数据丢失问题求助
我使用FastAPI接收OpenTelemetry Collector(otelcol)的追踪数据,但存在数据丢失问题。具体来说,每次API调用应生成一条追踪记录,但当我用Shell脚本每秒约10次高频调用API时,FastAPI仅能接收3-4条追踪记录。
已尝试方案
- 使用Redis缓存
- 异步处理
- 多线程
以上方案均未能解决数据丢失问题。
环境信息
部署在Kubernetes环境中,相关配置及代码如下:
OpenTelemetry Collector配置
apiVersion: opentelemetry.io/v1alpha1 kind: OpenTelemetryCollector metadata: name: otelcol spec: mode: daemonset config: | receivers: otlp: protocols: grpc: endpoint: 0.0.0.0:4317 http: endpoint: 0.0.0.0:4318 cors: allowed_origins: - "http://*" - "https://*" processors: memory_limiter: check_interval: 1s limit_percentage: 75 spike_limit_percentage: 15 batch: send_batch_size: 10000 timeout: 10s exporters: # NOTE: Prior to v0.86.0 use `logging` instead of `debug`. debug: otlp/jaeger/A: endpoint: "jaeger-collector.default.svc.cluster.local:14250" tls: insecure: true otlp/jaeger/B: endpoint: "jaeger-collector.default.svc.cluster.local:4317" tls: insecure: true otlp/python/C: endpoint: "10.20.1.230:8003" tls: insecure: true otlp/python/D: endpoint: "10.20.1.230:50051" tls: insecure: true service: pipelines: traces: receivers: [otlp] processors: [batch] exporters: [otlp/jaeger/A, otlp/jaeger/B, otlp/python/C, otlp/python/D, debug] metrics: receivers: [otlp] processors: [batch] exporters: [debug] logs: receivers: [otlp] processors: [batch] exporters: [debug]
Instrumentation配置
apiVersion: opentelemetry.io/v1alpha1 kind: Instrumentation metadata: name: instrumentation spec: exporter: endpoint: http://otelcol-collector.default.svc.cluster.local:4317 endpoint: http://otelcol-collector.default.svc.cluster.local:4318 propagators: - tracecontext - baggage - b3 java: image: ghcr.io/open-telemetry/opentelemetry-operator/autoinstrumentation-java:latest python: image: ghcr.io/open-telemetry/opentelemetry-operator/autoinstrumentation-python:latest env: - name: OTEL_LOGS_EXPORTER value: otlp_proto_http - name: OTEL_PYTHON_LOGGING_AUTO_INSTRUMENTATION_ENABLED value: 'true' - name: OTEL_TRACES_EXPORTER value: otlp_proto_http - name: OTEL_EXPORTER_OTLP_TRACES_ENDPOINT value: http://10.20.1.230:8004/v1/traces sampler: type: parentbased_traceidratio argument: "1"
FastAPI代码
import asyncio from fastapi import FastAPI, Request from opentelemetry.proto.collector.trace.v1 import trace_service_pb2 from google.protobuf.json_format import MessageToDict import aioredis import time from contextlib import asynccontextmanager import os REDIS_URL = "redis://127.0.0.1:6379" BATCH_SIZE = 100 BATCH_TIMEOUT = 1.0 # 1 second request_queue = asyncio.Queue() async def get_redis_connection(): return await aioredis.from_url(REDIS_URL) @asynccontextmanager async def lifespan(app: FastAPI): process_task = asyncio.create_task(process_batch()) yield process_task.cancel() try: await process_task except asyncio.CancelledError: pass app = FastAPI(lifespan=lifespan) async def update_stats(redis, worker_id): pipe = redis.pipeline() pipe.incr(f"worker:{worker_id}:count") pipe.incr("total_requests") results = await pipe.execute() return results[0], results[1] # worker_count, total_count @app.post("/v1/traces") async def receive_traces(request: Request): content = await request.body() worker_id = os.getpid() async with await get_redis_connection() as redis: worker_count, total_count = await update_stats(redis, worker_id) print({ "status": "received", "worker_id": worker_id, "worker_count": worker_count, "total_count": total_count }) print("total_count", total_count) await request_queue.put(content) return {"status": "received"} async def process_batch(): try: async with await get_redis_connection() as redis_conn: while True: batch = [] start_time = time.time() while len(batch) < BATCH_SIZE and time.time() - start_time < BATCH_TIMEOUT: try: content = await asyncio.wait_for(request_queue.get(), timeout=0.1) batch.append(content) except asyncio.TimeoutError: continue if batch: trace_data_list = [] for content in batch: trace_data = trace_service_pb2.ExportTraceServiceRequest() trace_data.ParseFromString(content) # print("trace_data", trace_data) print("trace_data") trace_data_list.append(MessageToDict(trace_data)) await redis_conn.set("latest_trace", str(trace_data_list[-1])) except Exception as e: print("ERROR", e) @app.get("/stats") async def get_stats(): async with await get_redis_connection() as redis: worker_keys = await redis.keys("worker:*:count") pipe = redis.pipeline() for key in worker_keys: pipe.get(key) pipe.get("total_requests") results = await pipe.execute() stats = {key.split(":")[1]: int(count) for key, count in zip(worker_keys, results[:-1])} total_count = int(results[-1]) if results[-1] else 0 return { "worker_stats": stats, "total_requests": total_count } if __name__ == "__main__": import uvicorn uvicorn.run("fast_scaler:app", host="0.0.0.0", port=8004, workers=16)
可能的数据丢失原因
1. OpenTelemetry Collector批处理配置不合理
Collector的batch处理器设置了send_batch_size: 10000和timeout: 10s,意味着只有攒够10000条数据或等待10秒才会批量导出。而每秒仅10次调用,10秒仅生成100条数据,远未达到阈值,导致数据被缓存,无法实时发送到FastAPI。
2. Instrumentation配置存在冲突
exporter字段重复设置endpoint,后一个地址会覆盖前一个,可能导致部分追踪数据无法正确上报到Collector。
3. FastAPI多Worker队列隔离问题
启动了16个Uvicorn Worker,每个Worker拥有独立的内存request_queue,但process_batch任务仅在主Worker运行,其他Worker的队列数据无法被处理,造成数据“丢失”。
4. 缺少队列溢出与超时处理
当内存队列满时,await request_queue.put(content)会阻塞请求,导致Collector超时重试或直接丢弃数据,且无任何错误提示。
解决措施
1. 调整Collector批处理参数
修改batch处理器配置,降低批量大小和超时时间,确保数据及时导出:
batch: send_batch_size: 100 timeout: 1s
2. 修复Instrumentation配置冲突
删除重复的endpoint字段,保留一个正确的Collector地址:
exporter: endpoint: http://otelcol-collector.default.svc.cluster.local:4318
3. 替换为分布式Redis队列
用Redis队列替代内存队列,实现多Worker数据共享:
# 替换内存队列为Redis队列操作 async def enqueue_trace(content): async with await get_redis_connection() as redis: await redis.lpush("trace_queue", content) # 设置队列最大长度,避免内存溢出 await redis.ltrim("trace_queue", 0, 9999) async def dequeue_trace(): async with await get_redis_connection() as redis: return await redis.rpop("trace_queue") # 修改receive_traces接口的入队逻辑 await enqueue_trace(content) # 修改process_batch的出队逻辑 while len(batch) < BATCH_SIZE and time.time() - start_time < BATCH_TIMEOUT: content = await dequeue_trace() if content: batch.append(content) else: await asyncio.sleep(0.01)
4. 添加请求超时与错误处理
在FastAPI接口中加入入队超时处理,避免请求阻塞:
from asyncio import timeout @app.post("/v1/traces") async def receive_traces(request: Request): content = await request.body() worker_id = os.getpid() async with await get_redis_connection() as redis: worker_count, total_count = await update_stats(redis, worker_id) print({ "status": "received", "worker_id": worker_id, "worker_count": worker_count, "total_count": total_count }) # 添加入队超时处理 try: async with timeout(0.5): await enqueue_trace(content) except TimeoutError: print("Queue is full, dropping trace data") return {"status": "received"}
5. 验证Collector Debug输出
通过debug exporter查看Collector是否接收并导出了所有追踪数据,排除Collector端的数据丢失问题。
内容的提问来源于stack exchange,提问作者Knut

