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

如何用Java的Mutiny实现带动态延迟的响应式事件发射与轮询

解决方案

核心修改点

  1. 实现周期性轮询:用Mutiny的Multi.createBy().repeating()替代单次查询,按设定间隔重复执行数据查询任务
  2. 修正延迟计算逻辑:原方法计算逻辑有误,改为计算录入时间+1秒与当前时间的差,确保事件在录入时间恰好1秒后发射
  3. 线程安全处理:用AtomicReference存储上次查询时间,避免多线程环境下的竞态问题
  4. 动态延迟绑定:通过onItem().delayIt().by()为每个事件绑定动态计算的延迟时长

完整代码实现

import io.smallrye.mutiny.Multi;
import java.time.Duration;
import java.time.LocalDateTime;
import java.util.List;
import java.util.concurrent.atomic.AtomicReference;

public class EventPoller {
    // 用AtomicReference保证多线程环境下的线程安全
    private final AtomicReference<LocalDateTime> timeOfLastQuery = new AtomicReference<>(LocalDateTime.MIN);
    private final EventMapper mapper; // 实体转Event的映射器

    public EventPoller(EventMapper mapper) {
        this.mapper = mapper;
    }

    public Multi<Event> getNewEvents() {
        // 轮询间隔可根据业务需求调整(示例设为5秒)
        Duration pollInterval = Duration.ofSeconds(5);

        return Multi.createBy()
                .repeating()
                // 每次轮询执行的任务:查询新事件并更新上次查询时间
                .supplier(() -> {
                    LocalDateTime lastQuery = timeOfLastQuery.get();
                    List<EventEntity> newEntities = getNewEventsFromDataSource(lastQuery);
                    // 更新为当前时间,确保下次查询只取新增数据
                    timeOfLastQuery.set(LocalDateTime.now());
                    return Multi.createFrom().iterable(newEntities);
                })
                .every(pollInterval)
                // 将每次轮询的结果流拼接为连续的Multi
                .concatenate()
                // 转换为目标Event类型
                .onItem().transform(mapper::toEvent)
                // 为每个事件设置动态延迟,保证在录入时间1秒后发射
                .onItem().delayIt().by(event -> {
                    long delayMillis = millisUntilOneSecondAfterEntry(event.getDate());
                    return Duration.ofMillis(delayMillis);
                });
    }

    // 计算从当前时间到「录入时间+1秒」的延迟时长,最小为0(立即发射)
    private Long millisUntilOneSecondAfterEntry(LocalDateTime timeOfEntry) {
        LocalDateTime targetEmitTime = timeOfEntry.plusSeconds(1);
        Duration delayDuration = Duration.between(LocalDateTime.now(), targetEmitTime);
        return coerceAtLeast(delayDuration.toMillis(), 0);
    }

    // 数据资源查询方法:获取上次查询时间之后的新事件实体
    private List<EventEntity> getNewEventsFromDataSource(LocalDateTime lastQueryTime) {
        // 替换为你的实际查询逻辑,比如数据库查询
        // return eventRepository.findByDateAfter(lastQueryTime);
        return List.of();
    }

    private static long coerceAtLeast(long x, long minimum) {
        return Math.max(x, minimum);
    }

    // 假设的事件实体类
    public static class EventEntity {
        private LocalDateTime date;
        // 其他属性...

        public LocalDateTime getDate() {
            return date;
        }
    }

    // 目标Event类
    public static class Event {
        private LocalDateTime date;
        // 其他属性...

        public LocalDateTime getDate() {
            return date;
        }
    }

    // 实体转Event的映射器接口
    public interface EventMapper {
        Event toEvent(EventEntity entity);
    }
}

关键细节说明

  • 轮询间隔:every(pollInterval)可根据业务实时性需求调整,比如高实时场景设为1秒
  • 动态延迟:delayIt().by()接收函数为每个事件单独计算延迟,确保精准发射时机
  • 线程安全:AtomicReference避免了多线程环境下修改查询时间的竞态问题
  • 流合并兼容:返回的Multi可直接与其他Multi流合并(如Multi.createBy().merging().streams(...)),满足父类合并流的需求

内容的提问来源于stack exchange,提问作者Schallabajzer

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 09:23:35