如何在Kafka Streams中记录被过滤掉的数据?
解决方案
你可以通过两种方式实现对被过滤记录的日志记录:
方法一:先使用peek()加条件判断记录,再执行过滤
直接在filterNot()之前添加peek()操作,在peek()里判断记录是否符合被过滤的条件(即ExampleProperty为null),符合条件时记录日志即可。这种方式简单直接,适合仅需记录日志的场景:
streamsBuilder.stream(topicName) // 先对要过滤的记录执行日志操作 .peek((k, v) -> { if (v.getExampleProperty() == null) { // 这里替换成你的日志逻辑,比如用日志框架输出 System.out.printf("过滤掉ExampleProperty为null的记录:key=%s, value=%s%n", k, v); } }) .filterNot((k, v) -> v.getExampleProperty() == null) // 后续业务处理逻辑 ...
方法二:用branch()分流处理
如果除了记录日志,你还需要对被过滤的记录做其他操作(比如发送到专门的死信topic),可以用branch()将流分成两个分支:一个分支保留需要继续处理的记录,另一个分支处理被过滤的记录:
// 将原始流拆分为两个分支:分支0是被过滤的记录,分支1是需要保留的记录 KStream<Object, YourValueClass>[] branches = streamsBuilder.stream(topicName) .branch( (k, v) -> v.getExampleProperty() == null, (k, v) -> v.getExampleProperty() != null ); // 处理被过滤的分支,执行日志记录 branches[0].peek((k, v) -> { System.out.printf("过滤掉ExampleProperty为null的记录:key=%s, value=%s%n", k, v); // 若需要,可在这里添加其他操作,比如发送到死信topic // .to("dead-letter-topic"); }); // 继续处理需要保留的分支 branches[1] // 后续业务处理逻辑 ...
内容的提问来源于stack exchange,提问作者Cr4zyTun4
相关产品推荐
相关产品推荐

