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

Azure Functions向Service Bus主题发消息是否需权限及排障

问题

我有一个基于Service Bus主题触发的Azure Function(v3),需要向Service Bus主题发送消息。Service Bus与Azure Function处于同一订阅和资源组中。

  1. 是否需要为Azure Function授予权限以允许其向Service Bus主题发送消息?
  2. 我已在Azure Function中配置了Service Bus连接字符串,但尝试发送消息时,消息无法发出且未抛出任何错误。

发送消息的AzureServiceBus类代码:

internal class AzureServiceBus : IMessageBus
{
    #region private fields
    private const string ContentType = "application/json";
    private readonly ServiceBusSender _serviceBusSender;
    #endregion

    #region ctors
    internal AzureServiceBus(ServiceBusSender serviceBusSender)
    {
        this._serviceBusSender = serviceBusSender;
    }
    #endregion

    #region public methods
    /// <summary>
    /// Publish message into Service Bus Topic 
    /// </summary>
    /// <typeparam name="T"></typeparam>
    /// <param name="message"></param>
    /// <returns></returns>
    public async Task PublishMessageAsync<T>(T message, string messageId)
    {
        var serviceBusMessage = new ServiceBusMessage(Encoding.UTF8.GetBytes(message.AsJson()))
        {
            ContentType = ContentType
        };
        serviceBusMessage.MessageId = messageId;
        await this._serviceBusSender.SendMessageAsync(serviceBusMessage);
    }

    /// <summary>
    /// Publish message into Service Bus Topic 
    /// </summary>
    /// <typeparam name="T"></typeparam>
    /// <param name="message"></param>
    /// <param name="messageId"></param>
    /// <param name="filterPropertyName"></param>
    /// <param name="filterPropertyValue"></param>
    /// <returns></returns>
    public async Task PublishMessageAsync<T>(T message, string messageId, string filterPropertyName, string filterPropertyValue)
    {
        try
        {
            var serviceBusMessage = new ServiceBusMessage(Encoding.UTF8.GetBytes(message.AsJson()))
            {
                ContentType = ContentType
            };
            serviceBusMessage.MessageId = messageId;
            serviceBusMessage.ApplicationProperties.Add(filterPropertyName, filterPropertyValue);
            serviceBusMessage.CorrelationId = filterPropertyValue;

            await this._serviceBusSender.SendMessageAsync(serviceBusMessage);
        }
        catch (Exception ex)
        {
            Console.WriteLine(ex?.StackTrace);
            throw;
        }
    }

    public async Task<bool> PublishMessageReturnTrueAsync<T>(T message, string messageId, string filterPropertyName, string filterPropertyValue)
    {
        try
        {
            var serviceBusMessage = new ServiceBusMessage(Encoding.UTF8.GetBytes(message.AsJson()))
            {
                ContentType = ContentType
            };
            serviceBusMessage.MessageId = messageId;
            serviceBusMessage.ApplicationProperties.Add(filterPropertyName, filterPropertyValue);
            serviceBusMessage.CorrelationId = filterPropertyValue;

            await this._serviceBusSender.SendMessageAsync(serviceBusMessage);
            
            return true;
        }
        catch (Exception ex)
        {
            Console.WriteLine(ex?.StackTrace);
            return false;
        }
    }
    #endregion

    #region internal methods
    internal static IMessageBus Create(ServiceBusSender sender)
    {
        return new AzureServiceBus(sender);
    }
    #endregion
}

AzureServiceBusFactory工厂类代码:

public class AzureServiceBusFactory : IMessageBusFactory
{
    #region private fields
    private readonly ILogger<AzureServiceBusFactory> _logger;
    private readonly ServiceBusConfigOption _serviceBusConfigOption;
    private readonly object _lockObject = new object();
    private readonly ConcurrentDictionary<string, ServiceBusClient> _clients = new ConcurrentDictionary<string, ServiceBusClient>();
    private readonly ConcurrentDictionary<string, ServiceBusSender> _senders = new ConcurrentDictionary<string, ServiceBusSender>();

    #endregion

    public AzureServiceBusFactory(ILogger<AzureServiceBusFactory> logger,
        IOptions<ServiceBusConfigOption> serviceBusConfigOption)
    {
        _logger = logger;
        _serviceBusConfigOption = serviceBusConfigOption.Value;
    }

    #region public methods
    /// <summary>
    /// Get ServiceBusClient
    /// </summary>
    /// <param name="senderName"></param>
    /// <returns></returns>
    public IMessageBus GetClient(string senderName)
    {
        _logger.LogInformation("{class} -> {method} -> Start",
             nameof(AzureServiceBusFactory), nameof(AzureServiceBusFactory.GetClient));

        var connectionString = _serviceBusConfigOption.ConnectionString;

        var key = $"{connectionString}-{senderName}";

        if (this._senders.ContainsKey(key) && !this._senders[key].IsClosed)
        {
            _logger.LogInformation("{class} -> {method} -> End",
                   nameof(AzureServiceBusFactory), nameof(AzureServiceBusFactory.GetClient));

            return AzureServiceBus.Create(this._senders[key]);
        }

        var client = this.GetServiceBusClient(connectionString);

        lock (this._lockObject)
        {
            if (this._senders.ContainsKey(key) && this._senders[key].IsClosed)
            {
                if (this._senders[key].IsClosed)
                {
                    this._senders[key].DisposeAsync().GetAwaiter().GetResult();
                }

                _logger.LogInformation("{class} -> {method} -> End",
                    nameof(AzureServiceBusFactory), nameof(AzureServiceBusFactory.GetClient));

                return AzureServiceBus.Create(this._senders[key]);
            }

            var sender = client.CreateSender(senderName);

            this._senders[key] = sender;

            _logger.LogInformation($"ServiceBusClient created connection string:{connectionString}.");
        }

        _logger.LogInformation("{class} -> {method} -> End",
            nameof(AzureServiceBusFactory), nameof(AzureServiceBusFactory.GetClient));

        return AzureServiceBus.Create(this._senders[key]);
    }
    #endregion

    #region protected methods
    /// <summary>
    /// Create ServiceBusClient
    /// </summary>
    /// <param name="connectionString"></param>
    /// <returns></returns>
    protected virtual ServiceBusClient GetServiceBusClient(string connectionString)
    {
        var key = $"{connectionString}";

        lock (this._lockObject)
        {
            if (this.ClientDoesnotExistOrIsClosed(connectionString))
            {
                try
                {
                    var client = new ServiceBusClient(connectionString, new ServiceBusClientOptions
                    {
                        TransportType = ServiceBusTransportType.AmqpTcp
                    });

                    this._clients[key] = client;
                }
                catch (Exception ex)
                {
                    _logger.LogError(ex, $"Failed to create ServiceBusClient, {ex.Message}");
                    throw;
                }
            }

            return this._clients[key];
        }
    }
    #endregion

    #region private methods
    /// <summary>
    /// Check ClientDoesnotExistOrIsClosed
    /// </summary>
    /// <param name="connectionString"></param>
    /// <returns></returns>
    private bool ClientDoesnotExistOrIsClosed(string connectionString)
    {
        return !this._clients.ContainsKey(connectionString) || this._clients[connectionString].IsClosed;
    }
    #endregion
}
解答

权限问题

即使Service Bus与Azure Function处于同一订阅和资源组,仍需授予发送权限,有两种常用实现方式:

  • SAS连接字符串:如果使用命名空间的RootManageSharedAccessKey连接字符串,自带发送权限,但生产环境不建议(权限范围过大)。更安全的做法是为目标主题创建仅含Send权限的SAS策略,使用该策略的连接字符串。
  • 托管标识:为Azure Function启用系统分配/用户分配的托管标识,在Service Bus主题(或命名空间)的IAM中添加Azure Service Bus Data Sender角色,无需存储连接字符串,安全性更高。

消息静默失败的排查与修复

从代码来看,主要问题点及修复方式如下:

  1. 异步方法未正确等待
    若调用PublishMessageAsync等方法时未使用await,会导致Azure Function提前终止,消息未完成发送就被中断。确保所有异步调用都加上await,比如触发函数本身必须是异步的:

    [FunctionName("MyServiceBusTrigger")]
    public async Task Run([ServiceBusTrigger("mytopic", "mysubscription", Connection = "ServiceBusConnection")] string mySbMsg)
    {
        var messageBus = _messageBusFactory.GetClient("target-topic");
        await messageBus.PublishMessageAsync(mySbMsg, Guid.NewGuid().ToString());
    }
    
  2. ServiceBusSender缓存逻辑错误
    在GetClient方法中,当检测到sender已关闭时,仅dispose但未重新创建新的sender,直接返回已失效的实例,调用其SendMessageAsync会静默失败。修复代码如下:

    lock (this._lockObject)
    {
        if (this._senders.ContainsKey(key) && this._senders[key].IsClosed)
        {
            // 释放旧的sender并移除缓存
            this._senders[key].DisposeAsync().GetAwaiter().GetResult();
            this._senders.TryRemove(key, out _);
            // 重新创建新的sender
            var sender = client.CreateSender(senderName);
            this._senders[key] = sender;
            
            _logger.LogInformation("Recreated closed ServiceBusSender for key: {Key}", key);
        }
        else if (!this._senders.ContainsKey(key))
        {
            var sender = client.CreateSender(senderName);
            this._senders[key] = sender;
            
            _logger.LogInformation("Created new ServiceBusSender for key: {Key}", key);
        }
    
        _logger.LogInformation("{class} -> {method} -> End",
            nameof(AzureServiceBusFactory), nameof(AzureServiceBusFactory.GetClient));
    
        return AzureServiceBus.Create(this._senders[key]);
    }
    
  3. 日志缺失导致无法排查
    当前代码仅在异常时打印StackTrace到Console,Azure Function生产环境无法获取该输出。建议为AzureServiceBus注入ILogger,添加发送前后的日志:

    internal class AzureServiceBus : IMessageBus
    {
        private const string ContentType = "application/json";
        private readonly ServiceBusSender _serviceBusSender;
        private readonly ILogger<AzureServiceBus> _logger; // 新增
    
        internal AzureServiceBus(ServiceBusSender serviceBusSender, ILogger<AzureServiceBus> logger) // 修改构造函数
        {
            this._serviceBusSender = serviceBusSender;
            this._logger = logger;
        }
    
        public async Task PublishMessageAsync<T>(T message, string messageId)
        {
            _logger.LogInformation("准备发送消息,ID: {MessageId}", messageId);
            var jsonContent = message.AsJson();
            _logger.LogDebug("消息内容: {JsonContent}", jsonContent);
            
            var serviceBusMessage = new ServiceBusMessage(Encoding.UTF8.GetBytes(jsonContent))
            {
                ContentType = ContentType,
                MessageId = messageId
            };
            
            await this._serviceBusSender.SendMessageAsync(serviceBusMessage);
            _logger.LogInformation("消息发送成功,ID: {MessageId}", messageId);
        }
    }
    
  4. 序列化方法验证
    确保AsJson()扩展方法能正确序列化对象,避免返回空或无效JSON。可通过日志记录序列化结果,验证是否符合预期。


内容的提问来源于stack exchange,提问作者Rakesh Kumar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 16:04:53