【问题标题】:Creating stream of elements from collection从集合中创建元素流
【发布时间】:2017-09-12 23:02:39
【问题描述】:

我正在使用 Junit 5 动态测试。 我的意图是从集合中创建一个元素流,以将其传递给 JUnit5 进行测试。 但是,使用此代码,我只能运行 1000 条记录。如何使这项工作无缝无阻塞。

    MongoCollection<Document> collection = mydatabase.getCollection("mycoll");
    final List<Document> cache = Collections.synchronizedList(new ArrayList<Document>());

    FindIterable<Document> f = collection.find().batchSize(1000);
    f.batchCursor(new SingleResultCallback<AsyncBatchCursor<Document>>() {

        @Override
        public void onResult(AsyncBatchCursor<Document> t, Throwable thrwbl) {
            t.next(new SingleResultCallback<List<Document>>() {

                @Override
                public void onResult(List<Document> t, Throwable thrwbl) {
                    if (thrwbl != null) {
                        th.set(thrwbl);
                    }
                    cache.addAll(t);
                    latch.countDown();;

                }
            });
        }
    });
    latch.await();
    return cache.stream().map(batch->process(batch));

更新代码

@ParameterizedTest
@MethodSource("setUp")
void cacheTest(MyClazz myclass) throws Exception {
    assertTrue(doTest(myclass));
}
public static MongoClient getMongoClient() {
 // get client here
}

private static Stream<MyClazz> setUp() throws Exception {
    MongoDatabase mydatabase = getMongoClient().getDatabase("test");
    List<Throwable> failures = new ArrayList<>();
    CountDownLatch latch = new CountDownLatch(1);
    List<MyClazz> list = Collections.synchronizedList(new ArrayList<>());
            mydatabase.getCollection("testcollection").find()
            .toObservable().subscribe(
            document -> {
                list.add(process(document));
            },
            throwable -> {
                failures.add(throwable);
            },
            () -> {
                latch.countDown();
            });
    latch.await();
    return list.stream();
}

public boolean doTest(MyClazz myclass) { 
// processing goes here
}
public MyClazz process(Document doc) { 
// doc gets converted to MyClazz
   return MyClazz;
}

即使是现在,我也看到所有数据都已加载,之后会进行单元测试。 我认为这是因为latch.await()。但是,如果我删除它,则可能不会运行任何测试用例,因为数据库可能正在加载集合。

我的用例是:我在 mongo 中有数百万条记录,并且正在使用它们运行某种集成测试用例。将它们全部加载到内存中是不可行的,因此我正在尝试流式解决方案。

【问题讨论】:

    标签: java mongodb unit-testing junit5 mongodb-asyc-driver


    【解决方案1】:

    我认为我不完全理解您的用例,但鉴于您的问题标有java 和mongo-asyc-driver,这个要求肯定是可以实现的:

    从集合中创建一个元素流以将其传递给测试...使这项工作无缝无阻塞

    以下代码:

    • 使用 MongoDB RxJava 驱动程序查询集合
    • 从该集合中创建一个 Rx Observable
    • 订阅Observable
    • 记录例外
    • 标记完成

      CountDownLatch latch = new CountDownLatch(1);
      List<Throwable> failures = new ArrayList<>();
      collection.find()
              .toObservable().subscribe(
              // on next, this is invoked for each document returned by your find call
              document -> {
                  // presumably you'll want to do something here to meet this requirement: "pass it on to test in JUnit5" 
                  System.out.println(document);
              },
              /// on error
              throwable -> {
                  failures.add(throwable);
              },
              // on completion
              () -> {
                  latch.countDown();
              });
      // await the completion event
      latch.await(); 
      

    注意事项:

    • 这需要使用MongoDB RxJava driver(即com.mongodb.rx.client 命名空间中的类...org.mongodb::mongodb-driver-rx Maven 工件)
    • 在您的问题中,您正在调用 collection.find().batchSize(),这清楚地表明您当前没有使用 Rx 驱动程序(因为 batchSize 不能是 Rx 友好的概念:)
    • 以上代码通过v1.4.0的MongoDB RxJava驱动和v1.1.10的io.reactive::rxjava验证

    更新 1:根据您对问题的更改(遵循我的原始答案),您问:“我看到所有数据都已加载,之后进行单元测试.我认为这是因为latch.await()"?我认为您正在[从可观察的流中填充列表,并且只有在可观察的用尽之后您才开始调用doTest()。这种方法涉及 (1) 来自 MongoDB 的流式处理结果; (2) 将这些结果存储在内存中; (3) 为每个存储的结果运行doTest()。如果你真的想一路流式传输,那么你应该在你的 observable 订阅中调用doTest()。例如:

    mydatabase.getCollection("testcollection").find()
            .toObservable().subscribe(
            document -> {
                doTest(process(document));
            },
            throwable -> {
                failures.add(throwable);
            },
            () -> {
                latch.countDown();
            });
    latch.await();
    

    上面的代码将调用doTest(),因为它从 MongoDB 接收每个文档,当整个 observable 用完时,latch 将递减,您的代码将完成。

    【讨论】:

    • 谢谢,这很好。但我认为我不能立即调用 doTest,因为这些是动态测试。我需要能够创建一个 mongo db 结果流并返回它们。有什么办法可以做到吗?如果我的用例有不同的解决方案,我很好。最困难的部分似乎是批处理动态测试。我相信有办法,只是我不知道。
    • 听起来这里可能有些混乱。我的答案中的代码确实“创建了一个mongo db结果流”。如果您随后选择在返回批次之前对它们进行批处理,那么您将不再进行流式传输。或者换一种说法;如果我的答案中的代码不能满足您的目的,那么您实际上并不想要 MongoDB 结果流。
    • 我同意您提供的代码确实创建了 mongo 结果流。它只是不适用于我正在编写的那种junit和那种用例。我想知道我的用例是否有解决方案,因此正在考虑批处理
    猜你喜欢
    • 2013-06-11
    • 1970-01-01
    • 2017-02-28
    • 1970-01-01
    • 2019-11-05
    • 1970-01-01
    • 2016-12-24
    • 1970-01-01
    • 2014-02-16
    相关资源
    最近更新 更多