【问题标题】:Different queues in celery芹菜中的不同队列
【发布时间】:2023-01-05 13:30:51
【问题描述】:

我有一个项目,我在其中使用文件 (python main.py) 启动我的 FastAPI:

import uvicorn
from configuration import API_HOST, API_PORT

if __name__ == "__main__":
    uvicorn.run("endpoints:app", host="localhost", port=8811, reload=True, access_log=False)

在 endpoints.py 里面我有:

from celery import Celery
from fastapi import FastAPI
import os
import time

# Create object for fastAPI
app = FastAPI(
    title="MYFASTAPI",
    description="MYDESCRIPTION",
    version=1.0,
    contact="ME!",
)

celery = Celery(__name__)
celery.conf.broker_url = os.environ.get("CELERY_BROKER_URL", "redis://localhost:6379")
celery.conf.result_backend = os.environ.get("CELERY_RESULT_BACKEND", "redis://localhost:6379")
celery.conf.task_track_started = True
celery.conf.task_serializer = pickle
celery.conf.result_serializer = pickle
celery.conf.accept_content = ["pickle"]

# By defaul celery can handle as many threads as CPU cores have the instance. 
celery.conf.worker_concurrency = os.cpu_count()

# Start the celery worker. I start it in a separate thread, so fastapi can run in parallel
worker = celery.Worker()

def start_worker():
    worker.start()

ce = threading.Thread(target=start_worker)
ce.start()

@app.post("/taskA")
def taskA():
    task = ask_taskA.delay()
    return {"task_id": task.id}

@celery.task(name="ask_taskA", bind=True)
def ask_taskA(self):
    time.sleep(100)

@app.post("/get_results")
def get_results(task_id):
    task_result = celery.AsyncResult(task_id)
    return {'task_status': task_result.status}

鉴于此代码,我如何拥有两个不同的队列,为每个搜索队列分配特定数量的工作人员并将特定任务分配给这些队列之一?

我读到人们习惯将芹菜执行为:

celery -A proj worker

但是由于一些进口,项目中有一个结构限制了我,最后我通过在不同的线程中启动 celery worker 来完成(工作完美)

【问题讨论】:

    标签: python redis queue celery fastapi


    【解决方案1】:

    根据官方芹菜文档https://docs.celeryq.dev/en/stable/userguide/routing.html#manual-routing[1],您可以按照此指定不同的队列。

    from kombu import Queue
    
    app.conf.task_default_queue = 'default'
    app.conf.task_queues = (
        Queue('default',    routing_key='task.#'),
        Queue('feed_tasks', routing_key='feed.#'),
    )
    app.conf.task_default_exchange = 'tasks'
    app.conf.task_default_exchange_type = 'topic'
    app.conf.task_default_routing_key = 'task.default'
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2014-06-09
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-12-29
      • 2017-12-20
      相关资源
      最近更新 更多