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

使用Lambda在KStream中处理Avro格式数据的差集查询问题

Kafka Streams Left Join Type Mismatch: Filter Records Where Key Doesn't Exist in Second Stream

I have two KStreams and need to retrieve records from Stream1 whose keys don't exist in Stream2 by joining them on the key. However, I'm hitting a type mismatch exception when trying to implement a left join, and I need help resolving this.

Stream Data Examples

Stream1

[KSTREAM-MAP-0000000004]: 1, {"id": 1, "name": "john", "age": 26}
[KSTREAM-MAP-0000000004]: 2, {"id": 2, "name": "jane", "age": 24}
[KSTREAM-MAP-0000000004]: 3, {"id": 3, "name": "julia", "age": 25}
[KSTREAM-MAP-0000000004]: 4, {"id": 4, "name": "jamie", "age": 22}
[KSTREAM-MAP-0000000004]: 5, {"id": 5, "name": "jenny", "age": 27}

Stream2

[KSTREAM-MAP-0000000004]: 1, {"id": 1, "name": "xxx", "age": 26}
[KSTREAM-MAP-0000000004]: 2, {"id": 2, "name": "yyy", "age": 24}
[KSTREAM-MAP-0000000004]: 31, {"id": 3, "name": "zzz", "age": 25}
[KSTREAM-MAP-0000000004]: 41, {"id": 4, "name": "uuu", "age": 22}
[KSTREAM-MAP-0000000004]: 51, {"id": 5, "name": "iii", "age": 27}

Expected Output

I want to get only the records from Stream1 where the key has no match in Stream2:

3, {"id": 3, "name": "julia", "age": 25}
4, {"id": 4, "name": "jamie", "age": 22}
5, {"id": 5, "name": "jenny", "age": 27}

Avro Schema Registry Definition

Here's the Avro schema I'm using:

{"namespace": "schema.avro", 
 "type": "record", 
 "name": "mysql", 
 "fields": [ 
     {"name": "id", "type": "int", "doc" : "id"}, 
     {"name": "name", "type": "string", "doc" : "name"}, 
     {"name": "age", "type": "int", "doc" : "age"} 
 ] 
}

My Attempted Code

I tried implementing a left join, but the code throws a type mismatch error:

final Serde<GenericRecord> genericAvroSerde = new GenericAvroSerde();
KStream<Integer,String> joined1 = psql_data.leftJoin(mysql_data,
    (leftValue, rightValue) -> "psql_data=" + leftValue + ", mysql_data=" + rightValue,
    JoinWindows.of(TimeUnit.MINUTES.toMillis(1)),
    Joined.with(
        Serdes.Integer(),
        genericAvroSerde,
        genericAvroSerde)
);

Error Message

The compilation error I'm seeing is:

[ERROR] /home/kafka-connect/confluent-4.1.0/kafka_streaming/src/main/java/com/aail/kafka_stream.java:[140,43] error: no suitable method found for leftJoin(KStream<Integer,mysql>,(leftValue[...]Value,JoinWindows,Joined<Integer,GenericRecord,GenericRecord>)
[ERROR] method KStream.<VO#1,VR#1>leftJoin(KStream<Integer,VO#1>,ValueJoiner<? super mysql,? super VO#1,? extends VR#1>,JoinWindows) is not applicable
[ERROR] (cannot infer type-variable(s) VO#1,VR#1
[ERROR] (actual and formal argument lists differ in length))
[ERROR] method KStream.<VO#2,VR#2>leftJoin(KStream<Integer,VO#2>,ValueJoiner<? super mysql,? super VO#2,? extends VR#2>,JoinWindows,Joined<Integer,mysql,VO#2>) is not applicable
[ERROR] (inferred type does not conform to equality constraint(s)
[ERROR] inferred: GenericRecord
[ERROR] equality constraints(s): GenericRecord,mysql)

My Hypothesis

I think the issue is with the Joined method — I'm using GenericAvroSerde, but I should be using a Serde that matches the generated mysql Avro class instead. However, I'm not sure how to implement this correctly. Can someone help me fix this join operation to get the expected output?

内容的提问来源于stack exchange,提问作者Sujitha Chinnu

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 07:39:22