Spring Cloud Stream序列化异常:消费修改后无法发布消息
问题场景
配置Spring Cloud Stream与Kafka集成,消费端开启batch-mode,自定义TransactionDeserializer和EnrichedTransactionSerializer,通过Function实现消息接收与 enrichment。能正常消费消息,但发布时触发序列化异常:
Caused by: org.apache.kafka.common.errors.SerializationException: Can't convert value of class reactor.core.publisher.FluxPeekFuseable to class org.apache.kafka.common.serialization.ByteArraySerializer specified in value.serializer
相关配置与代码
配置信息
spring: application: name: transaction-enricher-application integration: poller: fixed-delay: 5000 cloud: stream: kafka: binder: brokers: broker:9092 # Switch here to local instance when running on localhost bindings: enrichTransaction-in-0: consumer: batch-mode: true configuration: value: deserializer: com.example.demo.serdes.TransactionDeserializer destination: approvalRequest-out-0 enrichTransaction-out-0: producer: useNativeEncoding: true configuration: value: serializer: com.example.demo.serdes.EnrichedTransactionSerializer
原服务类代码
package com.example.demo.enricher; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Bean; import com.example.demo.service.EnrichmentService; import java.util.function.Function; import com.example.demo.domain.EnrichedTransaction; import com.example.demo.domain.Transaction; @Configuration public class CashCardTransactionEnricher { @Bean EnrichmentService enrichmentService() { return new EnrichmentService(); } @Bean public Function<Transaction, EnrichedTransaction> enrichTransaction(EnrichmentService enrichmentService) { return transaction -> { return enrichmentService.enrichTransaction(transaction); }; } }
根本原因
当开启batch-mode: true时,Spring Cloud Stream Kafka消费端会将批量消息封装为**List<Transaction>**(或响应式的Flux<Transaction>)传递给Function,但你定义的Function是Function<Transaction, EnrichedTransaction>,只能处理单个Transaction对象。这种类型不匹配导致框架无法正确解析批量消息,最终将原始的Flux对象直接传给了生产者序列化器,而序列化器只接受EnrichedTransaction类型,因此抛出类型转换异常。
解决方案
修改Function的输入输出类型,适配批量消费的场景:
1. 批量集合处理
将Function的输入改为List<Transaction>,输出改为List<EnrichedTransaction>,对应批量消息的处理逻辑:
package com.example.demo.enricher; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Bean; import com.example.demo.service.EnrichmentService; import java.util.function.Function; import com.example.demo.domain.EnrichedTransaction; import com.example.demo.domain.Transaction; import java.util.List; import java.util.stream.Collectors; @Configuration public class CashCardTransactionEnricher { @Bean EnrichmentService enrichmentService() { return new EnrichmentService(); } @Bean public Function<List<Transaction>, List<EnrichedTransaction>> enrichTransaction(EnrichmentService enrichmentService) { return transactions -> transactions.stream() .map(enrichmentService::enrichTransaction) .collect(Collectors.toList()); } }
2. (可选)响应式批量处理
如果使用响应式编程模型,也可以用Flux作为输入输出:
package com.example.demo.enricher; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Bean; import com.example.demo.service.EnrichmentService; import java.util.function.Function; import com.example.demo.domain.EnrichedTransaction; import com.example.demo.domain.Transaction; import reactor.core.publisher.Flux; @Configuration public class CashCardTransactionEnricher { @Bean EnrichmentService enrichmentService() { return new EnrichmentService(); } @Bean public Function<Flux<Transaction>, Flux<EnrichedTransaction>> enrichTransaction(EnrichmentService enrichmentService) { return flux -> flux.map(enrichmentService::enrichTransaction); } }
3. 配置验证
保持消费者的batch-mode: true配置不变,生产者的useNativeEncoding: true配置正确(该配置确保Spring Cloud Stream直接使用你指定的自定义序列化器,不额外包装)。
额外优化建议
反序列化器中ObjectMapper应改为类成员变量单例化,避免每次反序列化创建新实例,提升性能:
package com.example.demo.serdes; import com.fasterxml.jackson.core.JsonProcessingException; import com.fasterxml.jackson.databind.ObjectMapper; import com.example.demo.domain.Transaction; import org.apache.kafka.common.serialization.Deserializer; import java.util.Map; import com.example.demo.domain.CashCard; public class TransactionDeserializer implements Deserializer<Transaction> { private final ObjectMapper objectMapper = new ObjectMapper(); // 单例初始化 @Override public void configure(Map<String, ?> configs, boolean isKey) {} @Override public Transaction deserialize(String topic, byte[] data) { if (data == null) { return null; } try { return objectMapper.readValue(data, Transaction.class); } catch (JsonProcessingException e) { System.out.println("JsonProcessingException occured while deserializing Transaction" + e.getMessage()); } catch (Exception e) { System.out.println("Error deserializing Transaction: " + e.getMessage()); return new Transaction(0L, new CashCard(0L, "Error", 0.0)); } return new Transaction(1L, new CashCard(1L, "Test Owner", 3.14)); // just to avoid any errors } @Override public void close() {} }
内容的提问来源于stack exchange,提问作者Mostafa Hamid

