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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 07:30:02