Pour exécuter des requêtes Spark SQL sur des données stockées dans ApsaraDB for HBase Performance-enhanced Edition, vous devez configurer la connexion entre Spark et le cluster, puis mapper le schéma de la table HBase vers une table Spark. Cette rubrique explique comment réaliser cette configuration et interroger une table HBase avec Spark SQL.
Prérequis
Avant de commencer, assurez-vous de disposer des éléments suivants :
Un cluster ApsaraDB for HBase Performance-enhanced Edition en version 2.4.3 ou ultérieure. Pour vérifier ou mettre à jour la version, consultez la section Mises à jour des versions mineures.
L'adresse IP du client ajoutée à la liste d'autorisation du cluster. Consultez la section Configurer les listes d'autorisation d'adresses IP et les groupes de sécurité.
Le Java API endpoint du cluster, disponible sur la page Database Connection de la console ApsaraDB for HBase.
Notes d'utilisation
Accès Internet : Remplacez le client HBase open source par le client ApsaraDB for HBase avant d'accéder au cluster via Internet. Consultez la section Mettre à niveau le SDK ApsaraDB for HBase pour Java.
-
Accès VPC depuis une instance Elastic Compute Service (ECS) : Assurez-vous que le cluster et l'instance ECS remplissent les deux conditions suivantes :
Ils sont déployés dans la même région. Déployez-les dans la même zone pour réduire la latence réseau.
Ils appartiennent au même Virtual Private Cloud (VPC).
Processus global
Pour utiliser Spark SQL afin d'interroger les données d'un cluster ApsaraDB for HBase Performance-enhanced Edition :
Configurez la connexion au cluster en définissant la propriété
hbase.zookeeper.quorumavec le Java API endpoint.Créez une table HBase et insérez des données de test.
Définissez un catalogue JSON qui mappe le schéma Spark à la table HBase.
Créez une table Spark à l'aide de la définition du catalogue.
Exécutez des requêtes Spark SQL sur la table Spark mappée.
Configurer la connexion
Définissez la propriété hbase.zookeeper.quorum avec le Java API endpoint de votre cluster. Utilisez l'une des méthodes suivantes.
Méthode 1 : Fichier de configuration
Ajoutez les lignes suivantes au fichier hbase-site.xml :
<configuration>
<!--
The Java API endpoint of the cluster.
Get the endpoint on the Database Connection page in the ApsaraDB for HBase console.
-->
<property>
<name>hbase.zookeeper.quorum</name>
<value>ld-bp150tns0sjxs****-proxy-hbaseue.hbaseue.rds.aliyuncs.com:30020</value>
</property>
</configuration>
Méthode 2 : Objet de configuration
Définissez la propriété par programmation dans un objet Configuration :
// Create a Configuration object.
Configuration conf = HBaseConfiguration.create();
// The Java API endpoint of the cluster.
// Get the endpoint on the Database Connection page in the ApsaraDB for HBase console.
conf.set("hbase.zookeeper.quorum", "ld-bp150tns0sjxs****-proxy-hbaseue.hbaseue.rds.aliyuncs.com:30020");
Exemple
L'exemple suivant crée une table HBase, insère des données de test, mappe la table à une table Spark SQL à l'aide d'une définition de catalogue JSON, et exécute une requête COUNT(*).
Les classesHTableDescriptoretHColumnDescriptorsont obsolètes dans Apache HBase 2.x. UtilisezTableDescriptorBuilderetColumnFamilyDescriptorBuilderdans votre code de production.
test(" test the spark sql count result") {
// Step 1: Configure the connection.
var conf = HBaseConfiguration.create
conf.set("hbase.zookeeper.quorum", "ld-bp150tns0sjxs****-proxy-hbaseue.hbaseue.rds.aliyuncs.com:30020")
// Step 2: Create an HBase table.
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)
// Step 3: Insert test data.
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)
// Step 4: Create a Spark table using a catalog JSON definition.
// The catalog maps Spark column names to HBase column families and qualifiers.
// - "table": specifies the HBase namespace and table name.
// - "rowkey": identifies which Spark column maps to the HBase row key.
// - "columns": maps each Spark column to an HBase column family ("cf"),
// column qualifier ("col"), and data type ("type").
// The row key is also listed here as a named column with cf "rowkey".
val sparkTableName = "spark_hbase"
val createCmd = s"""CREATE TABLE ${sparkTableName} USING org.apache.hadoop.hbase.spark
| OPTIONS ('catalog'=
| '{"table":{"namespace":"${namespace}", "name":"${hbaseTableName}"},"rowkey":"rowkey",
| "columns":{
| "col0":{"cf":"rowkey", "col":"rowkey", "type":"string"},
| "col1":{"cf":"f", "col":"a", "type":"string"},
| "col2":{"cf":"f", "col":"b", "type":"string"}}}'
| )""".stripMargin
println("createCmd: \n" + createCmd + " rows : " + rowsCount)
sparkSession.sql(createCmd)
// Step 5: Run a Spark SQL COUNT(*) query.
val result = sparkSession.sql("select count(*) from " + sparkTableName)
val sparkCounts = result.collect().apply(0).getLong(0)
println("sparkCounts : " + sparkCounts)
}
Étapes suivantes
Pour accéder directement au cluster à l'aide de l'API Java, consultez la section Mettre à niveau le SDK ApsaraDB for HBase pour Java.
Pour gérer les versions du cluster, consultez la section Mises à jour des versions mineures.