【问题标题】:Datastax Cassandra Driver Asynchronous ResultSet Fetching Does Not WorkDatastax Cassandra 驱动程序异步结果集获取不起作用
【发布时间】:2020-06-06 20:28:56
【问题描述】:

我想加载大行数据,所以我的计划是将语句分成几部分,除以时间戳,然后异步运行。

...
// List to save ResultSets
List<CompletableFuture<AsyncResultSet>> pending = new ArrayList<>();

for(Range range : ranges) {
    System.out.println("Asynchronous execute query will be called soon!");
    pending.add(executeQuery(session, preparedStatement, range));
}

...

private static CompletableFuture<AsyncResultSet> executeQuery(CqlSession session, 
    PreparedStatement preparedStatement, Range range) {

return session
    .executeAsync(preparedStatement.bind()
        .setInstant("startDateTime", range.getStartDateTime().toInstant())
        .setInstant("endDateTime", range.getEndDateTime().toInstant())
        .setPageSize(1000000))
    .toCompletableFuture()
    .whenCompleteAsync((asyncResultSet, throwable) -> {
        if (throwable == null) {
            System.out.println("Range " + range.getStart() + " to " + range.getEnd() + 
                " has " + asyncResultSet.remaining() + " records.");

            fetchResultSet(asyncResultSet, throwable);

            if(asyncResultSet.hasMorePages()) {
                asyncResultSet.fetchNextPage().whenComplete(LoadCassandraAsync::fetchResultSet);
            }
        } else {
            throwable.printStackTrace();
        }
    }, Executors.newFixedThreadPool(4))
    .exceptionally(throwable -> {
        throwable.printStackTrace();
        return null;
    });
}

我将随机获得退出代码 0(不是来自 main 方法),表示它已关闭。或者,在一些获取之后我什么也得不到,就像有一个线程在运行但什么都不做一样。

如果我评论了“获取行”部分,我得到:

...
Asynchronous execute query will be called soon!
Asynchronous execute query will be called soon!
Asynchronous execute query will be called soon!
Asynchronous execute query will be called soon!
Range 2020-02-14 00:00:00+0700 to 2020-02-14 01:00:00+0700 has 102974 records.
Range 2020-02-14 01:00:00+0700 to 2020-02-14 02:00:00+0700 has 98201 records.
Range 2020-02-14 06:00:00+0700 to 2020-02-14 07:00:00+0700 has 104529 records.
Range 2020-02-14 08:00:00+0700 to 2020-02-14 09:00:00+0700 has 105257 records.
...

我认为这意味着executeQuery() 方法运行良好。

我做错了什么?

【问题讨论】:

    标签: asynchronous cassandra datastax resultset


    【解决方案1】:

    根据查询的数量,您可能会耗尽 cassandra 线程 - concurrent_reads(如果我没记错的话,默认数量是 250)。
    如果您查看日志 (/var/log/cassandra/system.log),应该会出现与该问题相关的消息。要解决此问题,请在发送 200 个查询后添加一个人工 Thread.wait。

    【讨论】:

    • “查询”是什么意思?如果我遍历 ResultSet 并在每一行中执行几个 get,它算作 1 个查询吗?
    • 您以异步方式执行查询。这意味着如果您有 1000 个查询,那么所有这些查询都会同时执行。
    • 这是否意味着异步执行大数据是不明智的?我必须加载一年前的历史数据。我的计划是按小时划分,并异步处理。仅供参考,一小时数据包含大约 300 万行。有什么建议吗?
    猜你喜欢
    • 1970-01-01
    • 2016-02-27
    • 2017-01-20
    • 2013-10-31
    • 2017-04-23
    • 2016-03-29
    • 2017-06-23
    • 2019-01-31
    • 1970-01-01
    相关资源
    最近更新 更多