如何用Java的Mutiny实现带动态延迟的响应式事件发射与轮询
解决方案
核心修改点
- 实现周期性轮询:用Mutiny的
Multi.createBy().repeating()替代单次查询,按设定间隔重复执行数据查询任务 - 修正延迟计算逻辑:原方法计算逻辑有误,改为计算录入时间+1秒与当前时间的差,确保事件在录入时间恰好1秒后发射
- 线程安全处理:用
AtomicReference存储上次查询时间,避免多线程环境下的竞态问题 - 动态延迟绑定:通过
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
相关产品推荐
相关产品推荐

