【问题标题】:HBase Connection InstanceHBase 连接实例
【发布时间】:2016-07-07 17:30:34
【问题描述】:

我有以下代码:

DStream.map {
      _.message()
}.foreachRDD { rdd =>
    rdd.foreachPartition{iter =>
        val conf = HBaseUtils.configureHBase("iemployee")
        val connection = ConnectionFactory.createConnection(conf)
        val table = connection.getTable(TableName.valueOf("""iemployee"""))
        iter.foreach{elem =>
        /* loop through the records in the partition and push them out to the DB */
    }
}

谁能告诉我这里创建的连接对象val connection = ConnectionFactory.createConnection(conf)是否与每个分区中使用的连接对象相同(因为我从不关闭它)还是会为每个分区创建一个新的连接对象?

【问题讨论】:

    标签: scala hadoop apache-spark database-connection hbase


    【解决方案1】:

    每个分区的新连接实例..

    请看下面的code & documentation of Connection Factory。还提到它的调用者有责任关闭连接。

     /**
       * Create a new Connection instance using the passed <code>conf</code> instance. Connection
       * encapsulates all housekeeping for a connection to the cluster. All tables and interfaces
       * created from returned connection share zookeeper connection, meta cache, and connections
       * to region servers and masters.
       * <br>
       * The caller is responsible for calling {@link Connection#close()} on the returned
       * connection instance.
       *
       * Typical usage:
       * <pre>
       * Connection connection = ConnectionFactory.createConnection(conf);
       * Table table = connection.getTable(TableName.valueOf("table1"));
       * try {
       *   table.get(...);
       *   ...
       * } finally {
       *   table.close();
       *   connection.close();
       * }
       * </pre>
       *
       * @param conf configuration
       * @param user the user the connection is for
       * @param pool the thread pool to use for batch operations
       * @return Connection object for <code>conf</code>
       */
      public static Connection createConnection(Configuration conf, ExecutorService pool, User user)
      throws IOException {
        if (user == null) {
          UserProvider provider = UserProvider.instantiate(conf);
          user = provider.getCurrent();
        }
    
        String className = conf.get(ClusterConnection.HBASE_CLIENT_CONNECTION_IMPL,
          ConnectionImplementation.class.getName());
        Class<?> clazz;
        try {
          clazz = Class.forName(className);
        } catch (ClassNotFoundException e) {
          throw new IOException(e);
        }
        try {
          // Default HCM#HCI is not accessible; make it so before invoking.
          Constructor<?> constructor =
            clazz.getDeclaredConstructor(Configuration.class,
              ExecutorService.class, User.class);
          constructor.setAccessible(true);
          return (Connection) constructor.newInstance(conf, pool, user);
        } catch (Exception e) {
          throw new IOException(e);
        }
      }
    }
    

    希望这会有所帮助!

    【讨论】:

    • 非常感谢!
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-02-15
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多