Spring应用中通过Spring Cloud Azure Stream Binder续期Azure Bus消息锁
使用Spring Cloud Azure Stream Binder实现Azure Service Bus消息锁续期
我们的应用通过Spring Cloud Azure Stream Binder消费Azure Service Bus消息,现有消费者代码如下:
import com.azure.spring.messaging.checkpoint.Checkpointer; [...] import static com.azure.spring.messaging.AzureHeaders.CHECKPOINTER; @SpringBootApplication public class ServiceBusApplication { [...] @Bean public Consumer<Message<String>> consume() { return message -> { Checkpointer checkpointer = (Checkpointer) message.getHeaders().get(CHECKPOINTER); checkpointer.success() .doOnSuccess(s -> LOGGER.info("Message '{}' successfully checkpointed", message.getPayload())) .doOnError(e -> LOGGER.error("Error found", e)) .block(); }; } }
当消息处理耗时较长时,需要编程方式续期消息锁,原本可通过com.azure.messaging.servicebus.ServiceBusReceiverAsyncClient.renewMessageLock()实现,但无法直接获取该客户端实例。请问无需直接重写Azure Java SDK代码,仅通过Spring Binder能否完成该操作?
可以通过Spring Cloud Azure Stream Binder提供的原生能力实现,无需直接操作底层SDK客户端,具体实现方式如下:
从消息头获取
ServiceBusReceivedMessage实例
Spring Cloud Azure Binder会自动将ServiceBusReceivedMessage注入到消息头中,该实例自带锁续期方法,无需额外获取客户端:import com.azure.spring.messaging.servicebus.core.ServiceBusMessageHeaders; import com.azure.messaging.servicebus.ServiceBusReceivedMessage; // 在消费方法内获取 ServiceBusReceivedMessage receivedMessage = message.getHeaders() .get(ServiceBusMessageHeaders.RECEIVED_MESSAGE, ServiceBusReceivedMessage.class);调用锁续期方法
ServiceBusReceivedMessage提供了同步和异步两种续期方式,可根据业务场景选择:// 同步续期(适合短耗时场景) receivedMessage.renewLock(); // 异步续期(推荐,避免阻塞消费主线程) receivedMessage.renewLockAsync() .doOnSuccess(v -> LOGGER.info("消息锁续期成功,消息ID: {}", receivedMessage.getMessageId())) .doOnError(e -> LOGGER.error("消息锁续期失败", e)) .subscribe();完整消费方法示例
结合耗时处理场景,将续期逻辑放在异步线程中定时执行:import com.azure.spring.messaging.checkpoint.Checkpointer; import com.azure.spring.messaging.servicebus.core.ServiceBusMessageHeaders; import com.azure.messaging.servicebus.ServiceBusReceivedMessage; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.springframework.boot.SpringBootApplication; import org.springframework.context.annotation.Bean; import org.springframework.messaging.Message; import java.util.function.Consumer; import static com.azure.spring.messaging.AzureHeaders.CHECKPOINTER; @SpringBootApplication public class ServiceBusApplication { private static final Logger LOGGER = LoggerFactory.getLogger(ServiceBusApplication.class); @Bean public Consumer<Message<String>> consume() { return message -> { // 获取带锁信息的ReceivedMessage ServiceBusReceivedMessage receivedMessage = message.getHeaders() .get(ServiceBusMessageHeaders.RECEIVED_MESSAGE, ServiceBusReceivedMessage.class); // 异步线程定时续期锁(模拟耗时处理场景) new Thread(() -> { try { // 每30秒续期一次,共续期5次 for (int i = 0; i < 5; i++) { Thread.sleep(30000); receivedMessage.renewLock(); LOGGER.info("已为消息[{}]续期锁", message.getPayload()); } } catch (InterruptedException e) { Thread.currentThread().interrupt(); LOGGER.error("续期线程中断", e); } }).start(); // 原有检查点逻辑 Checkpointer checkpointer = (Checkpointer) message.getHeaders().get(CHECKPOINTER); checkpointer.success() .doOnSuccess(s -> LOGGER.info("Message '{}' successfully checkpointed", message.getPayload())) .doOnError(e -> LOGGER.error("Error found", e)) .block(); }; } }
说明
- 该方式兼容Spring Cloud Azure 4.x及以上版本,无需额外引入SDK依赖
- 续期逻辑建议放在异步线程中执行,避免阻塞消费主线程导致消息处理延迟
ServiceBusReceivedMessage由Binder自动注入,无需手动实例化
内容的提问来源于stack exchange,提问作者Emil Salageanu
相关产品推荐
相关产品推荐

