Hive Warehouse Connector (HWC) enables Spark applications to access and operate on Hive tables, including Hive managed transactional tables.
You can use HWC with Spark Shell, Spark Submit, and PySpark to run Hive queries and perform table operations.
The HWC installation includes connector JARs for Spark 2 and Spark 3 and a Python package for PySpark.
The HWC files are available in the following directory:
/usr/odp/<ODP_VERSION>/hive_warehouse_connector/
For example:
hive-warehouse-connector-spark2-assembly-1.0.0.jar
hive-warehouse-connector-spark3-assembly-1.0.0.jar
pyspark3_hwc-1.0.0.zip
pyspark_hwc-1.0.0.zip
Replace <ODP_VERSION> with the ODP version installed on your cluster.
To use HWC, specify the HiveServer2 JDBC URL, Kerberos principal, required Spark configurations, and the HWC JAR when you start Spark.
For example:
spark-shell \
--master yarn \
--conf spark.sql.hive.hiveserver2.jdbc.url.principal="hive/_HOST@EXAMPLE.COM" \
--conf spark.datasource.hive.warehouse.read.mode=JDBC_CLIENT \
--conf spark.sql.extensions="com.acceldata.spark.sql.rule.Extensions" \
--conf spark.kryo.registrator=com.qubole.spark.hiveacid.util.HiveAcidKyroRegistrator \
--conf spark.sql.hive.hiveserver2.jdbc.url="<HIVE_JDBC_URL>" \
--jars /usr/odp/<ODP_VERSION>/hive_warehouse_connector/hive-warehouse-connector-spark3-assembly-1.0.0.jar
Replace <HIVE_JDBC_URL> with the HiveServer2 JDBC URL for your environment.
After Spark Shell starts, import HiveWarehouseSession:
import com.acceldata.hwc.HiveWarehouseSession
import com.acceldata.hwc.HiveWarehouseSession._
Create the Hive Warehouse session:
val hive = HiveWarehouseSession.session(spark).build()
Use the hive session to run Hive operations.
Run:
hive.executeQuery(
"CREATE EXTERNAL TABLE <table_name> " +
"(<column_names> <data_type>) " +
"PARTITIONED BY (<partition_name> STRING) " +
"LOCATION '/tmp/<location>/'"
)
Run:
hive.executeQuery(
"CREATE TABLE <table_name> " +
"(<column_names> <data_type>) " +
"TBLPROPERTIES(" +
"'transactional'='true', " +
"'transactional_properties'='default')"
)
Run:
hive.executeQuery("SHOW TABLES").show()
Run:
hive.executeQuery(
"ALTER TABLE <old_table_name> RENAME TO <new_table_name>"
)
For a partitioned table, run:
hive.executeQuery(
"INSERT INTO <table_name> " +
"PARTITION(<partition_name>='<partition_value>') " +
"VALUES(<values>)"
)
HWC supports operations on partitioned Hive tables.
For example:
hive.executeQuery(
"CREATE EXTERNAL TABLE tp13 " +
"(id INT, name STRING) " +
"PARTITIONED BY (dt STRING) " +
"LOCATION '/tmp/tc12/'"
)
Verify that the table was created:
hive.executeQuery("SHOW TABLES").show()
For example:
hive.executeQuery(
"INSERT INTO tp13 PARTITION(dt='2002') VALUES(10,'ad10')"
)
hive.executeQuery(
"INSERT INTO tp13 PARTITION(dt='2003') VALUES(10,'ad10')"
)
Verify the data:
hive.executeQuery("SELECT * FROM tp13").show()
Run:
hive.executeQuery(
"ALTER TABLE default.tp13 RENAME TO default.tp14"
)
Verify the renamed table:
hive.showTables().show(100)
Verify the data in the partition:
hive.executeQuery(
"SELECT * FROM tp14 WHERE dt='2006'"
).show()
Drop the partition:
hive.executeQuery(
"ALTER TABLE default.tp14 DROP PARTITION (dt='2006')"
)
Verify that the partition was removed:
hive.executeQuery("SELECT * FROM tp14").show()
Run:
hive.executeQuery(
"ALTER TABLE tp14 ADD COLUMNS (s STRING, t TIMESTAMP)"
)
Verify the table:
hive.executeQuery("SELECT * FROM tp14").show()
For example:
hive.executeQuery(
"INSERT INTO tp14 PARTITION(dt='2010') " +
"VALUES(10,'ad10','four',current_timestamp())"
)
Verify the data:
hive.executeQuery("SELECT * FROM tp14").show(false)
Use REPLACE COLUMNS to change the table columns.
For example:
hive.executeQuery(
"ALTER TABLE tp14 REPLACE COLUMNS(id INT, name STRING, s STRING)"
)
Verify the table:
hive.executeQuery("SELECT * FROM tp14").show(false)
Run:
hive.executeQuery(
"DESCRIBE FORMATTED tp14"
).show(false)
You can create, modify, and delete Hive views through HWC.
Run:
hive.executeQuery(
"CREATE VIEW tp_view AS SELECT * FROM tp14"
)
Query the view:
hive.executeQuery("SELECT * FROM tp_view").show()
Run:
hive.executeQuery(
"ALTER VIEW tp_view AS SELECT id, name FROM tp14"
)
Verify the view:
hive.executeQuery("SELECT * FROM tp_view").show()
Run:
hive.executeQuery("DROP VIEW tp_view")
Verify that the view was removed:
hive.showTables().show(100)
To remove all rows from a table without deleting the table, run:
hive.executeQuery("TRUNCATE TABLE tp15")
Verify that the table is empty:
hive.executeQuery("SELECT * FROM tp15").show()
You can use the VALUES clause to construct rows directly in a query.
For example:
hive.executeQuery(
"SELECT * FROM (VALUES(4,5,6),(7,8,9)) AS t"
).show()
You can also specify values of different data types:
hive.executeQuery(
"SELECT * FROM " +
"(VALUES(1 AS c1, true AS c2, 'abc' AS c3)," +
"(100,false,'xyz')) AS t"
).show()
HWC supports operations on Hive transactional tables, including insert, update, delete, compaction, rename, and merge operations.
Run:
hive.executeQuery(
"CREATE TABLE tm11(a INT, b INT) " +
"TBLPROPERTIES(" +
"'transactional'='true', " +
"'transactional_properties'='default')"
)
Verify that the table was created:
hive.showTables().show(100)
Run:
hive.executeQuery(
"INSERT INTO tm11 VALUES(10,10),(20,20)"
)
Verify the data:
hive.executeQuery("SELECT * FROM tm11").show(false)
Run:
hive.executeQuery(
"UPDATE tm11 SET a=15 WHERE a=10"
)
Verify the updated data:
hive.executeQuery("SELECT * FROM tm11").show(false)
Run:
hive.executeQuery(
"DELETE FROM tm11 WHERE a=15"
)
Verify that the row was deleted:
hive.executeQuery("SELECT * FROM tm11").show(false)
Run a minor compaction:
hive.executeQuery(
"ALTER TABLE tm11 COMPACT 'minor'"
)
Run a major compaction:
hive.executeQuery(
"ALTER TABLE tm11 COMPACT 'major'"
)
Run:
hive.executeQuery(
"ALTER TABLE tm11 RENAME TO tm12"
)
Verify the renamed table:
hive.showTables().show(100)
Create the source and target tables:
hive.executeQuery(
"CREATE TABLE merge_demo1(" +
"ID INT, FirstName VARCHAR(100), LastName VARCHAR(100))"
)
hive.executeQuery(
"CREATE TABLE merge_demo2(" +
"ID INT, FirstName VARCHAR(100), LastName VARCHAR(100))"
)
Insert data into the target table:
hive.executeQuery(
"INSERT INTO merge_demo1 VALUES (1, 'aaaa', 'bbbb')"
)
hive.executeQuery(
"INSERT INTO merge_demo1 VALUES (2, 'cccc', 'dddd')"
)
hive.executeQuery(
"INSERT INTO merge_demo1 VALUES (3, 'eeee', 'ffff')"
)
hive.executeQuery(
"INSERT INTO merge_demo1 VALUES (4, 'gggg', 'hhhh')"
)
hive.executeQuery(
"INSERT INTO merge_demo1 VALUES (5, 'iiii', 'jjjj')"
)
Insert data into the source table:
hive.executeQuery(
"INSERT INTO merge_demo2 VALUES (2, 'cccc', 'kkkk')"
)
hive.executeQuery(
"INSERT INTO merge_demo2 VALUES (3, 'eeee', 'llll')"
)
Merge the source table into the target table:
hive.executeQuery(
"""MERGE INTO merge_demo1
|USING (SELECT * FROM merge_demo2) sub
|ON sub.id = merge_demo1.id
|WHEN MATCHED THEN
| UPDATE SET
| firstname = sub.firstname,
| lastname = sub.lastname
|WHEN NOT MATCHED THEN
| INSERT VALUES
| (sub.id, sub.firstname, sub.lastname)""".stripMargin
)
Verify the results:
hive.executeQuery(
"SELECT * FROM merge_demo1"
).show()
You can use spark-submit to run a PySpark application with HWC.
The following example creates a Spark session and then creates the required Hive Warehouse session:
from pyspark.sql import SparkSession
from pyspark.conf import SparkConf
from pyspark_llap import HiveWarehouseSession
settings = [
(
"spark.sql.hive.hiveserver2.jdbc.url",
"<HIVE_JDBC_URL>"
)
]
conf = (
SparkConf()
.setAppName("PySpark and Hive")
.setAll(settings)
)
spark = (
SparkSession.builder
.config(conf=conf)
.master("yarn")
.enableHiveSupport()
.getOrCreate()
)
hive = HiveWarehouseSession.session(spark).build()
hive.showDatabases().show()
hive.execute("SELECT 2 GROUP BY 1 ORDER BY 1").show()
spark.stop()
Important: Create a HiveWarehouseSession before you run HWC operations. Creating only a standard Spark session isn't sufficient for HWC operations.
Submit the application by using the HWC JAR and PySpark package:
spark-submit \
--conf spark.sql.hive.hiveserver2.jdbc.url="<HIVE_JDBC_URL>" \
--jars /usr/odp/<ODP_VERSION>/hive_warehouse_connector/hive-warehouse-connector-spark3-assembly-1.0.0.jar \
--conf spark.datasource.hive.warehouse.read.mode=JDBC_CLUSTER \
--conf spark.sql.hive.hiveserver2.jdbc.url.principal="hive/_HOST@EXAMPLE.COM" \
--py-files /usr/odp/<ODP_VERSION>/hive_warehouse_connector/pyspark3_hwc-1.0.0.zip \
/path/to/application.py
Replace <HIVE_JDBC_URL>, <ODP_VERSION>, and /path/to/application.py with values for your environment.
Start PySpark with the required HWC configuration:
pyspark \
--master yarn \
--conf spark.sql.hive.hiveserver2.jdbc.url="<HIVE_JDBC_URL>" \
--jars /usr/odp/<ODP_VERSION>/hive_warehouse_connector/hive-warehouse-connector-spark3-assembly-1.0.0.jar \
--conf spark.datasource.hive.warehouse.read.mode=JDBC_CLUSTER \
--conf spark.sql.hive.hiveserver2.jdbc.url.principal="hive/_HOST@EXAMPLE.COM" \
--py-files /usr/odp/<ODP_VERSION>/hive_warehouse_connector/pyspark3_hwc-1.0.0.zip
Import the required classes:
from pyspark.sql import SparkSession
from pyspark.conf import SparkConf
from pyspark_llap import HiveWarehouseSession
Create the HWC session:
hive = HiveWarehouseSession.session(spark).build()
Run:
hive.showDatabases().show()
For example:
hive.execute(
"SELECT 2 GROUP BY 1 ORDER BY 1"
).show()
You can use the HWC session to run supported Hive operations from your PySpark application.