使用leftJoin合并KStream Avro主题时的反序列化错误修复求助
我是Kafka新手,正在做个人项目:往两个不同的Avro主题写入数据,用leftJoin合并,后续还要把消息输出到KSQL DB(暂未实现)。目前用Kafka Template往Avro主题生产数据,转成KStream合并后,KafkaListener能正常打印原始主题的消息,但合并后的主题没有消息输出,还遇到两个问题:
- 移除KStream的
consumed.with()配置时,抛出默认Key Serde错误; - 保留该配置时,抛出反序列化错误。
已经在application.properties和main()的streamConfig中配置了默认序列化/反序列化,但问题仍未解决。想请教:
- 如何正确合并这两个Avro主题?
- 错误是否由Avro schema导致?
- 是否需要改用JSON?但消息值包含多字段,希望用schema管理。
消息格式示例:{Key : Value} = {company : {inventory_id, company, color, inventory}} = {Toyota : {0, RAV4, 50,000}}
相关代码片段
application.properties
spring.kafka.bootstrap-servers=localhost:9092 spring.kafka.schema-registry-url=http://localhost:8081 # 生产者配置 spring.kafka.producer.key-serializer=org.apache.kafka.common.serialization.StringSerializer spring.kafka.producer.value-serializer=io.confluent.kafka.serializers.KafkaAvroSerializer spring.kafka.producer.properties.schema.registry.url=${spring.kafka.schema-registry-url} # 消费者配置 spring.kafka.consumer.group-id=inventory-group spring.kafka.consumer.auto-offset-reset=earliest spring.kafka.consumer.key-deserializer=org.apache.kafka.common.serialization.StringDeserializer spring.kafka.consumer.value-deserializer=io.confluent.kafka.serializers.KafkaAvroDeserializer spring.kafka.consumer.properties.schema.registry.url=${spring.kafka.schema-registry-url} spring.kafka.consumer.properties.specific.avro.reader=true # Streams配置 spring.kafka.streams.application-id=inventory-streams-app spring.kafka.streams.properties.schema.registry.url=${spring.kafka.schema-registry-url} spring.kafka.streams.properties.default.key.serde=org.apache.kafka.common.serialization.Serdes$StringSerde spring.kafka.streams.properties.default.value.serde=io.confluent.kafka.streams.serdes.avro.SpecificAvroSerde spring.kafka.streams.properties.specific.avro.reader=true
InventoryStreamsApplication.java(原核心代码)
package com.example.inventorystreams; import com.example.inventorystreams.avro.Inventory; import com.example.inventorystreams.avro.Price; 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.JoinWindows; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.Produced; import org.springframework.boot.SpringApplication; import org.springframework.boot.autoconfigure.SpringBootApplication; import org.springframework.context.annotation.Bean; import org.springframework.kafka.annotation.EnableKafkaStreams; import org.springframework.kafka.support.serializer.JsonSerde; import java.time.Duration; import java.util.Properties; @SpringBootApplication @EnableKafkaStreams public class InventoryStreamsApplication { public static void main(String[] args) { SpringApplication.run(InventoryStreamsApplication.class, args); } @Bean public KafkaStreams kafkaStreams() { Properties streamsConfig = new Properties(); streamsConfig.put(StreamsConfig.APPLICATION_ID_CONFIG, "inventory-streams-app"); streamsConfig.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); streamsConfig.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); streamsConfig.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, io.confluent.kafka.streams.serdes.avro.SpecificAvroSerde.class); streamsConfig.put("schema.registry.url", "http://localhost:8081"); streamsConfig.put("specific.avro.reader", true); StreamsBuilder builder = new StreamsBuilder(); KStream<String, Inventory> inventoryStream = builder.stream("inventory-topic", Consumed.with(Serdes.String(), new JsonSerde<>(Inventory.class))); KStream<String, Price> priceStream = builder.stream("price-topic", Consumed.with(Serdes.String(), new JsonSerde<>(Price.class))); KStream<String, String> joinedStream = inventoryStream.leftJoin(priceStream, (inventory, price) -> "Inventory: " + inventory.getInventoryId() + ", Price: " + (price != null ? price.getPrice() : "N/A"), JoinWindows.of(Duration.ofMinutes(5)) ); joinedStream.to("joined-inventory-price-topic", Produced.with(Serdes.String(), Serdes.String())); return new KafkaStreams(builder.build(), streamsConfig); } }
问题分析与解决方案
1. 核心问题:Serde配置不匹配
你当前代码中用JsonSerde读取Avro格式的消息,这直接导致反序列化失败,KStream无法读取原始主题数据,合并后的主题自然没有输出。另外,默认Serde配置因代码中显式指定错误Serde而未生效。
2. 具体修复步骤
步骤1:替换错误的Serde为Avro专用Serde
生产的是Avro数据,读取时必须使用SpecificAvroSerde,修改KStream创建代码:
@Bean public KStream<String, String> kStream(StreamsBuilder builder) { // 初始化Avro Serde并配置schema registry地址 SpecificAvroSerde<Inventory> inventorySerde = new SpecificAvroSerde<>(); SpecificAvroSerde<Price> priceSerde = new SpecificAvroSerde<>(); Map<String, String> serdeConfig = Collections.singletonMap("schema.registry.url", "http://localhost:8081"); inventorySerde.configure(serdeConfig, false); // false表示作为value serde priceSerde.configure(serdeConfig, false); // 使用正确的Serde创建流 KStream<String, Inventory> inventoryStream = builder.stream("inventory-topic", Consumed.with(Serdes.String(), inventorySerde)); KStream<String, Price> priceStream = builder.stream("price-topic", Consumed.with(Serdes.String(), priceSerde)); // 执行LeftJoin并输出结果 KStream<String, String> joinedStream = inventoryStream.leftJoin(priceStream, (inventory, price) -> String.format("库存ID: %d, 品牌: %s, 库存数量: %d, 价格: %s", inventory.getInventoryId(), inventory.getCompany(), inventory.getInventory(), price != null ? String.valueOf(price.getPrice()) : "无数据"), JoinWindows.of(Duration.ofMinutes(5)) ); joinedStream.to("joined-inventory-price-topic"); return joinedStream; }
步骤2:统一Serde配置逻辑
Spring Boot的@EnableKafkaStreams会自动读取application.properties中的配置,无需手动创建KafkaStreams bean。用上述方式直接基于StreamsBuilder定义拓扑,可避免重复配置导致的冲突。
步骤3:确保Join的Key一致性
LeftJoin生效的前提是两个流的Key完全匹配(比如都用company字段作为Key)。生产数据时要确保Key设置正确:
// 生产Inventory消息时指定Key为品牌名 kafkaTemplate.send("inventory-topic", inventory.getCompany(), inventory); // 生产Price消息时指定Key为品牌名 kafkaTemplate.send("price-topic", price.getCompany(), price);
3. Avro vs JSON选择建议
完全不需要改用JSON,Avro更适合多字段消息的schema管理,能保证数据一致性并支持schema演进。你遇到的问题是Serde配置错误,而非Avro本身的问题。
4. 调试技巧
- 开启详细日志排查反序列化错误:
logging.level.org.apache.kafka.streams=DEBUG logging.level.io.confluent=DEBUG
- 用命令行工具验证原始主题数据:
kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic inventory-topic --from-beginning \ --property key.deserializer=org.apache.kafka.common.serialization.StringDeserializer \ --property value.deserializer=io.confluent.kafka.serializers.KafkaAvroDeserializer \ --property schema.registry.url=http://localhost:8081 \ --property specific.avro.reader=true
内容的提问来源于stack exchange,提问作者Arjun

