You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.02 19:41:05