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

Spring Boot中Spring Cloud Stream Kafka多输出主题实现问题

问题:Spring Boot + Avro 实现Kafka多主题输出

我有一个Spring Boot应用,通过指定Avro schema将Kafka事件发送到单个主题,示例代码如下:

@Configuration
@Slf4j
public class KafkaListener {

  @Bean
  public Function<List<DummyClass>, List<DummyClass>> acceptEvent() {
     return messages -> {
        List<DummyClass> output = new ArrayList<>();
        messages.forEach(msg -> {
          // do work and add items to the output field
        });
        // some more work
        return output;
     };
  }
}

其中DummyClass是由resources文件夹中的.avsc文件通过Avro插件生成的类。

现在需要将事件发送到两个主题而非单个主题,尝试两种方式均失败:

  1. 使用Tuple2返回双列表
@Configuration
@Slf4j
public class KafkaListener {

  @Bean
  public Function<List<DummyClass>, Tuple2<List<DummyClass>,List<DummyClass>>> acceptEvent() {
     return messages -> {
        List<DummyClass> output1 = new ArrayList<>();
        List<DummyClass> output2 = new ArrayList<>();
        messages.forEach(msg -> {
          // do work and add items to the output1 and output2 fields
        });
        // some more work
        return Tuples.of(output1,output2);
     };
  }
}

运行时抛出异常:

java.lang.UnsupportedOperationException: At the moment only Tuple-based function are supporting multiple arguments
  1. 使用自定义对象封装输出,抛出Avro类型不支持异常:
Unsupported Avro type. Supported types are null, Boolean, Integer, Long, Float, Double, String, byte[] and IndexedRecord

需要解决:如何实现将事件写入多个Kafka主题?现有实现问题在哪?还有哪些可行方案?


问题分析
  1. Tuple方案失败原因:Spring Cloud Stream中,Tuple多输出要求输入输出为流类型(如Flux/Message),而非集合类型List。用List作为输入输出不符合Tuple多绑定的要求,因此抛出不支持的异常。
  2. 自定义对象方案失败原因:Avro序列化仅支持IndexedRecord(Avro生成类的父接口)及基础类型,自定义封装类不是Avro生成的IndexedRecord子类,无法被Avro序列化器识别。

解决方案

方案1:Spring Cloud Stream 流类型(Flux)+ Tuple 实现多主题输出

将输入输出改为Flux,配合Tuple2实现双主题输出,同时在配置文件中绑定两个输出目标:

代码实现

@Configuration
@Slf4j
public class KafkaEventProcessor {

  @Bean
  public Function<Flux<DummyClass>, Tuple2<Flux<DummyClass>, Flux<DummyClass>>> processEvents() {
     return inputFlux -> {
        // 拆分流为两个分支
        Flux<DummyClass> output1 = inputFlux.filter(msg -> {
          // 替换为实际判断逻辑,筛选发送到第一个主题的消息
          return true; 
        }).doOnNext(msg -> {
          // 对第一个主题的消息执行业务处理
        });

        Flux<DummyClass> output2 = inputFlux.filter(msg -> {
          // 替换为实际判断逻辑,筛选发送到第二个主题的消息
          return false; 
        }).doOnNext(msg -> {
          // 对第二个主题的消息执行业务处理
        });

        return Tuples.of(output1, output2);
     };
  }
}

配置文件(application.yml)

spring:
  cloud:
    stream:
      bindings:
        processEvents-in-0:
          destination: input-topic # 输入主题名称
        processEvents-out-0:
          destination: output-topic-1 # 第一个输出主题
          producer:
            use-native-encoding: true # 启用Avro原生编码
        processEvents-out-1:
          destination: output-topic-2 # 第二个输出主题
          producer:
            use-native-encoding: true
      kafka:
        bindings:
          processEvents-in-0:
            consumer:
              use-native-decoding: true # 启用Avro原生解码
        binder:
          configuration:
            schema.registry.url: http://your-schema-registry-url # 替换为你的Schema Registry地址

方案2:使用KStream实现多主题输出

如果使用Spring Kafka Streams,可以通过分支(branch)操作将消息路由到不同主题,适合更复杂的流处理场景:

代码实现

@Configuration
@Slf4j
public class KafkaStreamProcessor {

  @Bean
  public KStream<String, DummyClass> kStream(StreamsBuilder streamsBuilder) {
    KStream<String, DummyClass> inputStream = streamsBuilder.stream("input-topic");

    // 按条件分支为两个流
    KStream<String, DummyClass>[] branches = inputStream.branch(
      (key, value) -> {
        // 替换为实际判断逻辑,筛选发送到第一个主题的消息
        return true;
      },
      (key, value) -> {
        // 替换为实际判断逻辑,筛选发送到第二个主题的消息
        return false;
      }
    );

    // 将分支流发送到对应主题
    branches[0].to("output-topic-1");
    branches[1].to("output-topic-2");

    return inputStream;
  }
}

配置文件(application.yml)

spring:
  kafka:
    streams:
      application-id: kafka-stream-app-id # 自定义流应用ID
      bootstrap-servers: your-kafka-bootstrap-servers # 替换为你的Kafka地址
      properties:
        schema.registry.url: http://your-schema-registry-url # 替换为你的Schema Registry地址
        default.key.serde: org.apache.kafka.common.serialization.Serdes$StringSerde
        default.value.serde: io.confluent.kafka.streams.serdes.avro.SpecificAvroSerde

关键注意事项

  • 所有输出的消息类型必须是Avro生成的IndexedRecord子类(如DummyClass),避免使用自定义非Avro类导致序列化失败。
  • 使用Spring Cloud Stream时,确保输入输出绑定名称与Bean方法名对应(processEvents-in-0对应方法processEvents的输入,out-0/out-1对应Tuple的两个输出)。
  • 必须配置Schema Registry地址,确保Avro序列化器能正确获取schema信息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 12:27:47