Spring Cloud Stream Kafka应用中如何按条件发送消息到多主题之一
基于Spring Cloud Function实现消息动态路由到多Topic
针对你的需求,我们可以通过两种方式实现根据校验逻辑将消息发送到指定Topic,以下是具体实现方案:
方案一:使用StreamBridge动态发送(推荐)
这种方式无需修改原有输出绑定配置,通过StreamBridge直接动态指定目标Topic,灵活性更高。
步骤1:修改SinkService,注入StreamBridge并调整处理逻辑
import org.springframework.cloud.stream.function.StreamBridge; import org.springframework.messaging.Message; import org.springframework.messaging.support.MessageBuilder; import org.springframework.stereotype.Service; import lombok.extern.slf4j.Slf4j; @Slf4j @Service("SinkService") public class SinkService<T> { private final StreamBridge streamBridge; // 构造注入StreamBridge public SinkService(StreamBridge streamBridge) { this.streamBridge = streamBridge; } public void processMessage(Message<SourceMessage> message) { log.info("Message consumed at {} \n{}", message.getHeaders().getTimestamp(), message.getPayload()); try { SourceMessage payload = message.getPayload(); if (payload.isManaged()) { int type = payload.getType(); // 假设SourceMessage提供getType()方法获取类型值 SinkMessage sinkMessage = new SinkMessage(); sinkMessage.setPayload(payload); Message<SinkMessage> outputMsg = MessageBuilder.withPayload(sinkMessage).build(); if (type == 2) { // 发送到配置中producerBean-out-0对应的topic1 streamBridge.send("producerBean-out-0", outputMsg); } else if (type == 4) { // 发送到配置中producerBean-out-1对应的topic2 streamBridge.send("producerBean-out-1", outputMsg); } else { log.warn("Unsupported message type: {}, no message sent", type); } } } catch (Exception e) { log.error("Failed to process message", e); } } }
步骤2:修改Function定义为Consumer
因为不再通过函数的输出绑定发送消息,将原Function改为Consumer:
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.messaging.Message; import java.util.function.Consumer; @Configuration public class FunctionConfig { @Bean("producerBean") public Consumer<Message<SourceMessage>> producerBean(SinkService<SourceMessage> sinkService) { return sinkService::processMessage; } }
配置说明
原有application.properties无需修改,StreamBridge会通过绑定名称(如producerBean-out-0)自动匹配配置中的目标Topic。
方案二:使用多输出函数绑定
通过定义返回多结果的Function,将消息对应到不同的输出绑定(out-0/out-1),适合固定多分支路由场景。
步骤1:修改SinkService的返回类型和逻辑
import org.springframework.messaging.Message; import org.springframework.messaging.support.MessageBuilder; import org.springframework.stereotype.Service; import lombok.extern.slf4j.Slf4j; @Slf4j @Service("SinkService") public class SinkService<T> { public Message<SinkMessage>[] processMessage(Message<SourceMessage> message) { // 数组长度对应输出绑定数量(out-0和out-1) Message<SinkMessage>[] result = new Message[2]; log.info("Message consumed at {} \n{}", message.getHeaders().getTimestamp(), message.getPayload()); try { SourceMessage payload = message.getPayload(); if (payload.isManaged()) { int type = payload.getType(); SinkMessage sinkMessage = new SinkMessage(); sinkMessage.setPayload(payload); Message<SinkMessage> outputMsg = MessageBuilder.withPayload(sinkMessage).build(); if (type == 2) { // 对应out-0,发送到topic1 result[0] = outputMsg; } else if (type == 4) { // 对应out-1,发送到topic2 result[1] = outputMsg; } else { log.warn("Unsupported message type: {}, no message sent", type); } } } catch (Exception e) { log.error("Failed to process message", e); } return result; } }
步骤2:修改Function的返回类型
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.messaging.Message; import java.util.function.Function; @Configuration public class FunctionConfig { @Bean("producerBean") public Function<Message<SourceMessage>, Message<SinkMessage>[]> producerBean(SinkService<SourceMessage> sinkService) { return sinkService::processMessage; } }
配置说明
原有application.properties保持不变,Spring Cloud Stream会自动将数组中第1个元素发送到producerBean-out-0,第2个元素发送到producerBean-out-1,null值会被忽略不发送。
注意事项
- 确保
SourceMessage类提供getType()方法,用于获取消息类型值; - 替换
printStackTrace()为日志框架记录异常,便于排查问题; - 使用
MessageBuilder构建消息是Spring官方推荐的方式,比直接实例化GenericMessage更灵活。
内容的提问来源于stack exchange,提问作者Dharita Chokshi
相关产品推荐
相关产品推荐

