
本文介绍在 Quarkus 或 Reactive Java 项目中,如何利用 Mutiny 的 Multi 流对数据库轮询获取的事件进行精准时间对齐——使每个事件在其 timeOfEntry 时间点后恰好 1 秒才被下游消费。核心在于结合 call() 与 onItem().delayIt().by() 实现每项独立、非阻塞的动态延迟。
本文介绍在 quarkus 或 reactive java 项目中,如何利用 mutiny 的 `multi` 流对数据库轮询获取的事件进行精准时间对齐——使每个事件在其 `timeofentry` 时间点后恰好 1 秒才被下游消费。核心在于结合 `call()` 与 `onitem().delayit().by()` 实现每项独立、非阻塞的动态延迟。
在响应式编程中,「按业务时间而非处理时间发射数据」是一个典型且易被误用的场景。例如,你从数据库批量拉取一批带时间戳(LocalDateTime timeOfEntry)的事件,希望它们不是立即发出,而是各自延迟到 timeOfEntry.plusSeconds(1) 这一精确时刻才进入流——模拟“事件真实发生后 1 秒才被系统感知”。Mutiny 并未提供开箱即用的 delayByDuration(Function
以下是推荐实现方案:
private LocalDateTime timeOfLastQuery = LocalDateTime.MIN;
public Multi<event> getNewEvents() {
return Multi.createFrom()
.iterable(getNewEvents(timeOfLastQuery))
.onItem()
.transform(mapper::toEvent) // 先转换为 Event 对象
.call(event -> {
// 计算该 event 应延迟的毫秒数:目标时间 = event.timeOfEntry + 1s,当前时间为 now()
// 延迟 = max(0, (event.timeOfEntry + 1s) - now())
LocalDateTime targetTime = event.timeOfEntry.plusSeconds(1);
long delayMs = Math.max(0, Duration.between(LocalDateTime.now(), targetTime).toMillis());
return Uni.createFrom().nullItem()
.onItem().delayIt().by(Duration.ofMillis(delayMs));
})
.onItem().transformToUni(ignore -> Uni.createFrom().item(event)) // 重新注入原事件
.onItem().transformToMulti(ignore -> Multi.createFrom().item(event)); // 转回 Multi 流(关键!)
}</event>
⚠️ 注意:上述写法存在冗余转换。更简洁、高效且符合 Mutiny 最佳实践的写法是直接使用 onItem().transformToUniAndConcatenate(或 andMerge),避免中间 null 流:
public Multi<event> getNewEvents() {
List<event> rawEvents = getNewEvents(timeOfLastQuery);
timeOfLastQuery = LocalDateTime.now(); // 更新查询时间点(注意线程安全)
return Multi.createFrom()
.iterable(rawEvents)
.onItem()
.transform(mapper::toEvent)
.onItem()
.transformToUniAndConcatenate(event -> {
LocalDateTime scheduledTime = event.timeOfEntry.plusSeconds(1);
long delayMs = Math.max(0, Duration.between(LocalDateTime.now(), scheduledTime).toMillis());
return Uni.createFrom().item(event)
.onItem().delayIt().by(Duration.ofMillis(delayMs));
});
}</event></event>
✅ 关键要点说明:
- transformToUniAndConcatenate 确保每个事件的延迟是串行执行且彼此隔离的(即不会因前一个延迟长而挤压后一个);
- 使用 Duration.between(LocalDateTime.now(), scheduledTime) 计算正向延迟,自动处理“已过期事件”(返回负值 → Math.max(0, ...) 截断为 0,即立即发射);
- 务必在 getNewEvents(...) 调用后立即更新 timeOfLastQuery,否则下次轮询可能漏掉新数据(注意:若该方法被多线程调用,需加锁或改用原子引用);
- 此方案完全非阻塞,不占用 IO 线程,延迟由 Mutiny 内部定时器调度,适合高吞吐场景。
总结:Mutiny 的 call() 和 transformToUni* 系列操作符是实现“每项动态延迟”的黄金组合。它将延迟逻辑下沉至单个元素生命周期内,既保持了 Multi 的流式语义,又满足了事件驱动架构中对时间精度的严苛要求。










