使用Durable Functions推送Service Bus消息出现数量不可靠异常
问题根因
你遇到的消息数量超过预期的核心原因是Durable Functions Activity函数的默认自动重试机制:当提交消息到Service Bus时出现瞬时网络抖动、Service Bus限流响应等异常时,Durable Functions会自动重试执行Activity函数,你的Activity每次执行都会提交一条新消息,没有做幂等校验,因此重试次数越多,总消息数就比预期的3000越高,数值随机是因为每次运行的瞬时故障次数不固定。
另外你的代码还有两处不符合Durable Functions最佳实践的问题,也可能放大故障概率:
- Orchestrator函数中直接使用传入的ILogger打日志,没有使用重放安全的日志实例,Orchestrator重放时会重复生成日志,干扰问题排查
- 一次性Fan-Out 3000个并行Activity,瞬时请求量过高容易触发Service Bus限流,进一步提高重试概率
解决方案
1. 自定义Activity重试策略,限制重试次数
将Orchestrator中调用Activity的方法从context.CallActivityAsync改为context.CallActivityWithRetryAsync,明确设置重试规则,如果你不需要重试可以直接把最大重试次数设为1:
// 定义重试策略:最多尝试1次(即不重试) var retryOptions = new RetryOptions(TimeSpan.FromSeconds(2), 1) { Handle = ex => ex is ServiceBusException // 可指定仅针对特定异常重试 }; for (i_batch = 0; i_batch < restInputs.Count; i_batch++) { // 改用带重试规则的调用方法 parallelTasks.Add(context.CallActivityWithRetryAsync("EmailQueueSubmitter_ActivitySendMessageBatchSingleton", retryOptions, i_batch.ToString())); }
如果需要保留重试能力,配合下面的幂等方案使用即可。
2. 实现提交操作幂等,避免重复消息
方案A:配合Service Bus重复检测功能
- 首先在Azure Portal中开启对应Service Bus队列的重复检测功能,设置合适的检测时间窗口(比如1小时)
- 修改消息MessageId的生成规则,确保同一个任务对应的消息ID固定,不要每次生成新的Guid:
// 在Activity中获取Orchestration实例ID和当前批次号,拼接为唯一MessageId var instanceId = activityContext.InstanceId; var sbMsg = new ServiceBusMessage(Payload) { // 用实例ID+批次号作为唯一标识,同一个任务重试时MessageId不变 MessageId = $"{instanceId}_{Payload}", ContentType = "text/plain" };
这样即使Activity被重试多次,Service Bus也会自动丢弃MessageId重复的消息,不会产生多余数据。
方案B:提前做批次去重
你也可以在提交前自行记录已经提交成功的批次号,重试时跳过已经提交成功的批次。
3. 优化Orchestrator代码实践
使用重放安全的日志实例,避免Orchestrator重放时重复打印日志:
// 在Orchestrator函数入口替换ILogger log = context.CreateReplaySafeLogger(log);
如果3000个并行Activity压力过大,可以改为分批并行,比如每次跑200个,跑完一批再跑下一批,降低限流概率。
内容的提问来源于stack exchange,提问作者147
相关产品推荐
相关产品推荐

