【发布时间】:2020-03-17 11:33:09
【问题描述】:
我正在尝试实现以下行为:
- 定期轮询/生成事件流(持续时间短,例如 1 秒)
- 然后根据某些内部特征对事件进行分组。
- 每组事件立即写入匹配文件(这是我要保持的行为的关键)
- 预计文件将在后续事件中重复用于匹配组(具有相同的密钥),直到它们被密封/轮换
- 如果持续时间较长(例如 5 秒),文件将被密封/轮换,并在使用其他订阅者时采取行动
我编写了以下示例代码来实现上述行为:
private static final Integer EVENTS = 3;
private static final Long SHORTER = 1L;
private static final Long LONGER = 5L;
private static final Long SLEEP = 100000L;
public static void main(final String[] args) throws Exception {
val files = new DualHashBidiMap<Integer, File>();
Observable.just(EVENTS)
.flatMap(num -> Observable.fromIterable(ThreadLocalRandom.current().ints(num).boxed().collect(Collectors.toList())))
.groupBy(num -> Math.abs(num % 2))
.repeatWhen(completed -> completed.delay(SHORTER, TimeUnit.SECONDS))
.map(group -> {
val file = files.computeIfAbsent(group.getKey(), Unchecked.function(key -> File.createTempFile(String.format("%03d-", key), ".txt")));
group.map(Object::toString).toList().subscribe(lines -> FileUtils.writeLines(file, StandardCharsets.UTF_8.name(), lines, true));
return file;
})
.buffer(LONGER, TimeUnit.SECONDS)
.flatMap(Observable::fromIterable)
.distinct(File::getName)
.doOnNext(files::removeValue)
.doOnNext(file -> System.out.println("File - '" + file + "', Lines - " + FileUtils.readLines(file, StandardCharsets.UTF_8)))
.subscribe();
Thread.sleep(SLEEP);
}
虽然它按预期工作(暂时搁置地图访问的线程安全问题,我使用来自commons-collections4 的双向地图只是为了演示),我想知道是否有办法实现纯 RX 形式的上述功能,不依赖外部地图访问?
请注意,在创建组时立即写入文件至关重要,这意味着我们必须使文件在生成的事件组范围之外存在
提前致谢。
【问题讨论】:
标签: java rxjs rx-java reactive-programming rx-java2