【问题标题】:Terminate and get results from an event bus终止事件总线并从事件总线获取结果
【发布时间】: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 个事件
  • 感谢您的提问!正如您在示例中看到的那样,我想从两个“扫描”垂直中收集所有名称并将它们加入“打印”垂直...当然这是简化的示例。

标签: java rx-java vert.x


【解决方案1】:

您的假设是正确的,reduce() 运营商依赖 onComplete() 事件来知道在哪里停止收集。
scan() 是累加器,它为源流的每次新发射发出累积值,因此扫描将为您提供中间结果。
这里的问题是,你使用EventBus,从概念上讲,EventBus 是无限流,所以永远不应该调用onComplete()。从技术上讲,当您close() `EventBus 时可能会调用它,但我想您不应该也不想依赖它。

  • 如果你确实想使用 EventBus,你应该以某种方式结束你的 PrintVerticle 流,这样你就会有 onComplete() 事件。例如,您可以使用EventBus 发布一些事件,这些事件发布了两个观察到的顶点中的每一个结束。然后您可以在发出这 2 个事件后将其压缩并在 PrintVerticle 处结束您的流(使用 takeUntil())。
    像这样的:

    Observable vertice1EndStream = vertx.eventBus()
        .<JsonObject>consumer("vertice1_end")
        .bodyStream()
        .toObservable();
    
    Observable vertice2EndStream = vertx.eventBus()
        .<JsonObject>consumer("vertice2_end")
        .bodyStream()
        .toObservable()
    
    Observable bothVerticesEndStream = Observable.zip(vertice1EndStream, vertice2EndStream, (vertice1, vertice1) -> new Object());
    
    vertx.eventBus()
        .<JsonObject>consumer("address")
        .bodyStream()
        .toObservable()
        .takeUntil(bothVerticesEndStream)
        .reduce(new JsonArray(), JsonArray::add)
        .subscribe(j -> System.out.println(j.toString()));
    
  • 另一个选项是直接使用 2 个流而不是 EventBus。将它们合并在一起,然后使用reduce(),因为这两个流是有限的,所以你不会有问题。

【讨论】:

  • 感谢您的建议...我以为我会选择EventBus,但只使用Observables 实际上更简单,我觉得它更“纯粹”。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2020-12-07
  • 1970-01-01
  • 2022-11-11
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多