Hive 테이블
Hive 테이블 (Hive Tables)
Spark SQL로 Apache Hive에 저장된 데이터를 읽고 쓰는 방법을 정리한 문서예요. Hive 지원을 활성화한 SparkSession 생성, Hive 테이블의 스토리지 포맷 지정, 그리고 다양한 버전의 Hive Metastore와 상호작용하는 방법을 알아볼게요.
출처: 문서
본문
Spark SQL은 Apache Hive에 저장된 데이터의 읽기와 쓰기도 지원해요. 하지만 Hive는 의존성이 많기 때문에 이러한 의존성은 기본 Spark 배포판에 포함되지 않아요. Hive 의존성이 classpath에서 발견되면 Spark가 자동으로 로드해요. 이 Hive 의존성들은 모든 워커 노드에도 있어야 한다는 점을 기억하세요. Hive에 저장된 데이터에 접근하려면 Hive 직렬화·역직렬화 라이브러리(SerDes)에 접근해야 하기 때문이에요.
Hive의 구성은 hive-site.xml, core-site.xml(보안 구성용), hdfs-site.xml(HDFS 구성용) 파일을 conf/에 배치하면 돼요.
Hive로 작업할 때는 Hive 지원을 포함한 SparkSession을 인스턴스화해야 해요. 여기에는 영구 Hive metastore에 대한 연결, Hive serdes 지원, Hive 사용자 정의 함수 지원이 포함돼요. 기존 Hive 배포가 없는 사용자도 Hive 지원을 활성화할 수 있어요. hive-site.xml로 구성되지 않은 경우, 컨텍스트는 현재 디렉터리에 metastore_db를 자동으로 만들고 spark.sql.warehouse.dir로 구성된 디렉터리를 생성해요. 기본값은 Spark 애플리케이션이 시작된 현재 디렉터리의 spark-warehouse 디렉터리예요. 참고로 hive-site.xml의 hive.metastore.warehouse.dir 속성은 Spark 2.0.0부터 폐기됐어요. 대신 spark.sql.warehouse.dir을 사용해 웨어하우스에서 데이터베이스의 기본 위치를 지정하세요. Spark 애플리케이션을 시작하는 사용자에게 쓰기 권한을 부여해야 할 수 있어요.
from os.path import abspath
from pyspark.sql import SparkSession
from pyspark.sql import Row
# warehouse_location points to the default location for managed databases and tables
warehouse_location = abspath('spark-warehouse')
spark = SparkSession \
.builder \
.appName("Python Spark SQL Hive integration example") \
.config("spark.sql.warehouse.dir", warehouse_location) \
.enableHiveSupport() \
.getOrCreate()
# spark is an existing SparkSession
spark.sql("CREATE TABLE IF NOT EXISTS src (key INT, value STRING) USING hive")
spark.sql("LOAD DATA LOCAL INPATH 'examples/src/main/resources/kv1.txt' INTO TABLE src")
# Queries are expressed in HiveQL
spark.sql("SELECT * FROM src").show()
# +---+-------+
# |key| value|
# +---+-------+
# |238|val_238|
# | 86| val_86|
# |311|val_311|
# ...
# Aggregation queries are also supported.
spark.sql("SELECT COUNT(*) FROM src").show()
# +--------+
# |count(1)|
# +--------+
# | 500 |
# +--------+
# The results of SQL queries are themselves DataFrames and support all normal functions.
sqlDF = spark.sql("SELECT key, value FROM src WHERE key < 10 ORDER BY key")
# The items in DataFrames are of type Row, which allows you to access each column by ordinal.
stringsDS = sqlDF.rdd.map(lambda row: "Key: %d, Value: %s" % (row.key, row.value))
for record in stringsDS.collect():
print(record)
# Key: 0, Value: val_0
# Key: 0, Value: val_0
# Key: 0, Value: val_0
# ...
# You can also use DataFrames to create temporary views within a SparkSession.
Record = Row("key", "value")
recordsDF = spark.createDataFrame([Record(i, "val_" + str(i)) for i in range(1, 101)])
recordsDF.createOrReplaceTempView("records")
# Queries can then join DataFrame data with data stored in Hive.
spark.sql("SELECT * FROM records r JOIN src s ON r.key = s.key").show()
# +---+------+---+------+
# |key| value|key| value|
# +---+------+---+------+
# | 2| val_2| 2| val_2|
# | 4| val_4| 4| val_4|
# | 5| val_5| 5| val_5|
# ...
전체 예제 코드는 Spark 저장소의 examples/src/main/python/sql/hive.py 에서 찾을 수 있어요.
import java.io.File
import org.apache.spark.sql.{Row, SaveMode, SparkSession}
case class Record(key: Int, value: String)
// warehouseLocation points to the default location for managed databases and tables
val warehouseLocation = new File("spark-warehouse").getAbsolutePath
val spark = SparkSession
.builder()
.appName("Spark Hive Example")
.config("spark.sql.warehouse.dir", warehouseLocation)
.enableHiveSupport()
.getOrCreate()
import spark.implicits._
import spark.sql
sql("CREATE TABLE IF NOT EXISTS src (key INT, value STRING) USING hive")
sql("LOAD DATA LOCAL INPATH 'examples/src/main/resources/kv1.txt' INTO TABLE src")
// Queries are expressed in HiveQL
sql("SELECT * FROM src").show()
// +---+-------+
// |key| value|
// +---+-------+
// |238|val_238|
// | 86| val_86|
// |311|val_311|
// ...
// Aggregation queries are also supported.
sql("SELECT COUNT(*) FROM src").show()
// +--------+
// |count(1)|
// +--------+
// | 500 |
// +--------+
// The results of SQL queries are themselves DataFrames and support all normal functions.
val sqlDF = sql("SELECT key, value FROM src WHERE key < 10 ORDER BY key")
// The items in DataFrames are of type Row, which allows you to access each column by ordinal.
val stringsDS = sqlDF.map {
case Row(key: Int, value: String) => s"Key: $key, Value: $value"
}
stringsDS.show()
// +--------------------+
// | value|
// +--------------------+
// |Key: 0, Value: val_0|
// |Key: 0, Value: val_0|
// |Key: 0, Value: val_0|
// ...
// You can also use DataFrames to create temporary views within a SparkSession.
val recordsDF = spark.createDataFrame((1 to 100).map(i => Record(i, s"val_$i")))
recordsDF.createOrReplaceTempView("records")
// Queries can then join DataFrame data with data stored in Hive.
sql("SELECT * FROM records r JOIN src s ON r.key = s.key").show()
// +---+------+---+------+
// |key| value|key| value|
// +---+------+---+------+
// | 2| val_2| 2| val_2|
// | 4| val_4| 4| val_4|
// | 5| val_5| 5| val_5|
// ...
// Create a Hive managed Parquet table, with HQL syntax instead of the Spark SQL native syntax
// `USING hive`
sql("CREATE TABLE hive_records(key int, value string) STORED AS PARQUET")
// Save DataFrame to the Hive managed table
val df = spark.table("src")
df.write.mode(SaveMode.Overwrite).saveAsTable("hive_records")
// After insertion, the Hive managed table has data now
sql("SELECT * FROM hive_records").show()
// +---+-------+
// |key| value|
// +---+-------+
// |238|val_238|
// | 86| val_86|
// |311|val_311|
// ...
// Prepare a Parquet data directory
val dataDir = "/tmp/parquet_data"
spark.range(10).write.parquet(dataDir)
// Create a Hive external Parquet table
sql(s"CREATE EXTERNAL TABLE hive_bigints(id bigint) STORED AS PARQUET LOCATION '$dataDir'")
// The Hive external table should already have data
sql("SELECT * FROM hive_bigints").show()
// +---+
// | id|
// +---+
// | 0|
// | 1|
// | 2|
// ... Order may vary, as spark processes the partitions in parallel.
// Turn on flag for Hive Dynamic Partitioning
spark.conf.set("hive.exec.dynamic.partition", "true")
spark.conf.set("hive.exec.dynamic.partition.mode", "nonstrict")
// Create a Hive partitioned table using DataFrame API
df.write.partitionBy("key").format("hive").saveAsTable("hive_part_tbl")
// Partitioned column `key` will be moved to the end of the schema.
sql("SELECT * FROM hive_part_tbl").show()
// +-------+---+
// | value|key|
// +-------+---+
// |val_238|238|
// | val_86| 86|
// |val_311|311|
// ...
spark.stop()
전체 예제 코드는 Spark 저장소의 examples/src/main/scala/org/apache/spark/examples/sql/hive/SparkHiveExample.scala 에서 찾을 수 있어요.
import java.io.File;
import java.io.Serializable;
import java.util.ArrayList;
import java.util.List;
import org.apache.spark.api.java.function.MapFunction;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Encoders;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
public static class Record implements Serializable {
private int key;
private String value;
public int getKey() {
return key;
}
public void setKey(int key) {
this.key = key;
}
public String getValue() {
return value;
}
public void setValue(String value) {
this.value = value;
}
}
// warehouseLocation points to the default location for managed databases and tables
String warehouseLocation = new File("spark-warehouse").getAbsolutePath();
SparkSession spark = SparkSession
.builder()
.appName("Java Spark Hive Example")
.config("spark.sql.warehouse.dir", warehouseLocation)
.enableHiveSupport()
.getOrCreate();
spark.sql("CREATE TABLE IF NOT EXISTS src (key INT, value STRING) USING hive");
spark.sql("LOAD DATA LOCAL INPATH 'examples/src/main/resources/kv1.txt' INTO TABLE src");
// Queries are expressed in HiveQL
spark.sql("SELECT * FROM src").show();
// +---+-------+
// |key| value|
// +---+-------+
// |238|val_238|
// | 86| val_86|
// |311|val_311|
// ...
// Aggregation queries are also supported.
spark.sql("SELECT COUNT(*) FROM src").show();
// +--------+
// |count(1)|
// +--------+
// | 500 |
// +--------+
// The results of SQL queries are themselves DataFrames and support all normal functions.
Dataset<Row> sqlDF = spark.sql("SELECT key, value FROM src WHERE key < 10 ORDER BY key");
// The items in DataFrames are of type Row, which lets you to access each column by ordinal.
Dataset<String> stringsDS = sqlDF.map(
(MapFunction<Row, String>) row -> "Key: " + row.get(0) + ", Value: " + row.get(1),
Encoders.STRING());
stringsDS.show();
// +--------------------+
// | value|
// +--------------------+
// |Key: 0, Value: val_0|
// |Key: 0, Value: val_0|
// |Key: 0, Value: val_0|
// ...
// You can also use DataFrames to create temporary views within a SparkSession.
List<Record> records = new ArrayList<>();
for (int key = 1; key < 100; key++) {
Record record = new Record();
record.setKey(key);
record.setValue("val_" + key);
records.add(record);
}
Dataset<Row> recordsDF = spark.createDataFrame(records, Record.class);
recordsDF.createOrReplaceTempView("records");
// Queries can then join DataFrames data with data stored in Hive.
spark.sql("SELECT * FROM records r JOIN src s ON r.key = s.key").show();
// +---+------+---+------+
// |key| value|key| value|
// +---+------+---+------+
// | 2| val_2| 2| val_2|
// | 2| val_2| 2| val_2|
// | 4| val_4| 4| val_4|
// ...
전체 예제 코드는 Spark 저장소의 examples/src/main/java/org/apache/spark/examples/sql/hive/JavaSparkHiveExample.java 에서 찾을 수 있어요.
Hive로 작업할 때는 Hive 지원이 포함된 SparkSession을 인스턴스화해야 해요. 이는 MetaStore에서 테이블을 찾고 HiveQL로 쿼리를 작성하는 지원을 추가해요.
# enableHiveSupport defaults to TRUE
sparkR.session(enableHiveSupport = TRUE)
sql("CREATE TABLE IF NOT EXISTS src (key INT, value STRING) USING hive")
sql("LOAD DATA LOCAL INPATH 'examples/src/main/resources/kv1.txt' INTO TABLE src")
# Queries can be expressed in HiveQL.
results <- collect(sql("FROM src SELECT key, value"))
전체 예제 코드는 Spark 저장소의 examples/src/main/r/RSparkSQLExample.R 에서 찾을 수 있어요.
Hive 테이블의 스토리지 포맷 지정하기 (Specifying storage format for Hive tables)
Hive 테이블을 만들 때, 이 테이블이 파일 시스템에서 데이터를 어떻게 읽고 쓰는지, 즉 "input format"과 "output format"을 정의해야 해요. 또한 이 테이블이 데이터를 행으로 역직렬화하거나 행을 데이터로 직렬화하는 방법, 즉 "serde"도 정의해야 해요. 다음 옵션을 사용해 스토리지 포맷("serde", "input format", "output format")을 지정할 수 있어요. 예: CREATE TABLE src(id int) USING hive OPTIONS(fileFormat 'parquet'). 기본적으로 테이블 파일을 일반 텍스트로 읽어요. 참고로 테이블 생성 시 Hive storage handler는 아직 지원되지 않아요. Hive 쪽에서 storage handler로 테이블을 만들고 Spark SQL로 읽을 수는 있어요.
| 속성 이름 (Property Name) | 의미 (Meaning) |
|---|---|
fileFormat |
파일포맷(fileFormat)은 "serde", "input format", "output format"을 포함한 스토리지 포맷 명세의 묶음(package)이에요. 현재 6개의 fileFormat을 지원해요: sequencefile, rcfile, orc, parquet, textfile, avro. |
inputFormat, outputFormat |
이 2개 옵션은 해당 InputFormat 및 OutputFormat 클래스의 이름을 문자열 리터럴로 지정해요. 예: org.apache.hadoop.hive.ql.io.orc.OrcInputFormat. 이 2개 옵션은 반드시 쌍으로 나타나야 하며, 이미 fileFormat 옵션을 지정했다면 지정할 수 없어요. |
serde |
이 옵션은 serde 클래스의 이름을 지정해요. fileFormat 옵션을 지정했다면, 주어진 fileFormat이 이미 serde 정보를 포함하는 경우에는 이 옵션을 지정하지 마세요. 현재 "sequencefile", "textfile", "rcfile"은 serde 정보를 포함하지 않으며, 이 3개의 fileFormat과 함께 이 옵션을 사용할 수 있어요. |
fieldDelim, escapeDelim, collectionDelim, mapkeyDelim, lineDelim |
이 옵션들은 "textfile" fileFormat에서만 사용할 수 있어요. 구분자로 구분된 파일을 행으로 읽는 방법을 정의해요. |
OPTIONS로 정의된 다른 모든 속성은 Hive serde 속성으로 간주돼요.
다양한 버전의 Hive Metastore와 상호작용하기 (Interacting with Different Versions of Hive Metastore)
Spark SQL의 Hive 지원에서 가장 중요한 부분 중 하나는 Hive metastore와의 상호작용인데, 이를 통해 Spark SQL이 Hive 테이블의 메타데이터에 접근할 수 있어요. Spark 1.4.0부터, 아래에서 설명하는 구성을 사용해 단일 바이너리 빌드의 Spark SQL로 다양한 버전의 Hive metastore를 쿼리할 수 있어요. 참고로 metastore와 대화하는 데 사용되는 Hive 버전과 무관하게, Spark SQL은 내부적으로 내장 Hive를 대상으로 컴파일하고 내부 실행(serdes, UDFs, UDAFs 등)에 그 클래스들을 사용해요.
메타데이터를 검색하는 데 사용되는 Hive 버전을 구성하는 데 다음 옵션을 사용할 수 있어요:
| 속성 이름 (Property Name) | 기본값 (Default) | 의미 (Meaning) | 버전부터 (Since Version) |
|---|---|---|---|
spark.sql.hive.metastore.version |
2.3.10 |
Hive metastore의 버전이에요. 사용 가능한 옵션은 2.0.0부터 2.3.10, 3.0.0부터 3.1.3, 4.0.0부터 4.1.0이에요. |
1.4.0 |
spark.sql.hive.metastore.jars |
builtin |
HiveMetastoreClient를 인스턴스화하는 데 사용할 jars의 위치예요. 이 속성은 다음 네 가지 옵션 중 하나일 수 있어요: - builtin: -Phive가 활성화되어 있을 때 Spark 어셈블리에 번들된 Hive 2.3.10을 사용해요. 이 옵션을 선택하면 spark.sql.hive.metastore.version은 2.3.10이거나 정의되지 않아야 해요.- maven: Maven 저장소에서 다운로드한 지정 버전의 Hive jars를 사용해요. 이 구성은 일반적으로 프로덕션 배포에 권장되지 않아요.- path: spark.sql.hive.metastore.jars.path로 쉼표 구분 형식으로 구성된 Hive jars를 사용해요. 로컬·원격 경로를 모두 지원해요. 제공된 jars는 spark.sql.hive.metastore.version과 같은 버전이어야 해요.- JVM의 표준 형식 classpath: 이 classpath는 올바른 버전의 Hadoop을 포함한 모든 Hive와 그 의존성을 포함해야 해요. 제공된 jars는 spark.sql.hive.metastore.version과 같은 버전이어야 해요. 이 jars는 드라이버에만 있으면 되지만, yarn cluster 모드로 실행 중이라면 애플리케이션과 함께 패키징해야 해요. |
1.4.0 |
spark.sql.hive.metastore.jars.path |
(empty) |
HiveMetastoreClient를 인스턴스화하는 데 사용되는 jars의 쉼표 구분 경로예요. 이 구성은 spark.sql.hive.metastore.jars가 path로 설정된 경우에만 유용해요. 경로는 다음 형식 중 하나일 수 있어요: - file://path/to/jar/foo.jar- hdfs://nameservice/path/to/jar/foo.jar- /path/to/jar/ (URI 스킴이 없는 경로는 conf fs.defaultFS의 URI 스킴을 따름)- [http/https/ftp]://path/to/jar/foo.jar참고로 1, 2, 3은 와일드카드를 지원해요. 예를 들어: - file://path/to/jar/*,file://path2/to/jar/*/*.jar- hdfs://nameservice/path/to/jar/*,hdfs://nameservice2/path/to/jar/*/*.jar |
3.1.0 |
spark.sql.hive.metastore.sharedPrefixes |
com.mysql.jdbc,org.postgresql,com.microsoft.sqlserver,oracle.jdbc |
Spark SQL과 특정 버전의 Hive 사이에 공유되는 classloader로 로드되어야 하는 클래스 접두사들의 쉼표 구분 목록이에요. 공유되어야 하는 클래스의 예는 metastore와 대화하는 데 필요한 JDBC 드라이버예요. 공유되어야 하는 다른 클래스는 이미 공유된 클래스와 상호작용하는 클래스들이에요. 예를 들어 log4j가 사용하는 커스텀 appender가 있어요. | 1.4.0 |
spark.sql.hive.metastore.barrierPrefixes |
(empty) |
Spark SQL이 통신하는 각 버전의 Hive에 대해 명시적으로 다시 로드되어야 하는 클래스 접두사들의 쉼표 구분 목록이에요. 예를 들어 일반적으로 공유되는 접두사(즉, org.apache.spark.*)에 선언된 Hive UDF가 있어요. |
1.4.0 |
더 알아보기 (Learn more)
- 아파치 스파크 Hive 테이블 (원문)
- SQL 데이터 소스 (원문) — 데이터 소스 개요
- Spark SQL 시작하기 (원문) — Spark SQL 가이드