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

Spring Boot Kafka Streams如何将带Schema的Payload反序列化为POJO

解决Kafka Connect消息转自定义POJO的反序列化问题

问题根源

Topic中的消息并非直接是AccountData的JSON结构,而是包含schema和payload两层的嵌套结构;同时amount是Base64编码的字节(对应Kafka Connect的Decimal类型),trDate是毫秒级时间戳。直接用普通JsonDeserializer<AccountData>无法匹配字段结构,也无法处理特殊类型转换,导致反序列化后字段全为null或默认值。

解决方案步骤

1. 定义外层消息DTO

先创建类匹配Topic消息的整体结构,用于解析外层的schema和payload:

@Data
public class KafkaConnectAccountMessage {
    private JsonNode schema; // 无需解析schema时可忽略,或保留用于后续逻辑
    private JsonNode payload;
}

2. 自定义反序列化器

实现Deserializer<AccountData>,手动解析payload并处理特殊字段转换:

import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.datatype.jsr310.JavaTimeModule;
import org.apache.kafka.common.serialization.Deserializer;
import java.math.BigDecimal;
import java.math.BigInteger;
import java.time.Instant;
import java.time.LocalDateTime;
import java.time.ZoneId;
import java.util.Base64;
import java.util.Map;

public class AccountDataDeserializer implements Deserializer<AccountData> {
    private final ObjectMapper objectMapper = new ObjectMapper();

    @Override
    public void configure(Map<String, ?> configs, boolean isKey) {
        objectMapper.registerModule(new JavaTimeModule()); // 注册Java 8时间模块
    }

    @Override
    public AccountData deserialize(String topic, byte[] data) {
        if (data == null) return null;
        try {
            // 解析外层消息结构
            KafkaConnectAccountMessage message = objectMapper.readValue(data, KafkaConnectAccountMessage.class);
            JsonNode payload = message.getPayload();
            
            AccountData accountData = new AccountData();
            // 映射普通字段
            accountData.setTrId(payload.get("trId").asInt());
            accountData.setFirstName(payload.has("firstName") && !payload.get("firstName").isNull() ? payload.get("firstName").asText() : null);
            accountData.setLastName(payload.has("lastName") && !payload.get("lastName").isNull() ? payload.get("lastName").asText() : null);
            accountData.setAccType(payload.has("accType") && !payload.get("accType").isNull() ? payload.get("accType").asText() : null);
            accountData.setAccNumber(payload.get("accNumber").asInt());
            accountData.setTrType(payload.has("trType") && !payload.get("trType").isNull() ? payload.get("trType").asText() : null);
            
            // 处理amount:Base64解码转BigDecimal
            if (payload.has("amount") && !payload.get("amount").isNull()) {
                String amountBase64 = payload.get("amount").asText();
                byte[] amountBytes = Base64.getDecoder().decode(amountBase64);
                // Kafka Connect Decimal用大端字节序存储,转BigInteger后再转BigDecimal,scale从schema参数取(此处为0)
                BigInteger amountInt = new BigInteger(1, amountBytes);
                accountData.setAmount(new BigDecimal(amountInt).setScale(0));
            }
            
            // 处理trDate:毫秒时间戳转LocalDateTime
            if (payload.has("trDate") && !payload.get("trDate").isNull()) {
                long timestamp = payload.get("trDate").asLong();
                accountData.setTrDate(LocalDateTime.ofInstant(Instant.ofEpochMilli(timestamp), ZoneId.systemDefault()));
            }
            
            return accountData;
        } catch (Exception e) {
            throw new RuntimeException("Failed to deserialize AccountData", e);
        }
    }

    @Override
    public void close() {}
}

3. 更新自定义Serde

替换原Serde中的反序列化器为自定义实现:

import org.apache.kafka.common.serialization.Serdes;
import com.fasterxml.jackson.databind.ser.std.ToStringSerializer;

public class AccountDataSerde extends Serdes.WrapperSerde<AccountData> {
    public AccountDataSerde() {
        super(new JsonSerializer<>(), new AccountDataDeserializer());
    }
}

备选方案(修改Kafka Connect配置)

如果可以调整Kafka Connect配置,在MySQL Connector的配置文件中添加:

value.converter.schemas.enable=false

这样Topic中的消息会直接输出payload的JSON结构,此时可直接用JsonDeserializer<AccountData>,但仍需:

  • 给amount字段添加Jackson自定义反序列化器,处理Base64转BigDecimal
  • 给trDate字段添加@JsonFormat(shape = JsonFormat.Shape.NUMBER_INT)或用JavaTimeModule处理时间戳转LocalDateTime

关键注意事项

  • null值处理:MySQL中的可空字段在payload中可能为null,需判断后再取值,避免NullPointerException
  • Decimal的scale:必须和Kafka Connect schema中定义的scale参数一致(示例中为0),否则数值会错误
  • 时区一致性:时间戳转LocalDateTime时要确保时区和业务需求一致,避免时间偏移

内容的提问来源于stack exchange,提问作者Coder

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 07:17:04