Kafka KStream关联Avro SpecificRecords失败问题求助
KStream关联Avro SpecificRecords编译错误解决方案
问题根源:编译器无法正确推断
Joined类的泛型参数,导致方法重载匹配失败;同时未显式指定Serde的泛型类型也会加剧推断歧义。具体修复步骤:
显式指定
Joined的泛型参数
在调用Joined.with()时明确泛型类型,帮助编译器匹配正确的join方法重载:Joined.<String, RecordOne, RecordTwo>with(Serdes.String(), recOneSpecificSerde, recTwoSpecificSerde)确保Serde的泛型类型正确声明
确认你的Avro Serde是带泛型的具体类型,而非原始类型:// 初始化RecordOne的Serde SpecificAvroSerde<RecordOne> recOneSpecificSerde = new SpecificAvroSerde<>(); Map<String, String> avroConfig = new HashMap<>(); avroConfig.put(AbstractKafkaSchemaSerDeConfig.SCHEMA_REGISTRY_URL_CONFIG, "http://your-schema-registry:8081"); recOneSpecificSerde.configure(avroConfig, false); // false表示作为value Serde使用 // 同理初始化RecordTwo的Serde SpecificAvroSerde<RecordTwo> recTwoSpecificSerde = new SpecificAvroSerde<>(); recTwoSpecificSerde.configure(avroConfig, false);明确Lambda表达式的参数类型
显式指定ValueJoiner的参数类型,进一步消除编译器的推断歧义:KStream<String, JoinedRecord> joinedRecord = recOneChangedKeyStream.join(recTwoStream, (RecordOne recOn, RecordTwo recTwo) -> { JoinedRecord jr = new JoinedRecord(); jr.setFieldOne(recOn.getFieldOne()); // 设置其他字段... return jr; }, JoinWindows.of(Duration.ofSeconds(60)), Joined.<String, RecordOne, RecordTwo>with(Serdes.String(), recOneSpecificSerde, recTwoSpecificSerde) );
- 额外注意事项
如果JoinedRecord也是Avro SpecificRecord类型,后续若要将关联后的流输出到主题,需为其配置对应的SpecificAvroSerde<JoinedRecord>,并通过Produced.with()指定,避免运行时序列化错误。
内容的提问来源于stack exchange,提问作者Syed Ali
相关产品推荐
相关产品推荐

