Spring Cloud Stream中为每个绑定配置自定义Key序列化/反序列化
为Spring Cloud Stream绑定单独配置Key序列化器/反序列化器
你已经掌握了全局Key Serdes的配置方式,要为每个绑定指定不同的Key序列化/反序列化器其实很简单——Spring Cloud Stream的Kafka binder支持针对单个绑定单独配置,不需要覆盖全局设置,只需要在对应绑定的Kafka专属配置节点下指定即可。
核心配置方式
针对每个绑定,你可以分别配置生产者的Key序列化器和消费者的Key反序列化器:
1. 生产者绑定配置Key序列化器
在spring.cloud.stream.kafka.bindings.<你的output绑定名>.producer节点下添加keySerializer属性,指定序列化器的全类名:
spring: cloud: stream: kafka: bindings: # 为名为output的绑定指定Integer类型的Key序列化器 output: producer: keySerializer: org.apache.kafka.common.serialization.IntegerSerializer
2. 消费者绑定配置Key反序列化器
在spring.cloud.stream.kafka.bindings.<你的input绑定名>.consumer节点下添加keyDeserializer属性,指定反序列化器的全类名:
spring: cloud: stream: kafka: bindings: # 为名为input的绑定指定Integer类型的Key反序列化器 input: consumer: keyDeserializer: org.apache.kafka.common.serialization.IntegerDeserializer
完整配置示例(混合不同Key类型)
比如你想让input绑定处理Integer类型的Key,output绑定处理String类型的Key,完整配置如下:
spring: cloud: stream: bindings: input: contentType: application/*+avro destination: user group: my-group output: contentType: application/*+avro destination: user producer: partition-count: 2 kafka: binder: brokers: default:9092 schemaRegistryClient: endpoint: http://default:8081 bindings: # 消费者input绑定:Key用Integer反序列化 input: consumer: keyDeserializer: org.apache.kafka.common.serialization.IntegerDeserializer # 生产者output绑定:Key用String序列化 output: producer: keySerializer: org.apache.kafka.common.serialization.StringSerializer # 全局配置可以保留(如果有其他绑定需要默认值),但单独绑定的配置会覆盖全局 kafka: consumer: keyDeserializer: org.apache.kafka.common.serialization.StringDeserializer producer: keySerializer: org.apache.kafka.common.serialization.StringSerializer
对应代码调整
消费者接收Integer类型Key
修改@Header注解的类型为Integer:
@StreamListener(Sink.INPUT) public void handle(@Payload UserValue user, @Headers Map<String, Object> headers, @Header(KafkaHeaders.RECEIVED_MESSAGE_KEY) Integer key) { System.out.println("Received: " + user + " with key: " + key + " and headers: " + headers); }
生产者发送Integer类型Key
如果你的output绑定配置了Integer序列化器,就可以直接发送Integer类型的Key:
UserValue user = UserValue.newBuilder().setName("Alessandro").setSurname("Dionisi").build(); // 发送Integer类型的Key output.send(MessageBuilder.withPayload(user).setHeader(KafkaHeaders.MESSAGE_KEY, 1).build());
注意事项
- 单独绑定的配置优先级高于全局的
spring.kafka.consumer/producer配置,只会影响指定的绑定 - 如果使用自定义的Serdes(比如自定义序列化逻辑),只需要把配置值换成你的自定义类全类名即可
- 要保证生产者的Key序列化器和对应消费者的Key反序列化器类型匹配,否则会出现反序列化失败的问题
内容的提问来源于stack exchange,提问作者Alessandro Dionisi
相关产品推荐
相关产品推荐

