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

如何反转Kafka压缩主题中的多对多关系?

这需求太常见了!用Kafka Streams就能完美实现这个多对多关系的反转,我给你捋清楚具体怎么做:

核心思路

原主题是作者ID→关联书籍列表的映射,我们要把它转换成书籍ID→关联作者列表的映射。本质上就是把每个作者记录里的书籍ID逐个拆分,然后针对每本书籍ID,维护对应的作者集合(自动去重,毕竟同一个作者和书籍的关联只需要保留最新状态)。

具体实现步骤

1. 先定义数据模型

首先得把Kafka里的JSON数据转换成代码里的对象,这里用Java举例:

  • 原主题的AuthorBooks类(对应作者的书籍列表):
public class AuthorBooks {
    private List<String> books;
    // 记得加getter、setter、构造器,还要配置JSON序列化/反序列化(比如用Jackson)
}
  • 目标主题的BookAuthors类(对应书籍的作者列表):
public class BookAuthors {
    private List<String> authors;
    // 同样需要序列化相关配置
}

2. 构建Kafka Streams拓扑

核心是用流处理操作拆分、分组、聚合数据:

// 先初始化Streams配置
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "author-book-relation-reverse");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "你的Kafka broker地址:9092");
props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName());
props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, JsonSerde.class.getName()); // 自定义JSON序列化器

// 开始构建处理拓扑
StreamsBuilder builder = new StreamsBuilder();

// 第一步:读取原作者主题的流
KStream<String, AuthorBooks> authorStream = builder.stream("authors");

// 第二步:把每个作者的书籍列表拆成(书籍ID, 作者ID)的键值对
KStream<String, String> bookAuthorStream = authorStream
    .flatMap((authorId, authorBooks) -> 
        authorBooks.getBooks().stream()
            .map(bookId -> new KeyValue<>(bookId, authorId))
            .collect(Collectors.toList())
    );

// 第三步:按书籍ID分组,聚合出对应的作者集合
KTable<String, Set<String>> bookAuthorsTable = bookAuthorStream
    .groupByKey()
    .aggregate(
        () -> new HashSet<>(), // 初始状态:空的作者集合
        (bookId, authorId, currentAuthors) -> {
            currentAuthors.add(authorId);
            return currentAuthors;
        },
        Materialized.as("book-authors-state-store") // 用状态存储维护每本书的作者列表
    );

// 第四步:把聚合后的集合转成List,序列化后写入目标主题
bookAuthorsTable
    .mapValues(authors -> new BookAuthors(new ArrayList<>(authors)))
    .toStream()
    .to("books", Produced.with(Serdes.String(), new JsonSerde<>(BookAuthors.class)));

// 启动流应用
KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();

// 优雅关闭的钩子
Runtime.getRuntime().addShutdownHook(new Thread(streams::close));

3. 处理「书籍移除」的场景

如果原主题的作者记录是全量更新(每次发送的是该作者当前所有书籍,不是增量),那上面的代码会有问题——比如作者删除了某本书,我们需要从对应书籍的作者列表里移除这个作者。这时候得加状态对比的逻辑:

// 把原主题读成KTable,这样拿到的是作者的最新状态
KTable<String, AuthorBooks> authorTable = builder.table("authors");

// 对比作者前后的书籍列表,找出新增/删除的书籍
KStream<String, KeyValue<String, String>> bookAuthorChanges = authorTable
    .toStream()
    .transformValues(() -> new ValueTransformerWithKey<String, AuthorBooks, Iterable<KeyValue<String, String>>>() {
        private KeyValueStore<String, List<String>> previousBooksStore;

        @Override
        public void init(ProcessorContext context) {
            previousBooksStore = (KeyValueStore<String, List<String>>) context.getStateStore("previous-author-books");
        }

        @Override
        public Iterable<KeyValue<String, String>> transform(String authorId, AuthorBooks currentBooks) {
            List<String> previousBooks = previousBooksStore.get(authorId);
            List<String> currentBookList = currentBooks.getBooks();
            List<KeyValue<String, String>> changes = new ArrayList<>();

            // 处理删除的书籍:从对应书籍的作者列表移除当前作者
            if (previousBooks != null) {
                for (String bookId : previousBooks) {
                    if (!currentBookList.contains(bookId)) {
                        changes.add(new KeyValue<>(bookId, "-" + authorId)); // 用前缀标记删除操作
                    }
                }
            }

            // 处理新增的书籍:把当前作者加入对应书籍的列表
            for (String bookId : currentBookList) {
                if (previousBooks == null || !previousBooks.contains(bookId)) {
                    changes.add(new KeyValue<>(bookId, authorId));
                }
            }

            // 更新状态存储,保存当前的书籍列表
            previousBooksStore.put(authorId, currentBookList);
            return changes;
        }

        @Override
        public void close() {}
    }, Materialized.as("previous-author-books"));

// 处理变更,更新书籍的作者集合
KTable<String, Set<String>> bookAuthorsTable = bookAuthorChanges
    .groupByKey()
    .aggregate(
        () -> new HashSet<>(),
        (bookId, change, currentAuthors) -> {
            if (change.startsWith("-")) {
                String authorToRemove = change.substring(1);
                currentAuthors.remove(authorToRemove);
            } else {
                currentAuthors.add(change);
            }
            return currentAuthors;
        },
        Materialized.as("book-authors-state-store")
    );

// 后续写入目标主题的步骤和之前一致
关键注意点
  • 状态存储:生产环境建议用RocksDB作为状态存储,避免应用重启后丢失状态数据;内存存储只适合测试场景。
  • 序列化配置:一定要确保JSON序列化器能正确处理你的数据模型,比如用Jackson的ObjectMapper自定义Serde。
  • 全量vs增量:如果原主题的作者记录是增量更新(只发新增/删除的书籍),那逻辑可以简化,直接处理对应的新增/删除事件即可。

内容的提问来源于stack exchange,提问作者Ted Naleid

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 07:58:51