如何在Java版Kafka Streams中对AVRO格式数据执行差集JOIN?
Kafka Streams 左排除连接实现指南
这个需求刚好对应Kafka Streams里的左排除连接(Left Anti Join)——这正是专门用来提取左流(STREAM-1)中在右流(STREAM-2)找不到匹配记录的操作,完全符合你要“仅保留STREAM-1未出现在STREAM-2中的记录”的诉求,匹配依据是两条记录的id字段相等。
核心原理
左排除连接会遍历左流的每一条记录,仅当右流中不存在相同id的记录时,才会将左流的这条记录保留到结果中。和普通左连接不同,它不会返回任何右流的匹配数据,只筛选左流的“独苗”记录。
Java 实现步骤
假设你已经通过Avro工具(比如avro-maven-plugin)生成了对应两个流的Java实体类Stream1Record和Stream2Record,下面是完整实现流程:
1. 初始化基础流与配置
首先读取两个Kafka主题的AVRO数据流,同时注意:原流的key是null,我们需要将流的key替换为id字段(转为String类型),这样才能基于id进行连接匹配:
import org.apache.kafka.common.serialization.Serdes; import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.StreamsBuilder; import org.apache.kafka.streams.kstream.Consumed; import org.apache.kafka.streams.kstream.JoinWindows; import org.apache.kafka.streams.kstream.Joined; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.Produced; import org.apache.kafka.streams.serdes.avro.SpecificAvroSerde; import java.time.Duration; import java.util.Properties; public class StreamAntiJoinExample { public static void main(String[] args) { // 1. 配置Kafka Streams参数 Properties streamsConfig = getKafkaStreamsConfig(); // 2. 初始化StreamsBuilder StreamsBuilder builder = new StreamsBuilder(); // 3. 读取STREAM-1主题,反序列化为KStream SpecificAvroSerde<Stream1Record> stream1Serde = new SpecificAvroSerde<>(); stream1Serde.configure(getAvroConfig(), false); KStream<String, Stream1Record> stream1 = builder.stream( "stream-1-topic", Consumed.with(Serdes.String(), stream1Serde) ); // 4. 读取STREAM-2主题,反序列化为KStream SpecificAvroSerde<Stream2Record> stream2Serde = new SpecificAvroSerde<>(); stream2Serde.configure(getAvroConfig(), false); KStream<String, Stream2Record> stream2 = builder.stream( "stream-2-topic", Consumed.with(Serdes.String(), stream2Serde) ); // 5. 将两个流的key替换为id字段(转为String),用于连接匹配 KStream<String, Stream1Record> stream1WithIdAsKey = stream1 .selectKey((key, record) -> String.valueOf(record.getId())); KStream<String, Stream2Record> stream2WithIdAsKey = stream2 .selectKey((key, record) -> String.valueOf(record.getId()));
2. 执行左排除连接
如果你的Kafka Streams版本 >= 2.4,可以直接用leftAntiJoin方法(更简洁直观);如果版本较低,可以用leftJoin+filter的组合实现:
// 方式1:使用leftAntiJoin(推荐,Kafka Streams >=2.4) KStream<String, Stream1Record> resultStream = stream1WithIdAsKey.leftAntiJoin( stream2WithIdAsKey, JoinWindows.of(Duration.ofDays(1)), // 定义足够大的窗口覆盖所有可能的匹配 Joined.with( Serdes.String(), stream1Serde, stream2Serde ) ); // 方式2:低版本兼容(Kafka Streams <2.4) // KStream<String, Stream1Record> resultStream = stream1WithIdAsKey.leftJoin( // stream2WithIdAsKey, // (stream1Rec, stream2Rec) -> stream1Rec, // 右流无匹配时返回左流记录 // JoinWindows.of(Duration.ofDays(1)), // Joined.with(Serdes.String(), stream1Serde, stream2Serde) // ).filter((key, rec) -> rec != null); // 过滤掉右流有匹配的记录
3. 输出结果到目标主题
最后将筛选后的结果输出到指定主题:
// 6. 将结果输出到目标主题 resultStream.to( "stream-result-topic", Produced.with(Serdes.String(), stream1Serde) ); // 7. 启动Kafka Streams应用 KafkaStreams streams = new KafkaStreams(builder.build(), streamsConfig); streams.start(); // 注册关闭钩子,优雅停止应用 Runtime.getRuntime().addShutdownHook(new Thread(streams::close)); } // 辅助方法:获取Kafka Streams基础配置 private static Properties getKafkaStreamsConfig() { Properties props = new Properties(); props.put("bootstrap.servers", "localhost:9092"); props.put("application.id", "stream-anti-join-example"); props.put("default.key.serde", Serdes.String().getClass().getName()); props.put("default.value.serde", SpecificAvroSerde.class.getName()); return props; } // 辅助方法:获取Avro序列化配置(Schema Registry地址) private static Properties getAvroConfig() { Properties avroProps = new Properties(); avroProps.put("schema.registry.url", "http://localhost:8081"); return avroProps; } }
结果验证
按照上述代码实现后,输出流会自动过滤掉STREAM-1中id为1、2的记录(这两个id在STREAM-2中存在),最终输出和你给出的预期结果完全一致:
[KSTREAM-SOURCE-0000000000]: null, {"id": 3, "name": "julia", "age": 25, "updated_at": 1525774480827} [KSTREAM-SOURCE-0000000000]: null, {"id": 4, "name": "jamie", "age": 22, "updated_at": 1525774480875} [KSTREAM-SOURCE-0000000000]: null, {"id": 5, "name": "jenny", "age": 27, "updated_at": 1525774482927} [KSTREAM-SOURCE-0000000000]: null, {"id": 6, "name": "kishore", "age": 27, "updated_at": 1525775063908} [KSTREAM-SOURCE-0000000000]: null, {"id": 7, "name": "purna", "age": 27, "updated_at": 1525775072006} [KSTREAM-SOURCE-0000000000]: null, {"id": 8, "name": "xxx", "age": 10, "updated_at": 1525783464123} [KSTREAM-SOURCE-0000000000]: null, {"id": 9, "name": "yyy", "age": 10, "updated_at": 1525783667644} [KSTREAM-SOURCE-0000000000]: null, {"id": 10, "name": "zzz", "age": 10, "updated_at": 1525783741814}
内容的提问来源于stack exchange,提问作者bunny sunny
相关产品推荐
相关产品推荐

