このトピックでは、Spark を使用して ApsaraDB for HBase Performance-enhanced Edition クラスタにアクセスする方法について説明します。
前提条件
ApsaraDB for HBase Performance-enhanced Edition クラスタのバージョンが 2.4.3 以降であること。 ApsaraDB for HBase Performance-enhanced Edition クラスタのバージョンを表示または更新する方法の詳細については、マイナーバージョンアップデート をご参照ください。
クライアントの IP アドレスが ApsaraDB for HBase Performance-enhanced Edition クラスタのホワイトリストに追加されていること。 ApsaraDB for HBase Performance-enhanced Edition クラスタのホワイトリストにクライアントを追加する方法の詳細については、IP アドレス許可リストとセキュリティグループの設定 をご参照ください。
ApsaraDB for HBase パフォーマンス強化版クラスターのエンドポイント(Java API エンドポイント)は、ApsaraDB for HBase コンソールで確認できます。
使用上の注意
インターネット経由で ApsaraDB for HBase Performance-enhanced Edition クラスタにアクセスするには、データアクセス操作を実行する前に、オープンソースの HBase クライアントを ApsaraDB for HBase クライアントに置き換えます。 詳細については、Java 用 ApsaraDB for HBase SDK のアップグレード をご参照ください。
アプリケーションが Elastic Compute Service(ECS)インスタンスにデプロイされており、仮想プライベートクラウド(VPC)経由で ApsaraDB for HBase Performance-enhanced Edition クラスタにアクセスする場合、ネットワーク接続を確保するために、ApsaraDB for HBase Performance-enhanced Edition クラスタと ECS インスタンスが以下の要件を満たしていることを確認してください。
ApsaraDB for HBase クラスタと ECS インスタンスが同じリージョンにデプロイされていること。 ネットワークレイテンシを削減するために、クラスタとインスタンスを同じゾーンにデプロイすることをお勧めします。
ApsaraDB for HBase Performance-enhanced Edition クラスタと ECS インスタンスが同じ VPC に属していること。
ApsaraDB for HBase Performance-enhanced Edition クラスタへの接続を確立するためのパラメータの設定
方法 1:設定ファイルにアクセス設定を追加する。
hbase-site.xmlという名前の設定ファイルに次の設定項目を追加します。<configuration> <!-- クラスタの Java API エンドポイント。エンドポイントは、ApsaraDB for HBase コンソールの [データベース接続] ページで取得できます。 --> <property> <name>hbase.zookeeper.quorum</name> <value>ld-bp150tns0sjxs****-proxy-hbaseue.hbaseue.rds.aliyuncs.com:30020</value> </property> </configuration>方法 2:Configuration オブジェクトにパラメータを追加する。
// Configuration オブジェクトを作成します。 Configuration conf = HBaseConfiguration.create(); // クラスタの Java API エンドポイント。エンドポイントは、ApsaraDB for HBase コンソールの [データベース接続] ページで取得できます。 conf.set("hbase.zookeeper.quorum", "ld-bp150tns0sjxs****-proxy-hbaseue.hbaseue.rds.aliyuncs.com:30020");
例
test(" Spark SQL の count 結果をテストする") {
// 1. HBaseue シェルのアクセス設定を追加します。
var conf = HBaseConfiguration.create
conf.set("hbase.zookeeper.quorum", "ld-bp150tns0sjxs****-proxy-hbaseue.hbaseue.rds.aliyuncs.com:30020")
// 2. ApsaraDB for HBase テーブルを作成します。
val hbaseTableName = "testTable"
val cf = "f"
val column1 = cf + ":a"
val column2 = cf + ":b"
var rowsCount: Int = -1
var namespace = "spark_test"
val admin = ConnectionFactory.createConnection(conf).getAdmin()
val tableName = TableName.valueOf(namespace, hbaseTableName)
val htd = new HTableDescriptor(tableName)
htd.addFamily(new HColumnDescriptor(cf))
admin.createTable(htd)
// 3. 作成したテーブルにテストデータを挿入します。
val rng = new Random()
val k: Array[Byte] = new Array[Byte](3)
val famAndQf = KeyValue.parseColumn(Bytes.toBytes(column))
val puts = new util.ArrayList[Put]()
var i = 0
for (b1 <- ('a' to 'z')) {
for (b2 <- ('a' to 'z')) {
for (b3 <- ('a' to 'z')) {
if(i < 10) {
k(0) = b1.toByte
k(1) = b2.toByte
k(2) = b3.toByte
val put = new Put(k)
put.addColumn(famAndQf(0), famAndQf(1), ("value_" + b1 + b2 + b3).getBytes())
puts.add(put)
i = i + 1
}
}
}
}
val conn = ConnectionFactory.createConnection(conf)
val table = conn.getTable(tableName)
table.put(puts)
// 4. Spark テーブルを作成します。
val sparkTableName = "spark_hbase"
val createCmd = s"""CREATE TABLE ${sparkTableName} USING org.apache.hadoop.hbase.spark
| OPTIONS ('catalog'=
| '{"table":{"namespace":"$${hbaseTableName}", "name":"${hbaseTableName}"},"rowkey":"rowkey",
| "columns":{
| "col0":{"cf":"rowkey", "col":"rowkey", "type":"string"},
| "col1":{"cf":"cf1", "col":"a", "type":"string"},
| "col2":{"cf":"cf1", "col":"b", "type":"String"}}}'
| )""".stripMargin
println(" createCmd: \n" + createCmd + " rows : " + rowsCount)
sparkSession.sql(createCmd)
// 5. COUNT() 関数を含むステートメントを実行します。
val result = sparkSession.sql("select count(*) from " + sparkTableName)
val sparkCounts = result.collect().apply(0).getLong(0)
println(" sparkCounts : " + sparkCounts)