You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.19 02:17:35