如何用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
相关产品推荐
相关产品推荐

