how to schedule tasks from a RabbitMQ queue using celery beat?

Viewed 143

I'm publishing a Message to "RequestReceiverQueue" RMQ which I need to schedule as task.So here message being published is variable so how that can be set as a scheduled task using celery beat or alternative if any?

app.py

scheduler_app = Celery(
    "schedulerApp",
    backend="rpc://",
    broker="pyamqp://guest:guest@localhost:5672/poc_vhost",
    include=["workers.tasks"]
)

request_receiver_exchange = Exchange("request_receiver_exchange", type="topic")

scheduler_app.conf.task_queues = (
    Queue(
        'RequestReceiverQueue',
        exchange=request_receiver_exchange,
        durable=True,
        routing_key="workers.tasks.*"
    ),
)

workers.tasks.py

from workers.app import scheduler_app

@scheduler_app.task
def execute_request_1(msg):
    print("You are in execute_request_1: ")
    print("processing message: ", msg)


@scheduler_app.task
def execute_request_2(request):
    print("You are in execute_request_2")
    print("Received request is: ", request)

celery beat configs: here I need to schedule the incoming message from "RequestReceiverQueue" queue and not the kwargs defined in below scheduler config.

scheduler_app.conf.beat_schedule = {
    "task-scheduler-1": {
        "task": "workers.tasks.execute_request_1",
        "schedule": 60.0,
        "kwargs": {"Msg": "Between 1945 and 1947, Turing lived in Hampton, London."}
    },
    "task-scheduler-2": {
        "task": "workers.tasks.execute_request_2",
        "schedule": 60.0,
        "kwargs": {"Msg": "Turing worked at the National Physical Laboratory."}
    }
}
0 Answers
Related