RRM*_*RRM 2 cassandra datastax-java-driver cassandra-2.0
我在cassandra有很多桌子,超过20亿行并且越来越多.行具有日期字段,并且遵循日期桶模式以限制每一行.
即便如此,我在特定日期的参赛作品超过一百万.
我想尽快读取和处理每一天的行.我正在做的是获取它的实例com.datastax.driver.core.ResultSet并从中获取迭代器并在多个线程之间共享该迭代器.
所以,基本上我想增加读取吞吐量.这是正确的方法吗?如果没有,请建议更好的方法.
不幸的是你不能这样做.原因是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,请参阅此文档.