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

Java(Maven) Azure函数中Service Bus Topic的Peek and Lock模式实现

在Azure Functions Java版ServiceBusTopicTrigger中手动实现Peek and Lock模式

Azure Service Bus触发器默认采用PeekLock模式,但默认行为是函数执行成功自动完成消息、失败自动放弃。要手动控制消息的完成/放弃(不创建新接收器),按以下步骤操作:

核心配置与代码修改

1. 关闭触发器自动完成

在@ServiceBusTopicTrigger注解中添加autoComplete = false,禁用自动完成逻辑,消息状态将由你手动控制:

2. 绑定原始消息与操作工具

将函数参数从List<Sales>改为List<ServiceBusReceivedMessage>,同时注入ServiceBusMessageActions——这是Functions内置的工具,无需自行创建接收器即可操作消息状态。

3. 手动反序列化消息

从ServiceBusReceivedMessage中读取消息体,自行反序列化为你的Sales POJO。

4. 按业务逻辑控制消息状态

根据处理结果,调用actions.complete()或actions.abandon()来决定消息是被删除还是回到队列等待重试。

完整代码示例:

import com.azure.messaging.servicebus.ServiceBusReceivedMessage;
import com.azure.messaging.servicebus.ServiceBusMessageActions;
import com.microsoft.azure.functions.*;
import com.microsoft.azure.functions.annotation.FunctionName;
import com.microsoft.azure.functions.annotation.ServiceBusTopicTrigger;
import com.fasterxml.jackson.databind.ObjectMapper;

import java.util.List;

public class ServiceBusTopicTriggerJava {

    private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper();

    @FunctionName("ServiceBusTopicTriggerJava")
    public void run(
            @ServiceBusTopicTrigger(
                    name = "messages",
                    topicName = "aftopic",
                    subscriptionName = "afSubscription",
                    connection = "ServiceBusConnectionString",
                    autoComplete = false // 关键:禁用自动完成
            ) final List<ServiceBusReceivedMessage> messages,
            final ServiceBusMessageActions actions,
            final ExecutionContext context
    ) {
        context.getLogger().info("收到" + messages.size() + "条待处理消息");

        for (ServiceBusReceivedMessage message : messages) {
            try {
                // 手动反序列化消息体到Sales对象
                Sales sales = OBJECT_MAPPER.readValue(message.getBody(), Sales.class);
                context.getLogger().info("开始处理消息ID:" + sales.getId());

                // 模拟业务判断逻辑
                if (sales.getQuantity() > 0) {
                    // 处理成功:完成消息,从订阅中移除
                    actions.complete(message);
                    context.getLogger().info("消息" + sales.getId() + "处理完成,已移除");
                } else {
                    // 处理不符合要求:放弃消息,回到队列等待重试
                    actions.abandon(message);
                    context.getLogger().info("消息" + sales.getId() + "不符合条件,已放回队列");
                }

            } catch (Exception e) {
                context.getLogger().severe("处理消息失败:" + e.getMessage());
                // 异常场景:可选择放弃重试,或调用deadLetter移至死信队列
                actions.abandon(message);
            }
        }
    }
}

// 你的Sales POJO保持不变
@Data
@Builder
public class Sales {
    private String id;
    private String itemId;
    private String itemName;
    private double quantity;
    private double price;
}

关键细节说明

  • autoComplete = false:必须设置,否则触发器会在函数执行结束后自动完成所有消息,覆盖你的手动操作。
  • ServiceBusMessageActions:这是Functions提供的原生工具,完全满足你“无需创建新接收器”的需求,支持complete、abandon、deadLetter等操作。
  • 消息状态说明:
    • complete(message):标记消息已处理,永久从订阅中删除。
    • abandon(message):释放消息锁,消息会回到订阅,等待下一次接收(重试次数受订阅的最大传递次数限制)。
    • deadLetter(message):将消息移至死信队列,不再参与正常接收流程,适合处理无法修复的错误消息。

内容的提问来源于stack exchange,提问作者Marko Bozic

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 00:00:13