如何预热Celery worker 实现首个任务处理前预加载数据
Celery gevent模式下预加载启动数据的正确实现
问题根因
你之前的写法跑不通是两个核心原因:
- 预加载执行时机过早:
warmup_data()运行时,gevent的猴子补丁还未生效,Celery的协程运行时也没初始化完成。这时候读磁盘、初始化GPU上下文生成的资源句柄、内存对象,和后续启动的gevent协程池完全不兼容,很容易出现协程死锁——也就是你观察到的“任务被正常接收、逻辑执行完成但永远不返回结果”,这类协程阻塞不会抛出异常,属于静默故障,排查难度很高。 - 误判Celery的钩子能力:Celery自带的生命周期信号完全可以实现Gunicorn pre/post fork类的自定义逻辑注入,不需要在入口文件硬写加载逻辑。
正确实现代码
首先注意:gevent的猴子补丁必须放在整个项目所有import语句的最前面,不然所有涉及IO、线程的逻辑都会出现不可预期的问题。
# !!!这几行必须放在文件最顶部,比所有其他业务、第三方库import都靠前 from gevent import monkey monkey.patch_all() import logging from celery import Celery, signals logger = logging.getLogger(__name__) # 初始化Celery app实例 app = Celery( "your_service", broker="替换为实际的broker地址", backend="替换为实际的结果后端地址" ) # 全局变量存储预加载资源,所有工作协程共享 WARMUP_DATA = None def warmup_data(): """磁盘数据加载、GPU上下文初始化逻辑""" global WARMUP_DATA logger.info("Start preloading large data from disk...") # 替换为实际的加载逻辑:读大文件、加载模型到GPU、构建缓存索引均可 WARMUP_DATA = load_your_own_data_and_gpu_ctx() logger.info("Preload finished, data is ready for tasks.") # 注册worker初始化钩子:该信号会在worker完成所有内部初始化、正式开始消费任务之前触发 # 无论使用gevent/eventlet/solo/prefork哪种运行池,该信号的触发时机都符合预加载要求 @signals.worker_init.connect def run_warmup_before_task_consume(**kwargs): warmup_data() # 业务任务示例 @app.task(acks_late=True) def process_your_task(task_input): # 直接调用全局预加载的数据即可,不需要重复加载 task_result = do_calculation(WARMUP_DATA, task_input) return task_result if __name__ == "__main__": # 入口处不要重复写warmup逻辑 app.worker_main( argv=[ "worker", "--loglevel=INFO", "--pool=gevent", "--concurrency=4", ] )
避坑说明
- 不要用错信号:
worker_process_init是专门为prefork进程池设计的,只有fork出子进程时才会触发,gevent/eventlet是单进程多协程模式,绑定这个信号不会执行任何逻辑。 - 协程安全注意:gevent模式下所有协程共享进程内存,只要任务对预加载的
WARMUP_DATA是只读操作,就不会有并发安全问题,不需要加锁,也没有多进程模式下的内存拷贝开销,GPU资源可以正常复用。 - 正确性验证方式:启动worker后观察日志顺序,必须是先打印预加载开始、预加载完成的日志,再出现
celery@xxx ready的启动成功日志,才说明预加载时机正确,此时发送任务即可正常返回结果。 - 如果预加载耗时很长(比如加载几十GB模型),可以在
warmup_data中添加分段进度日志,方便排查启动慢的问题。
内容的提问来源于stack exchange,提问作者Leopd
相关产品推荐
相关产品推荐

