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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 03:34:54