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

使用opentelemetry-instrumentation-psycopg2与psycopg2连接池时的递归追踪问题

问题描述

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

复现方案

以下是复现所需的配置文件、依赖及代码:

  • 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 07:30:32