【问题标题】:Modify collection with ResultSetFuture使用 ResultSetFuture 修改集合
【发布时间】:2018-02-07 11:40:24
【问题描述】:

我尝试使用实现异步查询 java-driver-async-queries。我正在修改 FutureCallback 中的列表,但似乎它不起作用 -

List<Product> products = new ArrayList<Product>();

for (// iterating over a Map) {
    key = entry.getKey();
    String query = "SELECT id,desc,category FROM products where id=?";
    ResultSetFuture future = session.executeAsync(query, key);
    Futures.addCallback(future,
        new FutureCallback<ResultSet>() {
            @Override public void onSuccess(ResultSet result) {
                Row row = result.one();
                if (row != null) {
                    Product product = new Product();
                    product.setId(row.getString("id"));
                    product.setDesc(row.getString("desc"));
                    product.setCategory(row.getString("category"));

                    products.add(product);
                }
            }

            @Override public void onFailure(Throwable t) {
                // log error
            }
        },
        MoreExecutors.sameThreadExecutor()
    );
}

System.out.println("Product List : " + products); // Not printing correct values. Sometimes print blank

还有其他方法吗?

根据我实施的 Mikhail Baksheev 回答,现在得到了正确的结果。 只是一个转折。我需要实现一些额外的逻辑。我想知道我是否可以使用 List&lt;MyClass&gt; 而不是 List&lt;ResultSetFuture&gt; 和 MyClass 作为 -

public class MyClass {

    private Integer         productCount;
    private Integer         stockCount;
    private ResultSetFuture result;
}

然后在迭代时将 FutureList 设置为 -

ResultSetFuture result = session.executeAsync(query, key.get());
MyClass allResult = new MyClass();
allResult.setInCount(inCount);
allResult.setResult(result);
allResult.setSohCount(values.size() - inCount);

futuresList.add(allResult);

【问题讨论】:

  • “不工作”的定义是什么?发布预期行为和实际行为。
  • 你不是在等待你的未来吗?看起来您创建了一堆期货并立即打印结果结构。 Result 结构尚未填充,因为大多数期货仍在进行中。
  • 谢谢。什么是修正。有没有可以参考的代码示例?

标签: java cassandra datastax futuretask


【解决方案1】:

正如@RussS 提到的,代码不会等待所有期货都完成。

有很多方法可以同步异步代码。例如,使用CountDownLatch:

编辑: 也请使用单独的线程进行回调,并使用并发收集产品。

ConcurrentLinkedQueue<Product> products = new ConcurrentLinkedQueue<Product>();
final Executor callbackExecutor = Executors.newSingleThreadExecutor();
final CountDownLatch doneSignal = new CountDownLatch(/*the Map size*/);
for (// iterating over a Map) {
    key = entry.getKey();
    String query = "SELECT id,desc,category FROM products where id=?";
    ResultSetFuture future = session.executeAsync(query, key);
    Futures.addCallback(future,
        new FutureCallback<ResultSet>() {
            @Override public void onSuccess(ResultSet result) {
                Row row = result.one();
                if (row != null) {
                    Product product = new Product();
                    product.setId(row.getString("id"));
                    product.setDesc(row.getString("desc"));
                    product.setCategory(row.getString("category"));

                    products.add(product);
                }
                doneSignal.countDown();

            }

            @Override public void onFailure(Throwable t) {
                // log error
                doneSignal.countDown();
            }
        },
        callbackExecutor
    );
}

doneSignal.await();           // wait for all async requests to finish
System.out.println("Product List : " + products); 

另一种方法是收集列表中的所有期货,并使用 guava 的 Futures.allAsList 将所有结果作为单个未来等待,例如:

List<ResultSetFuture> futuresList = new ArrayList<>( /*Map size*/);
        for (/* iterating over a Map*/) {
            key = entry.getKey();
            String query = "SELECT id,desc,category FROM products where id=?";
            futuresList.add( session.executeAsync( query, key ) );
        }

        ListenableFuture<List<ResultSet>> allFuturesResult = Futures.allAsList( futuresList );
        List<Product> products = new ArrayList<>();
        try {
            final List<ResultSet> resultSets = allFuturesResult.get();
            for ( ResultSet rs : resultSets ) {
                if ( null != rs ) {
                    Row row = rs.one();
                    if (row != null) {
                        Product product = new Product();
                        product.setId(row.getString("id"));
                        product.setDesc(row.getString("desc"));
                        product.setCategory(row.getString("category"));

                        products.add(product);
                    }
                }
            }
        } catch ( InterruptedException | ExecutionException e ) {
            System.out.println(e);
        }
        System.out.println("Product List : " + products);

编辑 2

我想知道我是否可以使用 List 而不是 List 和 MyClass as

技术上是的,但是在这种情况下你不能在Futures.allAsList 中传递List&lt;MyClass&gt; 或者MyClass 应该实现ListenableFuture 接口

【讨论】:

  • 谢谢。我在 cassandra 中进行嵌套查询时遇到问题,并按照您的回复 stackoverflow.com/questions/45471519/…。现在出现间歇性问题,因为有时我的产品列表变空。但是当我在 onSuccess 方法中打印列表时,它会给我结果。如果我使用 await() 那么方法需要很长时间。请提出建议。
  • @Saurabh,我已经更新了我的答案
  • 已接受。我得到了我问过的其他查询的解决方法。如果你能帮助我,只是最后一点。有时(非常罕见)我得到 com.datastax.driver.core.exceptions.BusyPoolException。因此,只需考虑我们是否可以安全地在多请求场景中使用此解决方案。我们需要什么最佳实践?
  • BusyPoolExcption 意味着您有太多的并发请求,并且驱动程序没有资源来处理所有请求。您需要调整驱动程序选项:docs.datastax.com/en/developer/java-driver/3.3/manual/pooling 或限制异步查询的数量stackoverflow.com/questions/30509095/…,如果查询会使 cassandra 集群过载。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2013-05-14
  • 2013-03-17
  • 1970-01-01
  • 2010-11-12
  • 2017-03-26
  • 2016-01-13
  • 1970-01-01
相关资源
最近更新 更多