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

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值会被忽略不发送。


注意事项

  1. 确保SourceMessage类提供getType()方法,用于获取消息类型值;
  2. 替换printStackTrace()为日志框架记录异常,便于排查问题;
  3. 使用MessageBuilder构建消息是Spring官方推荐的方式,比直接实例化GenericMessage更灵活。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 08:55:18