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

使用Azure队列存储REST API实现队列不存在时自动创建后调用PutMessage发消息

实现方案

你可以通过捕获PutMessage调用时抛出的队列不存在异常,触发创建队列后重试发送消息的逻辑实现需求,具体修改后的代码如下:

using System.Net;
using System.Globalization;
using System.Text;
using System.IO;

// 请将StorageAccountName、CreateAuthorizationHeader、CreateQueue替换为你实际的类变量/方法
public static string StorageAccountName { get; set; }
public static string CreateAuthorizationHeader(string stringToSign)
{
    // 你原有实现逻辑
    throw new NotImplementedException();
}
public static void CreateQueue(string queueName)
{
    // 你原有实现逻辑
    throw new NotImplementedException();
}

public static void PutMessage(String queueName, String message)
{
    bool queueCreated = false;
    // 最多重试1次,队列不存在创建后重试
    for (int retryCount = 0; retryCount < 2; retryCount++)
    {
        String requestMethod = "POST";
        String urlPath = $"{queueName}/messages";
        String storageServiceVersion = "2017-11-09";
        String dateInRfc1123Format = DateTime.UtcNow.ToString("R", CultureInfo.InvariantCulture);

        String messageText = $"<QueueMessage><MessageText>{message}</MessageText></QueueMessage>";
        UTF8Encoding utf8Encoding = new UTF8Encoding();
        Byte[] messageContent = utf8Encoding.GetBytes(messageText);
        Int32 messageLength = messageContent.Length;

        String canonicalizedHeaders = String.Format(
                "x-ms-date:{0}\nx-ms-version:{1}",
                dateInRfc1123Format,
                storageServiceVersion);
        String canonicalizedResource = $"/{StorageAccountName}/{urlPath}";
        String stringToSign = $"{requestMethod}\n\n\n{messageLength}\n\n\n\n\n\n\n\n\n{canonicalizedHeaders}\n{canonicalizedResource}";
                
        String authorizationHeader = CreateAuthorizationHeader(stringToSign);

        Uri uri = new Uri("https://" + StorageAccountName + ".queue.Azure.com/" + urlPath);
        HttpWebRequest request = (HttpWebRequest)WebRequest.Create(uri);
        request.Method = requestMethod;
        request.Headers.Add("x-ms-date", dateInRfc1123Format);
        request.Headers.Add("x-ms-version", storageServiceVersion);
        request.Headers.Add("Authorization", authorizationHeader);
        request.ContentLength = messageLength;

        try
        {
            using (Stream requestStream = request.GetRequestStream())
            {
                requestStream.Write(messageContent, 0, messageLength);
            }

            using (HttpWebResponse response = (HttpWebResponse)request.GetResponse())
            {
                String requestId = response.Headers["x-ms-request-id"];
                // 发送成功直接退出循环
                return;
            }
        }
        catch (WebException ex)
        {
            // 判断是否是队列不存在的404错误,且还没尝试创建过队列
            if (!queueCreated && ex.Response is HttpWebResponse response && response.StatusCode == HttpStatusCode.NotFound)
            {
                // 调用你已实现的创建队列方法
                CreateQueue(queueName);
                queueCreated = true;
                // 进入下一轮循环重试发送
                continue;
            }
            // 其他异常直接抛出
            throw;
        }
    }
}

注意事项

  • CreateQueue操作本身是幂等的,就算多个客户端同时触发队列创建也不会报错,不会产生并发冲突
  • 你可以根据业务需要调整重试次数,常规场景下创建队列后重试1次即可满足需求
  • 每次重试都重新构造请求、生成时间戳和签名,避免Azure存储的签名过期校验失败

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 09:06:01