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

Spring Cloud Stream序列化异常:消费修改后无法发布消息

Spring Cloud Stream Kafka批量消费发布消息序列化异常排查

问题场景

配置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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 15:14:51