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

如何用C#获取Azure Service Bus所有队列消息并存储为JSON

多队列消息获取并保存为JSON文件的改造方案

改造后完整代码

using System;
using System.Collections.Generic;
using System.Text;
using System.Threading.Tasks;
using Microsoft.Azure.ServiceBus;
using Newtonsoft.Json;
using System.IO;

namespace AzureServiceBusMultiQueue
{
    class Program
    {
        // 存储所有队列客户端,方便后续统一关闭
        private static readonly List<QueueClient> _queueClients = new List<QueueClient>();
        private const string SbConnectionString = "connection-string";
        // 替换为你的多个队列名称
        private static readonly List<string> _queueNames = new List<string> { "queue-name-1", "queue-name-2", "queue-name-3" };

        static void Main(string[] args)
        {
            try
            {
                foreach (var queueName in _queueNames)
                {
                    var queueClient = new QueueClient(SbConnectionString, queueName);
                    _queueClients.Add(queueClient);

                    var messageHandlerOptions = new MessageHandlerOptions(ExceptionReceivedHandler)
                    {
                        MaxConcurrentCalls = 1,
                        AutoComplete = false
                    };

                    // 注册消息处理器,同时传入队列名称用于区分存储文件
                    queueClient.RegisterMessageHandler((message, token) => ReceiveMessagesAsync(message, token, queueName), messageHandlerOptions);
                    Console.WriteLine($"已启动监听队列: {queueName}");
                }
            }
            catch (Exception ex)
            {
                Console.WriteLine($"初始化失败: {ex.Message}");
            }
            finally
            {
                Console.WriteLine("按任意键停止监听并退出...");
                Console.ReadKey();

                // 统一关闭所有队列客户端
                foreach (var client in _queueClients)
                {
                    if (client != null && client.IsClosedOrClosing == false)
                    {
                        client.CloseAsync().Wait();
                    }
                }
            }
        }

        static async Task ReceiveMessagesAsync(Message message, CancellationToken token, string queueName)
        {
            try
            {
                var messageContent = Encoding.UTF8.GetString(message.Body);
                Console.WriteLine($"从队列[{queueName}]收到消息: {messageContent}");

                // 将消息序列化为JSON格式(如果消息本身不是JSON,可根据需要调整结构)
                var messageData = new
                {
                    QueueName = queueName,
                    MessageId = message.MessageId,
                    Content = messageContent,
                    ReceivedTime = DateTime.UtcNow
                };
                var jsonString = JsonConvert.SerializeObject(messageData, Formatting.Indented);

                // 保存到对应队列的JSON文件,每次追加内容
                var filePath = $"{queueName}_messages.json";
                using (var writer = new StreamWriter(filePath, true, Encoding.UTF8))
                {
                    await writer.WriteLineAsync(jsonString);
                }

                // 标记消息为已处理
                await message.CompleteAsync();
            }
            catch (Exception ex)
            {
                Console.WriteLine($"处理队列[{queueName}]消息失败: {ex.Message}");
                // 可根据需求选择放弃消息或重新入队
                await message.AbandonAsync();
            }
        }

        static Task ExceptionReceivedHandler(ExceptionReceivedEventArgs exceptionReceivedEventArgs)
        {
            Console.WriteLine($"发生异常: {exceptionReceivedEventArgs.Exception.Message}");
            return Task.CompletedTask;
        }
    }
}

关键改造点说明

  • 多队列管理:用列表存储所有QueueClient实例,实现批量初始化和统一关闭,避免单个客户端的资源泄漏。
  • 消息与队列绑定:在注册消息处理器时传入队列名称,确保每个队列的消息能对应到专属的JSON文件。
  • JSON存储逻辑:将消息包装成包含队列名、消息ID、内容和接收时间的结构化对象,序列化为格式化后的JSON,以追加模式写入文件,保证消息不会被覆盖。
  • 异常处理增强:在消息处理方法内增加异常捕获,失败时可选择放弃消息或重新入队,避免单个消息异常影响整个队列的监听。

依赖说明

需要安装两个NuGet包:

  • Microsoft.Azure.ServiceBus:Azure Service Bus旧版SDK,与你现有代码兼容
  • Newtonsoft.Json:用于JSON序列化(也可替换为.NET内置的System.Text.Json)

内容的提问来源于stack exchange,提问作者Cwater

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 09:20:12