【问题标题】:Efficient way to write asynchronously into cassandra using datastax java driver?使用datastax java驱动程序异步写入cassandra的有效方法?
【发布时间】:2017-04-23 04:14:36
【问题描述】:

我正在使用 datastax java driver 3.1.0 连接到 cassandra 集群,我的 cassandra 集群版本是 2.0.10。我正在使用 QUORUM 一致性异步编写。

  public void save(final String process, final int clientid, final long deviceid) {
    String sql = "insert into storage (process, clientid, deviceid) values (?, ?, ?)";
    try {
      BoundStatement bs = CacheStatement.getInstance().getStatement(sql);
      bs.setConsistencyLevel(ConsistencyLevel.QUORUM);
      bs.setString(0, process);
      bs.setInt(1, clientid);
      bs.setLong(2, deviceid);

      ResultSetFuture future = session.executeAsync(bs);
      Futures.addCallback(future, new FutureCallback<ResultSet>() {
        @Override
        public void onSuccess(ResultSet result) {
          logger.logInfo("successfully written");
        }

        @Override
        public void onFailure(Throwable t) {
          logger.logError("error= ", t);
        }
      }, Executors.newFixedThreadPool(10));
    } catch (Exception ex) {
      logger.logError("error= ", ex);
    }
  }

下面是我的CacheStatement类:

public class CacheStatement {
  private static final Map<String, PreparedStatement> cache =
      new ConcurrentHashMap<>();

  private static class Holder {
    private static final CacheStatement INSTANCE = new CacheStatement();
  }

  public static CacheStatement getInstance() {
    return Holder.INSTANCE;
  }

  private CacheStatement() {}

  public BoundStatement getStatement(String cql) {
    Session session = CassUtils.getInstance().getSession();
    PreparedStatement ps = cache.get(cql);
    // no statement cached, create one and cache it now.
    if (ps == null) {
      synchronized (this) {
        ps = cache.get(cql);
        if (ps == null) {
          cache.put(cql, session.prepare(cql));
        }
      }
    }
    return ps.bind();
  }
}

我上面的save 方法将从多个线程调用,我认为BoundStatement 不是线程安全的。顺便说一句,StatementCache 类是线程安全的,如上所示。

  • 因为BoundStatement 不是线程安全的。如果我从多个线程异步编写,我上面的代码会有什么问题吗?
  • 其次,我在addCallback 参数中使用了Executors.newFixedThreadPool(10)。这样可以吗还是会有什么问题?或者我应该使用MoreExecutors.directExecutor。那么这两者有什么区别呢?最好的方法是什么?

以下是我使用 datastax java 驱动程序连接到 cassandra 的连接设置:

Builder builder = Cluster.builder();
    cluster =
        builder
            .addContactPoints(servers.toArray(new String[servers.size()]))
            .withRetryPolicy(new LoggingRetryPolicy(DowngradingConsistencyRetryPolicy.INSTANCE))
            .withPoolingOptions(poolingOptions)
            .withReconnectionPolicy(new ConstantReconnectionPolicy(100L))
            .withLoadBalancingPolicy(
                DCAwareRoundRobinPolicy
                    .builder()
                    .withLocalDc(
                        !TestUtils.isProd() ? "DC2" : TestUtils.getCurrentLocation()
                            .get().name().toLowerCase()).withUsedHostsPerRemoteDc(3).build())
            .withCredentials(username, password).build();

【问题讨论】:

  • 您每次调用 save 时都会创建一个新的线程池,而是创建线程池的静态或实例版本并重用它。
  • 是的,我在阅读以下答案后做到了。我已经在类的顶部将它声明为 final 然后使用它。一般来说,MoreExecutors.directExecutor()threadpool 之间有什么区别?
  • @ChrisLohfink 在回调中使用MoreExecutors.directExecutor()threadpool 有什么真正的好处吗?你能帮我理解一下吗?

标签: java multithreading cassandra datastax-java-driver


【解决方案1】:

我认为你正在做的很好。您可以通过在应用程序启动时准备所有语句来进一步优化,因此您已经缓存了所有内容,因此在“保存”时准备语句不会受到任何性能影响,并且您不会锁定工作流中的任何内容。

BoundStatement 不是线程安全的,但PreparedStatement 是的,并且每次调用getStatement 时都会返回一个新的BoundStatement。事实上,PreparedStatement.bind() 函数实际上是new BoundStatement(ps).bind() 的快捷方式。而且您没有从多个线程访问 same BoundStatement。所以你的代码很好。

相反,对于线程池,您实际上是在每个addCallback 函数上创建一个新线程池。这是对资源的浪费。我不使用这种回调方法,我更喜欢自己管理普通的FutureResultSet,但我在使用MoreExecutors.sameThreadExecutor() 而不是MoreExecutors.directExecutor() 的datastax 文档中看到examples

【讨论】:

  • 是的,我只是在编写异步代码之前才读过那篇文章。sameThreadExecutor 已被 directExecutor 弃用,这就是我问directExecutor 的原因...什么是最好的方法,什么是那么你推荐这里用于线程池的东西吗?
  • 请看this以掌握这个想法。
猜你喜欢
  • 2014-12-03
  • 1970-01-01
  • 1970-01-01
  • 2013-10-12
  • 2018-02-15
  • 2014-03-16
  • 1970-01-01
  • 1970-01-01
  • 2017-07-13
相关资源
最近更新 更多