如何反转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
相关产品推荐
相关产品推荐

