Spring Boot Kafka Stream自定义SerDe处理不同输入输出类型疑问
解决方案:针对不同输入输出类型分别处理序列化/反序列化
SerDe的设计初衷就是绑定单一数据类型的序列化与反序列化逻辑,所以不需要强行实现跨类型的SerDe,正确的做法是为InputType和OutputType分别配置对应的序列化/反序列化组件,具体有两种实现方式:
方式一:为每种类型单独实现完整SerDe
为InputType和OutputType各自编写独立的SerDe类,每个SerDe只负责自身类型的序列化与反序列化。
1. 实现单个类型的Serializer和Deserializer
以InputType为例:
// InputType序列化器 public class InputTypeSerializer implements Serializer<InputType> { private final ObjectMapper objectMapper = new ObjectMapper(); @Override public byte[] serialize(String topic, InputType data) { try { return objectMapper.writeValueAsBytes(data); } catch (JsonProcessingException e) { throw new SerializationException("序列化InputType失败", e); } } @Override public void configure(Map<String, ?> configs, boolean isKey) {} @Override public void close() {} } // InputType反序列化器 public class InputTypeDeserializer implements Deserializer<InputType> { private final ObjectMapper objectMapper = new ObjectMapper(); @Override public InputType deserialize(String topic, byte[] data) { try { return objectMapper.readValue(data, InputType.class); } catch (IOException e) { throw new SerializationException("反序列化InputType失败", e); } } @Override public void configure(Map<String, ?> configs, boolean isKey) {} @Override public void close() {} }
2. 封装为SerDe类
将上面的序列化器和反序列化器封装成SerDe:
public class InputTypeSerDe implements SerDe<InputType> { @Override public Serializer<InputType> serializer() { return new InputTypeSerializer(); } @Override public Deserializer<InputType> deserializer() { return new InputTypeDeserializer(); } @Override public void configure(Map<String, ?> configs, boolean isKey) { serializer().configure(configs, isKey); deserializer().configure(configs, isKey); } @Override public void close() { serializer().close(); deserializer().close(); } }
同理实现OutputTypeSerDe,之后在KStream代码中指定对应的SerDe:
// 消费输入主题时指定InputType的SerDe KStream<String, InputType> inputStream = builder.stream( "input-topic", Consumed.with(Serdes.String(), new InputTypeSerDe()) ); // 处理后输出时指定OutputType的SerDe KStream<String, OutputType> outputStream = inputStream.transform(/* 你的业务处理逻辑 */); outputStream.to( "output-topic", Produced.with(Serdes.String(), new OutputTypeSerDe()) );
方式二:直接使用独立的Serializer/Deserializer(无需封装SerDe)
如果不想编写完整的SerDe类,可以直接在消费/生产时分别指定所需的反序列化器/序列化器:
// 消费时只指定InputType的反序列化器 Consumed<String, InputType> consumedConfig = Consumed.with( Serdes.String(), Serdes.serdeFrom(null, new InputTypeDeserializer()) ); KStream<String, InputType> inputStream = builder.stream("input-topic", consumedConfig); // 生产时只指定OutputType的序列化器 Produced<String, OutputType> producedConfig = Produced.with( Serdes.String(), Serdes.serdeFrom(new OutputTypeSerializer(), null) ); outputStream.to("output-topic", producedConfig);
这里Serdes.serdeFrom()方法允许分别传入序列化器和反序列化器,由于消费场景只用到反序列化器、生产场景只用到序列化器,所以另一参数传null即可。
核心逻辑说明
输入和输出属于不同的业务对象,本身就应该使用独立的序列化/反序列化逻辑,强行将两种类型绑定到同一个SerDe中会违反单一职责原则,增加代码维护成本。通过上述两种方式,就能清晰地为不同类型配置对应的序列化组件,完全适配你的KStream处理流程。
内容的提问来源于stack exchange,提问作者Manish Kumar
相关产品推荐
相关产品推荐

