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

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 改为 KafkaStreamsApp
  • nme_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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 23:18:19