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
相关产品推荐
相关产品推荐

