如何在C#中使用HttpClient将Azure Service Bus主题消息移入死信队列
使用HttpClient将Azure Service Bus消息移入死信队列(适配X++场景)
我正在开发一个供X调用的C#类库,用于Azure Service Bus通信,已实现从主题订阅读取消息的功能。但需要将处理失败、附带错误信息的消息移入死信队列,由于X不支持ServiceBusClient对应的NuGet包,只能通过HttpClient调用Azure Service Bus REST API实现。以下是修正后的可运行代码及关键说明:
修正后的完整代码
using System; using System.Net.Http.Headers; using System.Net.Http; using System.Text; using System.Security.Cryptography; using Newtonsoft.Json; namespace TestAzureServiceBusConsole { internal class Program2 { private const string ServiceBusNamespace = "<my-service-bus-instance>"; private const string TopicName = "<topic-name>"; private const string SubscriptionName = "<subscription-name>"; static string SasToken = ""; static string keyName = "<key-name>"; static string key = "<key>"; // API版本保持与Azure Service Bus兼容 private const string ApiVersion = "2021-12"; static void Main(string[] args) { var receiveMessagePath = $"{TopicName}/subscriptions/{SubscriptionName}/messages/head?api-version={ApiVersion}&lock-duration=PT5M"; var baseAddress = new Uri($"https://{ServiceBusNamespace}.servicebus.windows.net/"); using (var httpClient = new HttpClient { BaseAddress = baseAddress }) { // 生成SAS Token,资源URI指向订阅路径,权限更精准 var resourceUri = $"{baseAddress.AbsoluteUri}{TopicName}/subscriptions/{SubscriptionName}/"; SasToken = GetSasToken(resourceUri); httpClient.DefaultRequestHeaders.Authorization = new AuthenticationHeaderValue("SharedAccessSignature", SasToken); // 读取订阅中的消息(获取锁,有效期5分钟) var receiveResponse = httpClient.PostAsync(receiveMessagePath, null).GetAwaiter().GetResult(); if (!receiveResponse.IsSuccessStatusCode) { Console.WriteLine($"读取消息失败: {receiveResponse.StatusCode}"); Console.WriteLine($"错误详情: {receiveResponse.Content.ReadAsStringAsync().GetAwaiter().GetResult()}"); Console.ReadKey(); return; } // 解析BrokerProperties获取消息元数据 MessageRead currentMessage = null; if (receiveResponse.Headers.TryGetValues("BrokerProperties", out var brokerProperties)) { foreach (var propJson in brokerProperties) { currentMessage = JsonConvert.DeserializeObject<MessageRead>(propJson); currentMessage.MessageContent = receiveResponse.Content.ReadAsStringAsync().GetAwaiter().GetResult(); } } if (currentMessage == null || string.IsNullOrEmpty(currentMessage.LockToken)) { Console.WriteLine("无法获取消息的LockToken,无法执行死信操作"); Console.ReadKey(); return; } Console.WriteLine($"收到消息内容: {currentMessage.MessageContent}"); // 构造死信请求:使用LockToken定位消息,传入错误原因和描述 var deadLetterPath = $"{TopicName}/subscriptions/{SubscriptionName}/messages/{currentMessage.LockToken}/movetodeadletterqueue?api-version={ApiVersion}"; var deadLetterParams = new { DeadLetterReason = "处理失败", DeadLetterErrorDescription = "消息内容不符合业务规则,无法处理" }; var deadLetterContent = new StringContent(JsonConvert.SerializeObject(deadLetterParams), Encoding.UTF8, "application/json"); // 死信操作使用PUT请求 var deadLetterResponse = httpClient.PutAsync(deadLetterPath, deadLetterContent).GetAwaiter().GetResult(); if (deadLetterResponse.IsSuccessStatusCode) { Console.WriteLine("消息已成功移入死信队列"); } else { Console.WriteLine($"死信操作失败: {deadLetterResponse.StatusCode}"); Console.WriteLine($"错误详情: {deadLetterResponse.Content.ReadAsStringAsync().GetAwaiter().GetResult()}"); } } Console.ReadKey(); } public static string GetSasToken(string resourceUri) { string encodedUri = Uri.EscapeDataString(resourceUri).ToLower(); DateTime now = DateTime.UtcNow; // 设置Token有效期为7天 int ttl = (int)(now.AddDays(7) - new DateTime(1970, 1, 1)).TotalSeconds; string signatureString = $"{encodedUri}\n{ttl}"; using (var hmac = new HMACSHA256(Encoding.UTF8.GetBytes(key))) { byte[] signatureBytes = Encoding.UTF8.GetBytes(signatureString); byte[] hash = hmac.ComputeHash(signatureBytes); string hashInBase64 = Convert.ToBase64String(hash); return $"sr={encodedUri}&sig={Uri.EscapeDataString(hashInBase64)}&se={ttl}&skn={keyName}"; } } } public class MessageRead { public int? DeliveryCount { get; set; } public int? EnqueuedSequenceNumber { get; set; } public DateTime? EnqueuedTimeUtc { get; set; } public string LockToken { get; set; } public DateTime? LockedUntilUtc { get; set; } public string MessageId { get; set; } public int? SequenceNumber { get; set; } public string State { get; set; } public long TimeToLive { get; set; } public string MessageContent { get; set; } } }
关键修改说明
- 死信请求核心修正:
- 替换原代码中用
MessageId定位消息的逻辑,改用LockToken——这是Azure Service Bus REST API要求的唯一标识锁定消息的参数,避免MessageId重复导致的错误。 - 死信操作从POST改为PUT,符合REST API的规范。
- 替换原代码中用
- SAS Token权限优化:资源URI指向具体订阅路径,而非主题路径,确保Token权限最小化,提升安全性。
- 锁时长配置:读取消息时通过
lock-duration=PT5M参数设置5分钟的锁有效期,避免死信操作时锁过期。 - 死信元数据补充:传入
DeadLetterReason和DeadLetterErrorDescription参数,方便后续排查死信原因。 - 代码冗余清理:移除重复的响应状态判断,优化消息元数据解析逻辑。
内容的提问来源于stack exchange,提问作者Sandy
相关产品推荐
相关产品推荐

