【发布时间】:2019-11-21 04:27:30
【问题描述】:
我使用 Rabbit MQ 按顺序处理消息。现在我有一个场景,不需要立即消费消息,但我想每 1(x) 分钟从 MQ 消费/读取消息,批量大小为 20(y)
所以我可以一次处理这 20 条消息并在一次调用中保存到数据库,而不是为每条消息调用 20 次。
那么如何在每 x 间隔批量接收/消费消息。
我看到了以下信息,Consume messages in batches - RabbitMQ 和 Rabbitmq retrieve multiple messages using single synchronous call using .NET
我试图实现第二个问题的实现(questions/32309155)但没有工作,不明白“consumer.Received”会收到 **_fetchSize ==> 20 ** 的意思,它会在单次读取中收到 20 条消息?或者它将如何工作,因为我将 fetchsize 更改为 10,但 consumer.received 正在接收单个消息。
using (var connection = factory.CreateConnection())
{
using (var channel = connection.CreateModel())
{
channel.BasicQos(0, 1, false);
channel.ExchangeDeclare("helloExchange", type:"direct");
channel.QueueDeclare(queue: "hello", durable: true, exclusive: false, autoDelete: false,
arguments: null);
channel.QueueBind("hello", "helloExchange", routingKey:"hello");
var consumer = new EventingBasicConsumer(channel);
consumer.Received += (model, ea) =>
{
bool canAck = false;
var retryCount = 0;
try
{
var body = ea.Body;
var message = Encoding.UTF8.GetString(body);
// DO PROCESS MESSAGE HERE
Console.WriteLine($"{typeof(MyConsumer).Name} Message consumed {message}");
canAck = true;
}
catch (Exception ex)
{canAck = false;
// LOG ERROR
}
try
{
if (canAck)
{
channel.BasicAck(ea.DeliveryTag, false);
}
else
{
channel.BasicNack(ea.DeliveryTag, false, false);
}
}
catch (AlreadyClosedException ex)
{
Console.WriteLine(ex.Message + " >> RabbitMQ is closed!");
}
};
channel.BasicConsume(queue: "hello", autoAck: false, consumer: consumer);
Console.WriteLine(" Press [enter] to exit.");
Console.ReadLine();
}
}
【问题讨论】: