Django Celery Prefork工作进程下OpenTelemetry指标异常的解决方案求助
我最近在给我的Django应用做OpenTelemetry的链路追踪和指标埋点,写了一个otel_config.py放在manage.py同级目录,内容如下:
import socket import os from urllib.parse import urljoin from opentelemetry import trace, metrics from opentelemetry.sdk.trace import TracerProvider from opentelemetry.sdk.trace.export import BatchSpanProcessor from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter from opentelemetry.sdk.metrics import MeterProvider from opentelemetry.sdk.metrics.export import PeriodicExportingMetricReader from opentelemetry.exporter.otlp.proto.http.metric_exporter import OTLPMetricExporter from opentelemetry.sdk.resources import Resource from opentelemetry.semconv.resource import SERVICE_NAME, SERVICE_INSTANCE_ID from opentelemetry.instrumentation.django import DjangoInstrumentor from opentelemetry.instrumentation.psycopg2 import Psycopg2Instrumentor from opentelemetry.instrumentation.celery import CeleryInstrumentor # 资源配置 def get_default_service_instance_id(): try: hostname = socket.gethostname() or "unknown-host" except Exception as e: hostname = "unknown-host" try: process_id = os.getpid() except Exception as e: process_id = "unknown-pid" return f"{hostname}-{process_id}" service_name = "my-service" otlp_endpoint = "http://otel-collector:4318" service_instance_id = get_default_service_instance_id() resource = Resource.create( { SERVICE_NAME: service_name, SERVICE_INSTANCE_ID: service_instance_id, } ) # 链路追踪配置 otlp_endpoint_traces = urljoin(otlp_endpoint, "/v1/traces") trace_exporter = OTLPSpanExporter(endpoint=otlp_endpoint_traces) span_processor = BatchSpanProcessor(trace_exporter) tracer_provider = TracerProvider(resource=resource) trace.set_tracer_provider(tracer_provider) trace.get_tracer_provider().add_span_processor(span_processor) # 指标配置 otlp_endpoint_metrics = urljoin(otlp_endpoint, "/v1/metrics") metric_exporter = OTLPMetricExporter(endpoint=otlp_endpoint_metrics) metric_reader = PeriodicExportingMetricReader(metric_exporter) meter_provider = MeterProvider(resource=resource, metric_readers=[metric_reader]) metrics.set_meter_provider(meter_provider) # 自动埋点 DjangoInstrumentor().instrument() Psycopg2Instrumentor().instrument() CeleryInstrumentor().instrument()
然后我在settings.py的末尾导入了这个配置:
import otel_config
大部分场景下链路和指标都正常,但在Celery用prefork模式运行时,指标出了问题:子进程会继承父进程的SERVICE_INSTANCE_ID,导致不同子进程的指标被当成同一个实例的,而每个子进程内存独立,指标值在Collector里频繁波动,完全不是所有子进程的聚合值。但用--pool=threads线程池模式就没问题,因为所有线程共享同一个指标实例,最终Collector里的聚合值是正确的。
我尝试了两种解决方案,但都有明显缺陷:
方案一:通过Celery信号分场景导入配置
我把settings.py里的import otel_config删掉,改成在Celery的信号处理器里导入:
from celery.signals import worker_init, worker_process_init @worker_init def worker_init_handler(sender, **kwargs): pool_cls = str(sender.pool_cls) if hasattr(sender, 'pool_cls') else None if "threads" in pool_cls: # 线程池模式下初始化 import otel_config @worker_process_init.connect def worker_process_init_handler(**kwargs): # Prefork子进程初始化时导入 import otel_config
同时为了Django API的埋点正常,我还在wsgi.py里加了import otel_config。
但这个方案的问题是:父进程的指标丢了。比如我有个task_received信号处理器会上报指标,这个信号只在父进程触发,现在父进程没有初始化OpenTelemetry,这部分指标就完全没了。而且我也不确定在多个地方初始化OpenTelemetry是不是好的实践,之前都是只在settings.py里导入一次的。
方案二:尝试在子进程中重新初始化OpenTelemetry
我想在worker_process_init信号里重新初始化配置,但发现OpenTelemetry的Python SDK是单例模式——MeterProvider、TracerProvider还有各种Instrumentor的instrument()方法一旦调用过,后续再设置就不会生效了。就算我强行修改内部的单例变量:
@worker_process_init.connect def worker_process_init_handler(**kwargs): # 省略重复的资源、追踪、指标配置代码... # 强行覆盖单例(很hack的写法) # trace.set_tracer_provider(tracer_provider) # 单例模式下这个方法不生效 trace._TRACER_PROVIDER = tracer_provider # metrics.set_meter_provider(meter_provider) # 同样不生效 metrics._internal._METER_PROVIDER = meter_provider
这种写法不仅非常不优雅,还需要修改代码里所有获取Meter的逻辑,完全不现实。
现在我想知道有没有人遇到过类似的问题,针对Celery Prefork模式下的OpenTelemetry指标异常,有没有更合理的最佳实践?
内容来源于stack exchange

