Mutiny轮询场景下如何向父Multi传播完成事件终止轮询
问题原因
- 你当前使用的
transformToMultiAndMerge算子只有在所有内部子流全部完成时,才会向上游传递完成信号,单轮轮询里的emitter调用complete()仅会结束当前轮次的子流,无法终止上游ticks()生成的无限定时流。 - 你在手动创建的emitter内部单独调用
subscribe()的写法,切断了Mutiny的订阅信号传播链路,完成、取消事件无法回溯到最上游的定时源,导致ticks调度任务会一直持续运行。 - 现有逻辑没有在找到目标交易后给上游发送终止指令,无限长的ticks流没有收到取消信号自然不会停止。
修复实现
直接使用Mutiny原生算子搭建完整的响应式链路,不要手动嵌套订阅、手动管理emitter,保证信号可以沿链路向上传递:
Multi.createFrom().ticks().every(Duration.ofSeconds(5)) .onItem().invoke(tick -> System.out.println("Tick:" + tick)) // 每轮tick触发交易查询,展开返回的交易列表 .onItem().transformToMultiAndMerge(tick -> service.getTransactions() .onItem().transformToMulti(transactions -> Multi.createFrom().iterable(transactions)) ) // 校验流程状态,状态异常直接抛错终止流 .onItem().invoke(transaction -> { if (!verification.isOngoing()) { throw new TransactionVerificationException(); } }) // 过滤匹配条件的目标交易 .select().where(transaction -> transaction.getAmount().stream() .anyMatch(amount -> "test".equals(amount.getQuantity())) ) // 取第一个匹配结果后立刻终止整个流,上游ticks会自动被取消 .select().first() .subscribe() .with( transaction -> log.info("匹配到目标交易:{}", transaction), Throwable::printStackTrace );
关键说明
- 去掉了手动创建emitter、内部手动订阅的非标准写法,整个流的生命周期由Mutiny统一管理,取消、完成、异常信号都可以沿链路向上传递到最上游的ticks定时源。
select().first()算子拿到第一个符合条件的元素后,会立刻向上游所有节点发送取消信号,ticks绑定的定时调度任务会自动停止,不会再产生新的轮询触发事件。- 如果你需要保证轮询串行执行(避免上一轮接口请求未返回、下一轮tick已经触发导致的并发查询问题),将
transformToMultiAndMerge替换为transformToMultiAndConcatenate即可。
内容的提问来源于stack exchange,提问作者Panos Karampis
相关产品推荐
相关产品推荐

