RxJava基于对象时间戳的按信号缓冲实现求助
解决RxJava按时间戳分钟分组的Buffering需求
我完全理解你想要实现的功能——基于事件自身的时间戳按分钟分组,而不是用固定时间间隔触发buffer,并且每组结束后输出分隔信号,最后一组除外。折腾一整天没搞定确实让人头疼,下面给你一个可行的实现方案:
第一步:定义带时间戳的事件类
首先我们需要一个承载时间戳的事件模型,Kotlin和Java版本分别如下:
// Kotlin版本 data class TimedEvent(val timestamp: Long)
// Java版本 public class TimedEvent { private final long timestamp; public TimedEvent(long timestamp) { this.timestamp = timestamp; } public long getTimestamp() { return timestamp; } }
第二步:核心实现逻辑
我们可以利用groupBy按事件所属的分钟分组,再通过concatMapSingle将每个分组转换成列表,最后统一处理输出格式:
Kotlin实现:
import java.time.LocalDateTime import java.time.ZoneOffset import java.time.format.DateTimeFormatter import io.reactivex.rxjava3.core.Observable // 生成事件所属分钟的起始时间戳(作为分组key) fun getMinuteKey(timestamp: Long): Long { return timestamp / (60 * 1000) * 60 * 1000 } fun main() { // 模拟你的测试事件流(转换为时间戳) val testEvents = listOf( "2018-03-26 19:14:13", "2018-03-26 19:14:44", "2018-03-26 19:14:59", "2018-03-26 19:15:23", "2018-03-26 19:15:27", "2018-03-26 19:16:01", "2018-03-26 19:16:04", "2018-03-26 19:16:40", "2018-03-26 19:16:45" ).map { str -> val dt = LocalDateTime.parse(str, DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss")) TimedEvent(dt.toInstant(ZoneOffset.UTC).toEpochMilli()) } // 创建源Observable val source = Observable.fromIterable(testEvents) source.groupBy { event -> getMinuteKey(event.timestamp) } .concatMapSingle { group -> group.toList() } // 将每个分组转为List .toList() // 收集所有分组到列表,方便控制最后一组不输出分隔符 .subscribe { groups -> groups.forEachIndexed { index, group -> // 打印当前分组的所有事件时间 group.forEach { event -> val dt = LocalDateTime.ofEpochMilli(event.timestamp, ZoneOffset.UTC) println(dt.format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"))) } // 不是最后一组的话,输出分隔信号 if (index != groups.size - 1) { println("---- signal ----") } } } }
Java实现:
import io.reactivex.rxjava3.core.Observable; import java.time.LocalDateTime; import java.time.ZoneOffset; import java.time.format.DateTimeFormatter; import java.util.List; public class MinuteBufferingExample { private static long getMinuteKey(long timestamp) { return timestamp / (60 * 1000) * 60 * 1000; } public static void main(String[] args) { // 模拟测试事件 List<String> testTimeStrings = List.of( "2018-03-26 19:14:13", "2018-03-26 19:14:44", "2018-03-26 19:14:59", "2018-03-26 19:15:23", "2018-03-26 19:15:27", "2018-03-26 19:16:01", "2018-03-26 19:16:04", "2018-03-26 19:16:40", "2018-03-26 19:16:45" ); Observable<TimedEvent> source = Observable.fromIterable(testTimeStrings) .map(timeStr -> { LocalDateTime dt = LocalDateTime.parse(timeStr, DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss")); return new TimedEvent(dt.toInstant(ZoneOffset.UTC).toEpochMilli()); }); source.groupBy(MinuteBufferingExample::getMinuteKey) .concatMapSingle(group -> group.toList()) .toList() .subscribe(groups -> { for (int i = 0; i < groups.size(); i++) { List<TimedEvent> group = groups.get(i); for (TimedEvent event : group) { LocalDateTime dt = LocalDateTime.ofEpochMilli(event.getTimestamp(), ZoneOffset.UTC); System.out.println(dt.format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"))); } if (i != groups.size() - 1) { System.out.println("---- signal ----"); } } }); } }
关键逻辑解释
getMinuteKey方法:将任意时间戳转换为对应分钟的起始时间戳(比如19:14:13会转为19:14:00的时间戳),这样同一分钟的事件会被分到同一个分组。groupBy分组:按分钟key将事件流拆分为多个子流,每个子流对应一分钟内的事件。concatMapSingle与toList:将每个分组的事件收集为列表,保证分组顺序和事件的时间顺序一致。- 最后一组处理:通过索引判断是否为最后一个分组,避免在最后一组后输出分隔信号,完全匹配你想要的输出格式。
注意事项
- 这个方案假设事件流是按时间顺序发射的,如果你的事件是乱序的,需要先对事件流按时间戳排序,否则会出现同一分钟的事件被分到不同分组的情况。
- 如果你的事件流是无限流(不会结束),可以去掉
toList(),改用concatMap直接处理每个分组,此时每个分组结束后都会输出分隔信号(因为无限流没有最后一组)。
内容的提问来源于stack exchange,提问作者user163588
相关产品推荐
相关产品推荐

