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

Celery+RabbitMQ RPC调用结果无法返回问题求助

问题背景

使用Celery 5.2.7 + RabbitMQ 3.9-management实现RPC服务,拆分至worker、rabbitmq、发起RPC调用的API三个独立容器,通过单份docker-compose部署在bridge网络中,整体运行于Ubuntu 20.04.3 LTS单台机器上。

工作流程

  • API接收用户请求
  • API向RabbitMQ提交GPU Worker任务并等待响应
  • GPU Worker监听任务队列,有新任务时执行计算
  • 计算完成后,GPU Worker将结果返回(推测通过另一个RabbitMQ队列)

核心问题

API始终无法收到返回结果。日志显示GPU Worker完成任务耗时<0.1秒,但API一直无法获取结果,最终触发超时。


配置信息

Worker代码

# celery_worker.py
app = Celery('celery_worker', broker= f'amqp://{rabbitmq_username}:{rabbitmq_password}@{rabbitmq_host}:{rabbitmq_port}/', backend='rpc://')

@app.task
def remote_procedure(a, b):
    # some heavy computations
    return res

API主代码

# API main.py
from proj.celery_worker import remote_procedure

RESPONSE_TIMEOUT_SECONDS = 1.5

@app.get('/search/') 
def search(a,b):
    try:
        res = remote_procedure.apply_async((a,b), expires = RESPONSE_TIMEOUT_SECONDS-0.3, retry=False).get(timeout=RESPONSE_TIMEOUT_SECONDS)
    except celery.exceptions.TimeoutError:
        print('timeout')
        return {"res":"timeout"}
    return {"res":res}


def main():
    """
    Main function
    """

    uvicorn.run(app, host='127.0.0.1', port=CONFIG['port'])


if __name__ == '__main__':
    main()

Docker Compose配置

gpu_worker:
    build:
      context: .
      dockerfile: Dockerfile.worker
    deploy:
      resources:
        limits:
          memory: 10gb
        reservations:
          devices:
            - capabilities: [ gpu ]
    networks:
      - search_net
  
api:
    build:
      context: .
      dockerfile: Dockerfile.api
    depends_on:
      - gpu_worker
    ports:
      - "8001:8000"
    environment:
      MAX_WORKERS: 32
    deploy:
      resources:
        limits:
          memory: 10gb
        reservations:
          devices:
            - capabilities: [ gpu ]
    networks:
      - search_net

rabbitmq:
    image: rabbitmq:3.9-management
    hostname: rabbitmq
    container_name: 'rabbitmq'
    ports:
      - 5672:5672
      - 15672:15672
    volumes:
      - rabbitmq_data:/var/lib/rabbitmq/
      - rabbitmq_log:/var/log/rabbitmq/
    restart: on-failure
    networks:
      - search_net

volumes:
  rabbitmq_data:
  rabbitmq_log:

networks:
  search_net:
    driver: bridge

Dockerfile.worker

...
ENTRYPOINT [ "celery", "-A", "proj.celery_worker", "worker", "--loglevel=DEBUG", "--pool=solo", "-c", "1"]

Dockerfile.api

...
EXPOSE 80

# Start the app
CMD ["uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000"]

排查方向与解决方案

1. 容器间网络连通性验证

Celery的rpc://后端依赖RabbitMQ传递结果,必须确保API和Worker能正确访问RabbitMQ容器:

  • 检查API和Worker代码中的rabbitmq_host是否设置为docker-compose中RabbitMQ的服务名rabbitmq(不能用localhost或127.0.0.1,否则会指向容器内部而非RabbitMQ容器)。
  • 进入API容器执行ping rabbitmq,确认网络连通性;若不通,检查bridge网络配置是否正常。

2. API服务绑定地址冲突修复

API主函数中uvicorn绑定了127.0.0.1,会覆盖Dockerfile中0.0.0.0的配置,导致Celery实例的网络监听范围受限。修改main函数中的启动代码:

uvicorn.run(app, host='0.0.0.0', port=CONFIG['port'])

3. RabbitMQ就绪状态依赖配置

当前depends_on仅保证Worker启动,但RabbitMQ可能未完全就绪就启动API,导致Celery初始化失败。给RabbitMQ添加健康检查,并调整API的依赖条件:

# 修改rabbitmq服务配置
rabbitmq:
    # ... 原有配置
    healthcheck:
      test: ["CMD", "rabbitmq-diagnostics", "-q", "ping"]
      interval: 30s
      timeout: 30s
      retries: 3

# 修改api服务的depends_on
api:
    # ... 原有配置
    depends_on:
      gpu_worker:
        condition: service_started
      rabbitmq:
        condition: service_healthy

4. RPC后端队列权限与状态检查

  • 登录RabbitMQ管理界面(http://主机IP:15672),查看是否有临时结果队列被创建,以及Worker完成任务后是否有消息发送到该队列。
  • 检查RabbitMQ日志,确认是否存在队列创建失败、权限不足的错误。

5. 任务过期时间调整

当前任务过期时间设置为1.2秒,虽短于超时时间,但网络延迟可能导致任务未被Worker接收就过期。可临时去掉expires参数测试,或延长至与超时时间一致。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 12:25:23