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?