RabbitMQ/RabbitMQConsumer/FanoutSub.cs
(完整代码) 2KBnamespace RabbitMQConsumer
{
using CommonLib.RabbitMQ;
using System;
internal static class FanoutConsumer
{
private static void Main(string[] args)
{
using RabbitMQHelper mq = new(new string[] { "192.168.181.191" });
mq.UserName = "guest";
mq.Password = "guest";
mq.Port = 5672;
Console.WriteLine("input queueName...");
var input = Console.ReadLine();
switch (input)
{
case "1":
mq.Received += (result) =>
{
Console.WriteLine($"message:{result.Body}");
result.Commit();
};
mq.Listen("test.fanout.queue1", new ConsumeQueueOptions { AutoAck = false });
Console.ReadLine();
break;
case "2":
mq.Received += (result) =>
{
Console.WriteLine($"message:{result.Body}");
result.Commit();
};
mq.Listen("test.fanout.queue2", new ConsumeQueueOptions { AutoAck = false });
Console.ReadLine();
break;
}
#if rabbitMQClient
ConnectionFactory factory = BaseConsumer.CreateRabbitMqConnection();
using IConnection connection = factory.CreateConnection();
using IModel channel = connection.CreateModel();
EventingBasicConsumer consumer = new(channel);
channel.BasicQos(0, 1, false);
channel.BasicConsume(queue: "test.fanout.queue1", autoAck: false, consumer: consumer);
// 绑定消息接收后的事件委托
consumer.Received += (model, message) =>
{
Console.WriteLine($"Message:{Encoding.UTF8.GetString(message.Body.ToArray())}");
channel.BasicAck(
deliveryTag: message.DeliveryTag,
// 是否一次性确认多条数据
multiple: false);
};
#endif
}
}
}
24 小时内免费解锁 3 个项目,之后 1 积分/个。 规则说明
AI 解读
登录后可用,每次 10 积分,解读结果公开显示在下面。
还没有人解读过这个文件。
