Celery publisher doesn't send messages to the intended rabbitmq queue

Viewed 181

I have a flask app running with celery in it. When a user inputs a message into app uri it needed to be sent to a rabbitmq queue as a message to a given queue. But it send messages to a default queue. How do I specifically send a message to a defined queue?

publisher.py

from flask import Flask,request,jsonify
from celery import Celery,bootsteps
from time import sleep
from kombu.common import QoS

class NoChannelGlobalQoS(bootsteps.StartStopStep):
    requires = {'celery.worker.consumer.tasks:Tasks'}

    def start(self, c):
        qos_global = False

        c.connection.default_channel.basic_qos(0, c.initial_prefetch_count, qos_global)
        def set_prefetch_count(prefetch_count):
            return c.task_consumer.qos(
                prefetch_count=prefetch_count,
                apply_global=qos_global,
            )
        c.qos = QoS(set_prefetch_count, c.initial_prefetch_count)

server = Flask(__name__)
broker_uri1="amqp://"
backend="mongodb+srv://"

app = Celery('TestApp', broker_uri=broker_uri1,backend=backend)
CELERY_TASK_ROUTES = {'app.tasks.*': {'queue': 'rabbit'},'app.task.*': {'queue': 'rabbit'}}

app.config_from_object('celeryconfig')
app.conf.task_default_exchange='rabbit'
app.conf.task_default_routing_key='rabbit'

app.steps['worker'].add(NoChannelGlobalQoS)

@server.route("/")
def print_hello():
    return "Flask server working"

@server.route("/sync-reverse")
def reverse_by_worker_with_results():
    input = request.args.get("text")
    app.select_queues(queues="rabbit")
    task = app.signature('tasks.reverse', kwargs={'text': input})
    result = task.delay()

    while result.status == 'PENDING':
        print(result.status)
        sleep(2)

    print(result.status)

    return "Reversed >><br />Input Text : " + input + "<br />Output Text : " + result.get()

@server.route("/async-reverse")
def reverse_by_worker():
    input = request.args.get("text")
    app.select_queues(queues="rabbit")     #supposed to add queues to the task
    task = app.signature('tasks.reverse', kwargs={'text': input})
    result = task.delay()

    if result.id:
        return 'Task add'
    else:
        return 'Failure in adding task'

# Press the green button in the gutter to run the script.
if __name__ == '__main__':
    server.run(host= '0.0.0.0')

worker.py

from celery import Celery,bootsteps
from time import sleep
from kombu.common import QoS

broker_uri='amqp:///'
backend_uri="mongodb+srv://"

app = Celery('TestApp', broker=broker_uri, backend=backend_uri)
app.config_from_object('celeryconfig')
app.conf.task_default_exchange='rabbit'
app.conf.task_default_routing_key='rabbit'

class NoChannelGlobalQoS(bootsteps.StartStopStep):
    requires = {'celery.worker.consumer.tasks:Tasks'}

    def start(self, c):
        qos_global = False

        c.connection.default_channel.basic_qos(0, c.initial_prefetch_count, qos_global)
        def set_prefetch_count(prefetch_count):
            return c.task_consumer.qos(
                prefetch_count=prefetch_count,
                apply_global=qos_global,
            )
        c.qos = QoS(set_prefetch_count, c.initial_prefetch_count)

app.steps['consumer'].add(NoChannelGlobalQoS)

@app.task
def reverse(text):
    sleep(10)
    return text[:-1]

I have used a celery config file to define the queue name and type as suggested

from kombu import Queue

task_queues = [Queue(name="rabbit", queue_arguments={"x-queue-type": "quorum"})]

task_routes = {
    'tasks.add': 'rabbit',
}
0 Answers
Related