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

使用leftJoin合并KStream Avro主题时的反序列化错误修复求助

Kafka Streams Avro主题LeftJoin问题求助

我是Kafka新手,正在做个人项目:往两个不同的Avro主题写入数据,用leftJoin合并,后续还要把消息输出到KSQL DB(暂未实现)。目前用Kafka Template往Avro主题生产数据,转成KStream合并后,KafkaListener能正常打印原始主题的消息,但合并后的主题没有消息输出,还遇到两个问题:

  • 移除KStream的consumed.with()配置时,抛出默认Key Serde错误;
  • 保留该配置时,抛出反序列化错误。

已经在application.properties和main()的streamConfig中配置了默认序列化/反序列化,但问题仍未解决。想请教:

  1. 如何正确合并这两个Avro主题?
  2. 错误是否由Avro schema导致?
  3. 是否需要改用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 20:50:24