使用opentelemetry-instrumentation-psycopg2与psycopg2连接池时的递归追踪问题
问题描述
当将opentelemetry-instrumentation-psycopg2库与psycopg2的ThreadedConnectionPool结合使用时,出现了递归追踪问题。该问题在并发场景下会概率性触发,以下是整理好的复现方案:

复现方案
以下是复现所需的配置文件、依赖及代码:
- otel-collector-config.yaml
receivers: otlp: protocols: grpc: http: exporters: logging: loglevel: debug jaeger: endpoint: jaeger-all-in-one:14250 tls: insecure: true processors: batch: service: pipelines: traces: receivers: [otlp] exporters: [logging, jaeger] processors: [batch]
- docker-compose.yaml
version: "3" services: jaeger-all-in-one: image: jaegertracing/all-in-one:1.42 restart: always environment: - COLLECTOR_OTLP_ENABLED=true ports: - "16686:16686" # server frontend - "14268:14268" # HTTP collector - "14250:14250" # gRPC collector otel-collector: image: otel/opentelemetry-collector:0.72.0 restart: always command: ["--config=/etc/otel-collector-config.yaml"] volumes: - ./otel-collector-config.yaml:/etc/otel-collector-config.yaml ports: - "4317:4317" # OTLP gRPC receiver - "4318:4318" # OTLP Http receiver depends_on: - jaeger-all-in-one # database postgres: image: postgres:13.2-alpine environment: - POSTGRES_USER=root - POSTGRES_PASSWORD=12345678 - POSTGRES_DB=example ports: - "5432:5432"
- requirements.txt
psycopg2==2.9.7 opentelemetry-instrumentation-psycopg2>=0.33b0 opentelemetry-exporter-otlp>=1.12.0
- index.py
import logging import threading import time from psycopg2 import pool import psycopg2 from opentelemetry import trace from opentelemetry.instrumentation.psycopg2 import Psycopg2Instrumentor from opentelemetry.sdk.resources import SERVICE_NAME, Resource from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter from opentelemetry.sdk.trace import TracerProvider from opentelemetry.sdk.trace.export import BatchSpanProcessor # 使用连接池的数据库类 class PostgresDatabasePool: def __init__(self, minconn, maxconn, host, database, user, password, port): self.db_pool = pool.ThreadedConnectionPool( minconn=minconn, maxconn=maxconn, host=host, database=database, user=user, password=password, port=port ) def execute_query(self, query): conn = None cursor = None try: conn = self.db_pool.getconn() cursor = conn.cursor() cursor.execute(query) logging.info(f"执行命令: {query}") except Exception as e: logging.error(f"SQL错误: {e}") finally: if cursor: cursor.close() if conn: self.db_pool.putconn(conn) # 使用单连接的数据库类 class PostgresDatabase: def __init__(self, host, database, user, password, port): self.conn = psycopg2.connect( host=host, database=database, user=user, password=password, port=port ) self.conn.autocommit = True self.cursor = self.conn.cursor() def execute_query(self, query): if not self.cursor: logging.warning("请先建立连接") return try: self.cursor.execute(query) logging.info(f"执行命令: {query}") except Exception as e: logging.error(f"SQL错误: {e}") def delete1_with_db_pool(thread_name): while True: with tracer.start_as_current_span("delete1_with_db_pool", kind=trace.SpanKind.INTERNAL): dbPool.execute_query("DELETE FROM public.test1;") dbPool.execute_query("DELETE FROM public.test1;") dbPool.execute_query("DELETE FROM public.test1;") dbPool.execute_query("DELETE FROM public.test1;") dbPool.execute_query("DELETE FROM public.test1;") time.sleep(5) def delete2_with_db_pool(thread_name): while True: with tracer.start_as_current_span("delete2_with_db_pool", kind=trace.SpanKind.INTERNAL): dbPool.execute_query("DELETE FROM public.test2;") dbPool.execute_query("DELETE FROM public.test2;") dbPool.execute_query("DELETE FROM public.test2;") dbPool.execute_query("DELETE FROM public.test2;") dbPool.execute_query("DELETE FROM public.test2;") time.sleep(5) def delete1(thread_name): while True: with tracer.start_as_current_span("delete1", kind=trace.SpanKind.INTERNAL): db.execute_query("DELETE FROM public.test1;") db.execute_query("DELETE FROM public.test1;") db.execute_query("DELETE FROM public.test1;") db.execute_query("DELETE FROM public.test1;") db.execute_query("DELETE FROM public.test1;") time.sleep(5) def delete2(thread_name): while True: with tracer.start_as_current_span("delete2", kind=trace.SpanKind.INTERNAL): db.execute_query("DELETE FROM public.test2;") db.execute_query("DELETE FROM public.test2;") db.execute_query("DELETE FROM public.test2;") db.execute_query("DELETE FROM public.test2;") db.execute_query("DELETE FROM public.test2;") time.sleep(5) # 追踪配置 resource = Resource(attributes={SERVICE_NAME: 'Demo_Bug'}) provider = TracerProvider(resource=resource) otlp_exporter = OTLPSpanExporter(endpoint='localhost:4317', insecure=True) provider.add_span_processor(BatchSpanProcessor(otlp_exporter)) trace.set_tracer_provider(provider) tracer = trace.get_tracer("test", "1.0.0") # 初始化Psycopg2链路追踪 Psycopg2Instrumentor().instrument() # 初始化数据库实例 dbPool = PostgresDatabasePool(1, 5, "localhost", "example", "root", "12345678", 5432) db = PostgresDatabase("localhost", "example", "root", "12345678", 5432) # 创建测试表 db.execute_query("CREATE TABLE IF NOT EXISTS public.test1 (id serial NOT NULL , text varchar(64) NOT NULL);") db.execute_query("CREATE TABLE IF NOT EXISTS public.test2 (id serial NOT NULL , text varchar(64) NOT NULL);") # 初始化线程 threads = [threading.Thread(target=delete1_with_db_pool, args=("Thread-1",)), threading.Thread(target=delete2_with_db_pool, args=("Thread-2",)), threading.Thread(target=delete1, args=("Thread-3",)), threading.Thread(target=delete2, args=("Thread-4",))] # 启动线程 for thread in threads: thread.start() # 等待线程结束 for thread in threads: thread.join()
执行命令:
docker-compose up -d pip install -r requirements.txt python index.py
打开Jaeger Dashboard即可观察到问题。
疑似根因
推测问题源于使用psycopg2的pool.ThreadedConnectionPool()复用数据库连接的场景,opentelemetry-instrumentation-psycopg2未适配该场景。已在opentelemetry-python-contrib仓库提交相关Issue,编号1925。
内容的提问来源于Stack Exchange,提问作者Changemyminds
相关产品推荐
相关产品推荐

