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

使用Spring Integration向Azure Service Bus主题发消息时如何获取生成的消息ID

要通过Spring Integration获取Azure Service Bus生成的消息ID,只需调整现有代码的网关配置和消息处理器配置即可,具体修改如下:

第一步:开启消息处理器的回复生成能力

DefaultMessageHandler默认是单向发送不返回结果,需要开启producesReply配置,让其在发送成功后返回包含服务端响应信息的回复消息,即可从回复消息的头部获取Azure生成的消息ID。

第二步:调整消息网关的返回值定义

将原来返回void的网关方法修改为返回Message类型,接收消息处理器返回的回复内容。

修改后的完整代码示例

import com.azure.spring.integration.servicebus.core.ServiceBusHeaders;
import org.springframework.integration.annotation.MessagingGateway;
import org.springframework.integration.annotation.ServiceActivator;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageHandler;
import org.springframework.stereotype.Service;

@Service
public class ServiceBusDemo {
    private static final String OUTPUT_CHANNEL = "topic.output";
    private static final String TOPIC_NAME = "my_topic";

    @Autowired
    TopicOutboundGateway messagingGateway;

    public String send(String message) {
        // 发送消息并接收回复
        Message<?> reply = this.messagingGateway.send(message);
        // 从头部获取Azure生成的消息ID
        String azureMessageId = reply.getHeaders().get(ServiceBusHeaders.MESSAGE_ID, String.class);
        System.out.println("Azure生成的消息ID: " + azureMessageId);
        return azureMessageId;
    }

    @Bean
    @ServiceActivator(inputChannel = OUTPUT_CHANNEL)
    public MessageHandler topicMessageSender(ServiceBusTopicOperation topicOperation) {
        DefaultMessageHandler handler = new DefaultMessageHandler(TOPIC_NAME, topicOperation);
        // 开启回复生成能力
        handler.setProducesReply(true);
        handler.setSendCallback(new ListenableFutureCallback<>() {
            @Override
            public void onSuccess(Void result) {
                System.out.println("Message was sent successfully to service bus.");
            }

            @Override
            public void onFailure(Throwable ex) {
                System.out.println("There was an error sending the message to service bus.");
            }
        });

        return handler;
    }

    @MessagingGateway(defaultRequestChannel = OUTPUT_CHANNEL)
    public interface TopicOutboundGateway {
        // 修改返回值为Message类型接收回复
        Message<?> send(String text);
    }
}

异步场景适配

如果需要使用异步发送逻辑,只需将网关方法的返回值修改为CompletableFuture<Message<?>>即可,无需调整其他核心配置。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 20:27:00