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

使用RabbitMQ和Celery实现Flask异步任务时任务接收未执行的问题

问题分析与解决方案

你的代码存在几个关键问题,导致Celery任务接收后无法执行:

核心问题

  1. 消费逻辑冲突:同时用pika手动监听RabbitMQ队列,又用Celery处理任务。Celery本身是基于消息队列的任务调度框架,自带成熟的消费逻辑,手动用pika消费会和Celery Worker抢消息,导致任务分发异常。
  2. Flask路由阻塞:channel.start_consuming()是阻塞调用,放在Flask路由里会导致该请求永久挂起,Flask服务无法处理其他请求,同时干扰Celery Worker的执行上下文。
  3. Celery配置与启动缺失:大概率没正确启动Celery Worker,或者Django发送消息的队列和Celery监听的队列不匹配,导致Worker收不到任务。

修复步骤

1. 重构代码,移除手动pika消费逻辑

让Celery直接负责任务的消费与处理,不需要自己编写pika消费代码。调整后的Flask+Celery示例代码:

from celery import Celery
from flask import Flask

app = Flask(__name__)
# 初始化Celery(单独配置,避免和Flask配置冲突)
cel = Celery(
    app.name,
    broker='amqp://guest:guest@localhost:5672/',
    backend='rpc://'
)
cel.conf.update(app.config)

# 定义Celery任务
@cel.task()
def process_task(data):
    print("Processing task with data:", data)
    import time
    time.sleep(5)
    print("Processing task completed:", data)
    return "TASK IS DONE"

# Flask仅负责提供API,不再处理队列消费
@app.route('/persons/')
def hello_world():
    return "Flask service is running"

if __name__ == '__main__':
    app.run(debug=True)

2. 调整Django端消息发送逻辑

让Django直接向Celery队列发送任务,而非自定义的test_micro队列。Django端代码示例:

from celery import Celery

# 配置Celery,和Flask端参数保持一致
cel = Celery(
    'django_tasks',
    broker='amqp://guest:guest@localhost:5672/',
    backend='rpc://'
)

# 在Django的POST视图中触发Celery任务
def post_handler(request):
    request_data = request.data.get('data')
    # 发送任务到Celery
    process_task.apply_async(args=(request_data,))
    return Response({"status": "task submitted"})

3. 正确启动Celery Worker

在Flask项目根目录下,执行以下命令启动Celery Worker:

celery -A 你的Flask模块名 worker --loglevel=info

比如你的Flask代码在app.py中,就替换为celery -A app worker --loglevel=info。

4. (可选)保留自定义队列的处理方式

如果Django必须发送到test_micro队列,需要让Celery监听该队列:

  • 在Celery配置中绑定队列:
cel.conf.task_routes = {
    '你的模块名.process_task': {'queue': 'test_micro'},
}
  • 启动Worker时指定监听该队列:
celery -A 你的Flask模块名 worker -Q test_micro --loglevel=info

关键注意事项

  • 确保RabbitMQ服务正常运行,本地5672端口可访问。
  • Flask服务和Celery Worker是两个独立进程,需要分别启动。
  • 禁止在Flask请求处理逻辑中加入阻塞操作(如pika的consume),会导致服务无法响应其他请求。

内容的提问来源于stack exchange,提问作者ek1669

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 15:32:43