响应式函数绑定消费Avro事件批处理:指定类型参数丢失Payload
在Spring Boot应用中,我需要消费Avro事件批处理。因业务需求需手动确认批处理,因此需要将入站批处理包装在Message中并访问其Header。我编写了如下绑定:
@Bean public Function<Flux<Message<List<MyAvroEvent>>>, Mono<Void>> consumeAvroEvents() { return flux -> flux .concatMap(events -> { Acknowledgment acknowledgment = events.getHeaders().get(KafkaHeaders.ACKNOWLEDGMENT, Acknowledgment.class); List<MyAvroEvent> events = events.getPayload(); // 此列表为空 // 保存事件到数据库 // 无超时异常时确认,否则不确认 }) .then(); }
我发现events.getPayload()为空,但移除类型参数改用通配符后,Payload不为空:
@Bean public Function<Flux<Message<?>>, Mono<Void>> consumeAvroEvents() { return flux -> flux .concatMap(events -> { Message<List<MyAvroEvents>> typedEvents = (Message<List<MyAvroEvents>>) events; Acknowledgment acknowledgment = typedEvents.getHeaders().get(KafkaHeaders.ACKNOWLEDGMENT, Acknowledgment.class); List<MyAvroEvent> events = typedEvents.getPayload(); // 此列表不为空 // 保存事件到数据库 // 无超时异常时确认,否则不确认 }) .then(); }
调试第一种写法Function<Flux<Message<List<MyAvroEvent>>>, Mono<Void>>时,发现Payload在SimpleFunctionRegistry.java的第1256行处丢失。请问这是Bug还是消费批处理仅允许使用通配符参数?官方文档中实际只提到了通配符。
使用版本:
- spring-cloud-stream:4.0.2
- spring-cloud-stream-binder-kafka:4.0.2
- spring-cloud-function-core:4.0.2
这不是Bug,而是Spring Cloud Function/Stream在处理批消息时的类型转换机制限制。
当使用Function<Flux<Message<List<MyAvroEvent>>>, ...>这种强类型声明时,框架会尝试对Message的Payload进行类型转换,但批处理场景下,框架的类型转换器无法正确将原始批数据映射为List<MyAvroEvent>,导致Payload被清空。
而使用通配符Message<?>时,框架不会主动进行类型转换,保留了原始的Payload数据,因此可以通过强制类型转换获取到正确的事件列表。
官方文档只提及通配符的写法,说明这是当前版本下消费批处理消息的推荐方式。如果需要强类型支持,可考虑在获取Payload后手动进行类型转换,或者升级到更高版本的Spring Cloud组件(后续版本可能修复该类型转换问题)。
内容的提问来源于stack exchange,提问作者jwpol

