如何用C#消费Node.js中Bull库生成的消息?
用C#消费Node.js Bull库生成的消息
Bull是基于Redis实现的队列库,只要在C#中连接同一Redis实例,按照Bull的队列键结构和消息格式处理,就能消费消息。以下是具体实现步骤:
1. 安装Redis客户端
使用NuGet安装StackExchange.Redis,这是C#生态中常用的Redis客户端:
dotnet add package StackExchange.Redis
2. 核心实现代码
以下代码实现了基础的队列监听和消息消费,同时模拟了Bull的消费确认机制(处理中、处理完成、处理失败的队列流转):
using StackExchange.Redis; using System.Text.Json; namespace BullQueueConsumer { class Program { static async Task Main(string[] args) { // 替换为你的Redis连接配置(需和Node.js Bull使用的Redis完全一致) var redisConn = ConnectionMultiplexer.Connect("localhost:6379,password=yourRedisPassword,defaultDatabase=0"); var redisDb = redisConn.GetDatabase(); // 要消费的Bull队列名称,必须和Node.js端定义的一致 string queueName = "user-notifications"; var queueKeys = new { Waiting = $"bull:{queueName}:waiting", Active = $"bull:{queueName}:active", Completed = $"bull:{queueName}:completed", Failed = $"bull:{queueName}:failed" }; Console.WriteLine($"开始监听Bull队列 [{queueName}]..."); while (true) { // 阻塞式从waiting队列尾部弹出消息(一直等待新消息) var rawMessage = await redisDb.ListRightPopAsync(queueKeys.Waiting, TimeSpan.FromSeconds(-1)); if (rawMessage.IsNull) continue; try { // 将消息移入active队列,标记为正在处理 await redisDb.ListLeftPushAsync(queueKeys.Active, rawMessage); // 反序列化Bull消息结构,提取业务数据 var bullMsg = JsonSerializer.Deserialize<BullMessage>(rawMessage); if (bullMsg != null) { Console.WriteLine($"[{DateTime.Now}] 收到消息ID: {bullMsg.Id}"); Console.WriteLine($"业务数据: {JsonSerializer.Serialize(bullMsg.Data)}"); // 这里编写你的业务处理逻辑 // HandleBusinessLogic(bullMsg.Data); } // 处理完成后,从active队列移除,移入completed队列 await redisDb.ListRemoveAsync(queueKeys.Active, rawMessage, count: 1); await redisDb.ListLeftPushAsync(queueKeys.Completed, rawMessage); } catch (Exception ex) { Console.WriteLine($"处理消息失败: {ex.Message}"); // 处理失败时,从active队列移除,移入failed队列 await redisDb.ListRemoveAsync(queueKeys.Active, rawMessage, count: 1); await redisDb.ListLeftPushAsync(queueKeys.Failed, rawMessage); } } } } // 匹配Bull存储的消息结构,可根据实际字段扩展 public class BullMessage { public string Id { get; set; } public object Data { get; set; } public int Timestamp { get; set; } public int Attempts { get; set; } public int Delay { get; set; } public string FailedReason { get; set; } } }
关键注意事项
- Redis一致性:必须保证C#端连接的Redis实例和Node.js Bull使用的完全一致(包括地址、端口、密码、数据库编号),否则无法读取队列消息。
- 延迟任务处理:如果你的队列包含延迟消息,需要监听
bull:<队列名>:delayed键(Bull用Redis有序集合存储延迟任务),定时通过ZRangeByScoreAsync获取到期消息,将其移到waiting队列。 - 消息结构兼容:Bull的消息字段可能随版本略有差异,可通过RedisInsight等工具查看实际存储的消息格式,调整
BullMessage类的字段。
内容的提问来源于stack exchange,提问作者Matheus Henrique
相关产品推荐
相关产品推荐

