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
相关产品推荐
相关产品推荐

