使用Lambda在KStream中处理Avro格式数据的差集查询问题
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

