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

