Kafka Streams修改Avro字段值未生效且未生产消息问题排查
Kafka Streams Avro字段转换无输出问题修复
场景说明
原始输入Avro结构
{ "num_enterprise_txn_entity_edition": 1, "num_enterprise_txn_entity_version": 1, "dte_event_occurred": "2022-06-30T18:42:49.533301Z", "nme_creator": "AvroProducer", "nme_event_type": "ReadyToSubmit" }
预期输出Avro结构
需要保留num_enterprise_txn_entity_edition、num_enterprise_txn_entity_version、dte_event_occurred三个字段的原始值,仅修改两个字段:
nme_creator改为KafkaStreamsAppnme_event_type改为NewSubmission
预期输出结构:
{ "num_enterprise_txn_entity_edition": 1, "num_enterprise_txn_entity_version": 1, "dte_event_occurred": "2022-06-30T18:42:49.533301Z", "nme_creator": "KafkaStreamsApp", "nme_event_type": "NewSubmission" }
现有代码的核心错误
- SpecificAvroSerde未配置必填参数:Confluent的SpecificAvroSerde必须配置Schema Registry连接地址等参数,仅指定类名无法完成Avro的序列化/反序列化,会直接导致消费时反序列化失败、生产时序列化失败,流线程异常退出,这是无消息输出的核心原因。
- Avro对象构建方式错误:Avro生成的SpecificRecord类的
newBuilder()是静态方法,代码中先new POCEntity()再调用builder的写法虽然不触发编译错误,但容易导致字段初始化异常;同时代码错误地将dte_event_occurred字段替换为当前时间,不符合需求中保留原始时间值的要求。 - 字段赋值不符合预期:代码中
nme_creator赋值为StreamsApp,和需求要求的KafkaStreamsApp不一致。 - 泛型类型错误:
mapValues中强制将ValueMapper的返回值声明为Object类型,Kafka Streams无法匹配到正确的序列化器,会触发序列化异常。 - 空值风险:filter逻辑中直接调用
value.getNmeEventType(),未做非空判断,遇到null值会直接抛出NPE导致流线程挂掉。 - Serde未显式绑定:调用
stream()和to()方法时未显式指定Serde,仅依赖全局默认配置,在配置不全的场景下很容易出现序列化/反序列化不匹配的问题。 - 冗余代码:存在未使用的import,加载配置文件的输入流未做关闭处理。
修复后可运行代码
import java.io.IOException; import java.io.InputStream; import java.util.HashMap; import java.util.Map; import java.util.Properties; import com.kinsaleins.avro.POCEntity; import io.confluent.kafka.streams.serdes.avro.SpecificAvroSerde; import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.StreamsConfig; import org.apache.kafka.streams.kstream.Consumed; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.Produced; public class StreamsApp { public static void main(String[] args) throws IOException { Properties properties = new Properties(); // 加载配置,使用try-with-resources自动关闭流 try (InputStream in = StreamsApp.class.getClassLoader().getResourceAsStream("kafkastream.properties")) { properties.load(in); } properties.put(StreamsConfig.APPLICATION_ID_CONFIG, "kafka-streams-app"); properties.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); properties.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, SpecificAvroSerde.class); final String inputTopic = properties.getProperty("producer.send.topic"); final String outputTopic = "SubmissionsTopic"; // 初始化Avro Serde,配置Schema Registry参数 SpecificAvroSerde<POCEntity> pocEntitySerde = new SpecificAvroSerde<>(); Map<String, Object> serdeConfig = new HashMap<>(); // 从配置文件中读取schema registry地址,确保kafkastream.properties中配置了schema.registry.url serdeConfig.put("schema.registry.url", properties.getProperty("schema.registry.url")); pocEntitySerde.configure(serdeConfig, false); StreamsBuilder builder = new StreamsBuilder(); // 消费时显式指定Serde KStream<String, POCEntity> firstStream = builder.stream(inputTopic, Consumed.with(Serdes.String(), pocEntitySerde)); firstStream.peek((key, value) -> System.out.println("Consumed original value: " + value)) // 增加空值判断,避免NPE,字符串比较把常量放前面避免空指针 .filter((key, value) -> value != null && "ReadyToSubmit".equals(value.getNmeEventType())) .mapValues(pocEntity -> { // 直接调用静态newBuilder方法,保留原始dte_event_occurred字段值 return POCEntity.newBuilder() .setIdtEnterpriseTxnEntity(pocEntity.getIdtEnterpriseTxnEntity()) .setNumEnterpriseTxnEntityEdition(pocEntity.getNumEnterpriseTxnEntityEdition()) .setNumEnterpriseTxnEntityVersion(pocEntity.getNumEnterpriseTxnEntityVersion()) // 保留原始事件时间,不重新生成 .setDteEventOccurred(pocEntity.getDteEventOccurred()) // 按要求赋值nme_creator .setNmeCreator("KafkaStreamsApp") .setNmeEventType("NewSubmission") .build(); }) .peek((key, value) -> System.out.println("Transformed value to send: " + value)) // 生产时显式指定Serde .to(outputTopic, Produced.with(Serdes.String(), pocEntitySerde)); KafkaStreams kafkaStreams = new KafkaStreams(builder.build(), properties); // 增加优雅关闭逻辑 Runtime.getRuntime().addShutdownHook(new Thread(kafkaStreams::close)); kafkaStreams.start(); } }
额外配置检查项
- 确保
kafkastream.properties中已经正确配置schema.registry.url参数,以及Kafka集群连接、SASL认证(如果开启)等必填配置。 - 确保运行环境中POCEntity的Avro schema和Schema Registry中存储的输入、输出topic的schema兼容。
- 启动后可以通过控制台peek日志判断消费、转换逻辑是否正常执行,若出现异常可以查看Kafka Streams的日志定位具体序列化/反序列化问题。
内容的提问来源于stack exchange,提问作者rb_9999
相关产品推荐
相关产品推荐

