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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 05:54:16