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

Psycopg在Celery应用上下文停止反序列化,类型注册问题排查

问题

我用psycopg的ConnectionPool扩展Flask应用,初始化时通过register_composite_types注册PostgreSQL自定义复合类型。用注入Flask上下文的自定义FlaskTask类配置Celery,shared_task在Flask应用内运行正常,但在Celery Worker中执行时抛出错误:AttributeError: 'str' object has no attribute 'field_idx'。这个错误一般是PostgreSQL自定义类型没在当前作用域注册导致的,请问我是不是没在相关范围内正确注册这些类型?

相关代码

初始化连接池代码

def init_pool(cfg, connection_class, name) -> ConnectionPool:

    pool = ConnectionPool(
        conninfo=f"postgresql://{cfg['user']}:{cfg['password']}@{cfg['host']}/{cfg['database']}",
        min_size=cfg['min_pool_size'],
        max_size=cfg['max_pool_size'],
        connection_class=connection_class,
        kwargs={"autocommit": True, "options": "-c idle_session_timeout=0"},
        name=name
    )
    pool.wait(timeout=30.0)

    with pool.connection() as conn:
        # retrieve a connection to register the user-defined types
        # in the global scope of the application (*default*)
        register_composite_types(conn)


    return pool

def register_composite_types(conn):
    """ Convert postgresql type -> python
        * Register the shared composite types in the conn scope.
        * Utilizes the built-in psycopg factory call
    """
    with conn.cursor() as cur:
        for (t,) in cur.execute("select shared_types()"):
            if not conn.broken:
                info = CompositeInfo.fetch(conn, t)
                register_composite(info, context=None)
            else:
                raise LevelsDbException(project_id=None, filename=None,
                                        message="Failed to initialize: Broken connection")

Celery初始化代码

def celery_init(app: Flask) -> Celery:

    class FlaskTask(Task):
        def __call__(self, *args: object, **kwargs: object) -> object:
            with app.app_context():
                return self.run(*args, **kwargs)

    # instantiated once for each task.  Each task serves multiple task requests.
    celery_app = Celery(app.name, task_cls=FlaskTask)
    celery_app.config_from_object(celeryconfig)
    celery_app.set_default()
    app.extensions["celery"] = celery_app

    return celery_app

任务代码

with current_app.inspection_db_pool.connection() as conn:
    field_data = db.get_file_fields( conn
                                   , project_id
                                   , path
                                   , hashstr
                                   )

# field_data (a Cursor returned from `conn.execute`)
fields: list[dict] = fields_from_levels_db(field_data)

# ... 
def fields_from_levels_db(file_fields: Iterable) -> list[dict]:
    """ file_fields is a Cursor """
    return [ dict(
        idx = field.field_idx,
        purpose = field.purpose,
        ...
        levels = field.levels

    ) for (field,) in file_fields ]
分析与解决

你当前的问题核心是Celery Worker进程没有注册PostgreSQL自定义复合类型,原因如下:

  • 连接池初始化时的类型注册仅在Flask主进程完成,Celery Worker是独立的进程,不会共享主进程的类型注册上下文。
  • register_composite(info, context=None)中的context=None是将类型注册到当前进程的全局上下文,但Worker进程启动时没有执行这个注册逻辑。

以下是两种可行的解决方法:

方法1:Worker启动时自动注册类型

修改Celery初始化逻辑,通过Worker启动信号触发类型注册,确保每个Worker进程启动时都完成自定义类型的注册:

from celery.signals import worker_init

def celery_init(app: Flask) -> Celery:
    class FlaskTask(Task):
        def __call__(self, *args: object, **kwargs: object) -> object:
            with app.app_context():
                return self.run(*args, **kwargs)

    celery_app = Celery(app.name, task_cls=FlaskTask)
    celery_app.config_from_object(celeryconfig)
    celery_app.set_default()
    app.extensions["celery"] = celery_app

    # Worker启动时执行类型注册
    @worker_init.connect(sender=celery_app)
    def init_worker(sender=None, **kwargs):
        cfg = app.config['DATABASE_CONFIG']
        # 建立临时连接完成类型注册,无需加入连接池
        with psycopg.connect(
            f"postgresql://{cfg['user']}:{cfg['password']}@{cfg['host']}/{cfg['database']}",
            autocommit=True,
            options="-c idle_session_timeout=0"
        ) as conn:
            register_composite_types(conn)

    return celery_app

方法2:获取连接时检查并注册

修改任务中获取连接的逻辑,每次从连接池拿连接时执行类型注册(重复注册不会产生问题):

# 移除init_pool函数中初始化时的注册代码,修改后的init_pool:
def init_pool(cfg, connection_class, name) -> ConnectionPool:
    pool = ConnectionPool(
        conninfo=f"postgresql://{cfg['user']}:{cfg['password']}@{cfg['host']}/{cfg['database']}",
        min_size=cfg['min_pool_size'],
        max_size=cfg['max_pool_size'],
        connection_class=connection_class,
        kwargs={"autocommit": True, "options": "-c idle_session_timeout=0"},
        name=name
    )
    pool.wait(timeout=30.0)
    # 这里不再执行register_composite_types
    return pool

# 任务中获取连接时添加注册逻辑
with current_app.inspection_db_pool.connection() as conn:
    # 执行类型注册,重复注册不影响
    register_composite_types(conn)
    field_data = db.get_file_fields(conn, project_id, path, hashstr)

关键说明

  • psycopg的复合类型注册是进程级的,每个独立进程(包括Celery Worker)都需要单独完成注册。
  • register_composite传入context=None时,会将类型注册到进程全局上下文,该进程内所有后续连接都会生效。
  • 重复注册同一复合类型不会报错,因此可以安全地在Worker启动或每次获取连接时执行注册操作。

内容的提问来源于stack exchange,提问作者Edmund's Echo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 12:59:52