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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 17:36:03