【发布时间】:2017-09-20 19:32:39
【问题描述】:
假设我有两个 Verticle 正在发现特殊文件名(例如。可以是任何东西)并将它们发布到 事件总线 例如一个是从 REST api 读取名称,另一个是从文件系统读取名称:
ScanRestVerticle.java
/**
* Reads file names through an REST API
*/
public class ScanRestVerticle extends AbstractVerticle {
@Override
public void start() throws Exception {
HttpClientRequest req = vertx.createHttpClient().request(HttpMethod.GET, "BASE_URL", "URL");
req.toObservable()
.flatMap(HttpClientResponse::toObservable)
.lift(unmarshaller(Model.class))
.subscribe(c -> vertx.eventBus().publish("address", c.specialName()));
req.exceptionHandler(Throwable::printStackTrace);
req.end();
}
}
ScanFsVerticle.java
/**
* Reads file names from a file
*/
public class ScanFsVerticle.java extends AbstractVerticle {
@Override
public void start() throws Exception {
StringObservable.byLine(vertx.fileSystem()
.rxReadFile("myFileNames.txt")
.map(Buffer::toString)
.toObservable())
.subscribe(c -> vertx.eventBus().publish("address", c), e -> System.err.println(e.getMessage()));
}
}
这里一切正常,但是,现在,我有一个 Verticle,它结合了来自事件总线的这些名称并将它们打印在 STD.OUT 上:
PrintVerticle.java
public class PrintVerticle extends AbstractVerticle {
@Override
public void start() throws Exception {
vertx.eventBus()
.<JsonObject>consumer("address")
.bodyStream()
.toObservable()
.reduce(new JsonArray(), JsonArray::add)
.subscribe(j -> System.out.println(j.toString()));
}
问题是这里的 reduce 从未真正完成,因为我认为事件总线正在制作无限流。
那么我该如何真正完成这个操作并打印两个verticles发布的名字呢?
注意:我是 vert.x 和 rx 的新手,我可能会遗漏一些内容或出错,所以请不要判断 :)
提前致谢。
编辑:我可以在这里调用scan() 而不是reduce(),以获得中间结果,但我如何获得scan().last()?
【问题讨论】:
-
你对完成的定义是什么?如果要收集 2 个来源发布的每 2 个名称,则需要区分这 2 个事件
-
感谢您的提问!正如您在示例中看到的那样,我想从两个“扫描”垂直中收集所有名称并将它们加入“打印”垂直...当然这是简化的示例。