Mutiny处理多个Uni并使用外部值 如何实现Kafka位点校验的纯响应式改造
纯响应式改造方案
这个场景完全可以实现纯响应式改造,复杂度不高,不需要保留同步阻塞的实现。
实现思路
- 以
consumer.getPositions()返回的Uni<Map<TopicPartition, Long>>作为链式调用的起点,使用flatMap承接positions返回后的后续逻辑 - 将positions的每一条entry转换为一个独立的校验
Uni<Boolean>:每个entry对应调用consumer.committed(partition)获取提交位点,和当前entry的position做比对 - 用
Uni.combine().all().unis()合并所有分区的校验结果,最终通过短路与操作汇总为全局的校验结果 - 全程无阻塞等待,所有异步操作通过响应式流的回调串联,最终返回
Uni<Boolean>符合需求
改造后代码
private Uni<Boolean> isAllProcessedForChannel(String channel) { KafkaConsumer<Object, Object> consumer = clientService.getConsumer(channel); // 第一步:获取所有分区的当前position return consumer.getPositions() .flatMap(positions -> { // 把每个entry转成单个分区的校验Uni List<Uni<Boolean>> partitionChecks = positions.entrySet().stream() .map(entry -> { TopicPartition partition = entry.getKey(); Long expectedPosition = entry.getValue(); // 单个分区的校验逻辑 return consumer.committed(partition) .map(committedMap -> { OffsetAndMetadata offsetAndMetadata = committedMap.get(partition); if (offsetAndMetadata == null) { // 无提交位点的情况默认符合要求,和原有逻辑对齐 return true; } long committedOffset = offsetAndMetadata.offset(); log.info("Offset/Position ({}): {}/{}", partition.partition(), committedOffset, expectedPosition); return committedOffset == expectedPosition; }); }) .toList(); // 合并所有校验结果,全部为true则返回true,否则false return Uni.combine().all().unis(partitionChecks) .combinedWith(results -> results.stream().allMatch(Boolean.TRUE::equals)); }); }
短路优化版本
如果需要提前终止校验(只要有一个分区校验不通过就立刻返回结果,不需要等待其他分区校验完成),可以改用Multi实现,性能更高:
private Uni<Boolean> isAllProcessedForChannelShortCircuit(String channel) { KafkaConsumer<Object, Object> consumer = clientService.getConsumer(channel); return consumer.getPositions() .onItem().transformToMulti(positions -> Multi.createFrom().iterable(positions.entrySet())) .onItem().transformToUniAndMerge(entry -> { TopicPartition partition = entry.getKey(); Long expectedPosition = entry.getValue(); return consumer.committed(partition) .map(committedMap -> { OffsetAndMetadata om = committedMap.get(partition); if (om == null) return true; long offset = om.offset(); log.info("Offset/Position ({}): {}/{}", partition.partition(), offset, expectedPosition); return offset == expectedPosition; }); }) // 只要有一个校验不通过就立刻终止流返回结果 .anyMatch(isValid -> !isValid) // 取反,没有不通过的分区即为全部校验通过 .map(hasInvalid -> !hasInvalid); }
内容的提问来源于stack exchange,提问作者Simon Wick
相关产品推荐
相关产品推荐

