Cassandra datastax驱动程序ResultSet在多个线程中共享以便快速读取

RRM*_*RRM 2 cassandra datastax-java-driver cassandra-2.0

我在cassandra有很多桌子,超过20亿行并且越来越多.行具有日期字段,并且遵循日期桶模式以限制每一行.

即便如此,我在特定日期的参赛作品超过一百万.

我想尽快读取和处理每一天的行.我正在做的是获取它的实例com.datastax.driver.core.ResultSet并从中获取迭代器并在多个线程之间共享该迭代器.

所以,基本上我想增加读取吞吐量.这是正确的方法吗?如果没有,请建议更好的方法.

And*_*ert 6

不幸的是你不能这样做.原因是ResultSet提供了一个内部分页状态,用于一次检索第1行.

但是你确实有选择.由于我认为您正在进行范围查询(跨多个分区的查询),因此您可以使用策略,使用token指令一次跨令牌范围提交多个查询.在通过无序分区器结果进行分页时记录了一个很好的例子.

java-driver 2.0.10和2.1.5各自提供了一种机制,用于从主机检索令牌范围并拆分它们.在TokenRangeIntegrationTest.java中的java-driver集成测试中有一个如何执行此操作的示例#Should_expose_token_ranges():

    PreparedStatement rangeStmt = session.prepare("SELECT i FROM foo WHERE token(i) > ? and token(i) <= ?");

    TokenRange foundRange = null;
    for (TokenRange range : metadata.getTokenRanges()) {
        List<Row> rows = rangeQuery(rangeStmt, range);
        for (Row row : rows) {
            if (row.getInt("i") == testKey) {
                // We should find our test key exactly once
                assertThat(foundRange)
                    .describedAs("found the same key in two ranges: " + foundRange + " and " + range)
                    .isNull();
                foundRange = range;
                // That range should be managed by the replica
                assertThat(metadata.getReplicas("test", range)).contains(replica);
            }
        }
    }
    assertThat(foundRange).isNotNull();
}
...
private List<Row> rangeQuery(PreparedStatement rangeStmt, TokenRange range) {
    List<Row> rows = Lists.newArrayList();
    for (TokenRange subRange : range.unwrap()) {
        Statement statement = rangeStmt.bind(subRange.getStart(), subRange.getEnd());
        rows.addAll(session.execute(statement).all());
    }
    return rows;
}
Run Code Online (Sandbox Code Playgroud)

您基本上可以生成语句并以异步方式提交它们,上面的示例只是一次迭代一个语句.

另一个选择是使用spark-cassandra-connector,它基本上是在封面下以非常有效的方式完成的.我觉得它很容易使用,你甚至不需要设置一个火花簇来使用它.有关如何使用Java API,请参阅此文档.