pika publish in celery task func

Viewed 158

in some reasons, I got some code using pika to publish message in a celery task. code like this:

@app.task
async def test_celery():
    with Rabbitmq_Helper() as connection:
        channel = connection.channel()
        channel.basic_qos(
            prefetch_size=0, prefetch_count=0
    )
        exchange_name = 'test'
        channel.exchange_declare(exchange=exchange_name, exchange_type='topic', 
durable=True)

    queue_name = "testtest"

    properties = pika.BasicProperties(delivery_mode=2)
    # body = msgpack.packb(clickhouse_utils.format_data(data, nodes), default=str)
    for i in range(100):
        body = "test data"
        print("send msg")
        channel.basic_publish(
            exchange=exchange_name, routing_key=queue_name, body=body, properties=properties
        )

In the formal scene, the celery task will delay many times. And a consumer func will watch the rabbitmq queue.

But in the test, I found the consumer can't receive the msg in the first time. And the celery can publish the msg in queue.

So I try to debug, I take the test file like this. But I found the celery task delay is right,but it can't publish the msg in queue.

what's happend?

0 Answers
Related