Hive Non-ACID ORC 테이블 Compaction 가이드
Compactor가 동작하지 않는 Hive Non-ACID ORC 테이블의 Small File을 PySpark로 Staging·Row Count 검증을 거쳐 안전하게 병합하는 방법과 운영 기준을 정리했습니다.
Hive ORC 테이블에 INSERT가 반복되면 하나의 Partition 아래에 수백~수천 개의 작은 ORC 파일이 쌓입니다. 작은 파일은 NameNode 메모리를 잠식하고, 쿼리마다 열어야 할 파일·Split 수를 늘려 Hive·Impala·Spark 모두를 느리게 만듭니다.
dt=2026-10-01
2 MB, 3 MB, 1 MB, 4 MB, 2 MB, ... → 1,000 filesACID(transactional) 테이블이라면 Hive Compactor가 ALTER TABLE ... COMPACT 'MAJOR'로 이를 정리해 줍니다. 하지만 Non-ACID(External 또는 transactional=false) 테이블에는 Compactor가 동작하지 않습니다. 이 글은 Non-ACID ORC 테이블의 Small File을 PySpark로 안전하게 병합하는 방법을 다룹니다. 전체 코드는 DataDynamics/hive-orc-table-compaction 저장소에 있습니다.
목표는 파일을 무조건 하나로 합치는 것이 아니라 적절한 크기(256 MB ~ 512 MB)의 ORC 파일로 재구성하는 것입니다. 예를 들어 약 2.4 GB(평균 2.4 MB × 1,000개)라면 5~10개 파일이 적당합니다.
먼저: ACID인지 Non-ACID인지 확인
SHOW TBLPROPERTIES dw.sales('transactional');hdfs dfs -ls <partition 경로>| 확인 결과 | 테이블 유형 |
|---|---|
transactional=true | ACID |
Partition 경로에 base_*, delta_*, delete_delta_* 디렉터리 존재 | ACID |
데이터 파일 이름이 bucket_00000 형태 | ACID |
위에 해당하지 않고 000000_0 등의 파일이 Partition 경로에 바로 존재 | Non-ACID (이 글의 대상) |
ACID 테이블에는 이 글의 PySpark 프로그램을 사용하면 안 됩니다. Spark는 Hive Warehouse Connector 없이 ACID 테이블을 올바르게 읽거나 쓸 수 없고, HDP 3.x / CDP처럼 Managed Table이 기본적으로 ACID로 생성되는 환경에서는 별도 지정 없이도 ACID일 수 있습니다. ACID 테이블은 Hive 자체 Compaction으로 정리하며, 자세한 방법은 Hive ACID ORC 테이블 Compaction 가이드에 따로 정리했습니다.
Hive만으로 할 수 있는 것: CONCATENATE와 hive.merge
Non-ACID ORC 테이블에는 Hive 자체 기능도 두 가지 있습니다.
1) ALTER TABLE ... CONCATENATE — ORC 파일을 Stripe 단위로 이어 붙입니다. 압축을 풀고 다시 쓰지 않아 빠르지만, 원래의 작은 Stripe가 그대로 남고 결과 파일 크기를 세밀하게 통제하기 어렵습니다.
ALTER TABLE dw.sales PARTITION (dt='2026-10-01') CONCATENATE;2) hive.merge.* 설정 — 애초에 작은 파일이 덜 생기도록 적재 시점에 병합 단계를 추가합니다. 이미 쌓인 파일을 정리하는 수단이 아니라 예방책입니다.
SET hive.merge.tezfiles=true; -- Tez 엔진
SET hive.merge.mapfiles=true;
SET hive.merge.mapredfiles=true;
SET hive.merge.smallfiles.avgsize=134217728; -- 평균 128 MB 미만이면 병합
SET hive.merge.size.per.task=536870912; -- 병합 결과 목표 512 MB이미 쌓인 Small File을 목표 크기로 다시 써서, 검증까지 거쳐 교체하고 싶다면 Spark 방식이 더 유연합니다. 이하 내용은 이 방식입니다.
전체 처리 구조
예시 테이블은 다음과 같습니다.
CREATE TABLE dw.sales (
id BIGINT,
customer_id STRING,
amount DECIMAL(18,2)
)
PARTITIONED BY (dt STRING)
STORED AS ORC;/user/hive/warehouse/dw.db/sales/
├── dt=2026-09-29
├── dt=2026-09-30
└── dt=2026-10-01
├── 000000_0
├── 000001_0
├── ...
└── 000999_0Hive가 만든 ORC 파일은 보통 000000_0처럼 확장자가 없고, Spark로 Compaction한 뒤에는 part-00000-<uuid>-c000.snappy.orc 형태로 바뀝니다.
원본에 바로 Overwrite하면 안 되는 이유
읽고 있는 Partition에 같은 DataFrame을 곧바로 Overwrite하면, Spark가 읽는 도중에 원본 디렉터리를 지우게 됩니다. 대부분 Spark가 다음 에러로 막아 주지만, 막히지 않는 경로라면 데이터 유실로 이어집니다.
Cannot overwrite a path that is also being read from.별도 HDFS 경로의 Staging에 먼저 쓰고, 검증한 뒤, Staging에서 원본 Partition으로 INSERT OVERWRITE 하는 것이 이 프로그램의 핵심 구조입니다.
목표 파일 수 계산과 Compaction 판단
파일 수를 target_files = 4처럼 고정하지 않고 Partition 크기로 계산합니다.
Target Files = ceil(Partition Size / Target File Size)
예) 10 GB / 512 MB = 10 × 1024 / 512 = 20그리고 다음 정책으로 Compaction 여부를 판단합니다.
현재 1,027개 파일, 10 GB, 목표 512 MB라면 1,027 → 약 20개로 병합합니다.
PySpark 프로그램
hive_orc_compaction.py는 다음 기능을 포함합니다.
- 안전장치: ACID 테이블·Bucket 테이블 차단, 모든 Partition Column 지정 여부 검증
- 판단: Hive Metastore에서 Partition Location 조회, HDFS 용량·파일 수 확인, 목표 파일 수 자동 계산, 최소 파일 수 미만이면 Skip
- 병합:
coalesce()또는repartition(), Staging ORC 생성 - 검증: Source/Staging Row Count 비교 →
INSERT OVERWRITE PARTITION→ 최종 Row Count 비교 - 정리: 성공 시에만 Staging 삭제 (실패 시 복구용으로 보존)
#!/usr/bin/env python3
"""Hive Non-ACID ORC Partition Small File Compaction (PySpark)"""
import argparse
import math
import uuid
from datetime import datetime
from pyspark.sql import SparkSession
spark = (
SparkSession.builder
.appName("HiveORCPartitionCompaction")
.enableHiveSupport()
.getOrCreate()
)
# ------------------------------------------------------------
# Utility
# ------------------------------------------------------------
def sql_quote(value):
return str(value).replace("'", "''")
def parse_partition_args(partition_args):
"""--partition dt=2026-10-01 (복수 지정 가능) → dict"""
result = {}
for item in partition_args:
if "=" not in item:
raise ValueError(f"잘못된 partition 형식: {item}")
key, value = item.split("=", 1)
key, value = key.strip(), value.strip()
if not key:
raise ValueError(f"Partition column이 비어 있습니다: {item}")
result[key] = value
return result
def make_partition_sql(partition_spec):
"""{"dt": "2026-10-01"} → `dt`='2026-10-01'"""
return ", ".join(f"`{k}`='{sql_quote(v)}'" for k, v in partition_spec.items())
def make_where_sql(partition_spec):
return " AND ".join(f"`{k}`='{sql_quote(v)}'" for k, v in partition_spec.items())
def human_size(size):
value = float(size)
for unit in ["B", "KB", "MB", "GB", "TB"]:
if value < 1024:
return f"{value:.2f} {unit}"
value /= 1024
return f"{value:.2f} PB"
# ------------------------------------------------------------
# Table Safety Check
# ------------------------------------------------------------
def validate_table(table_name):
"""
Spark 방식으로 처리하면 안 되는 테이블을 사전에 차단한다.
1. Hive ACID (transactional=true)
2. Hive Bucketed Table (CLUSTERED BY)
Spark 3.x의 SHOW CREATE TABLE은 Hive DDL이 아닌 Spark DDL을
출력하고 Transactional Table에서는 예외를 내므로,
Metastore 정보를 직접 조회한다.
"""
properties = {
str(row[0]).strip(): str(row[1]).strip()
for row in spark.sql(f"SHOW TBLPROPERTIES {table_name}").collect()
}
if properties.get("transactional", "").lower() == "true":
raise RuntimeError(
"Hive ACID Transactional Table입니다. "
"Spark Merge 대신 Hive MAJOR COMPACTION을 사용하십시오."
)
for row in spark.sql(f"DESCRIBE FORMATTED {table_name}").collect():
key = str(row[0]).strip() if row[0] is not None else ""
if key == "Num Buckets":
value = str(row[1]).strip()
if value.lstrip("-").isdigit() and int(value) > 0:
raise RuntimeError(
"Hive Bucketed Table입니다. "
"CLUSTERED BY 테이블에는 이 프로그램을 사용하지 마십시오."
)
print("Table validation OK")
def validate_partition_spec(table_name, partition_spec):
"""
모든 Partition Column이 지정되었는지 확인한다.
year/month/day 테이블에 year=2026만 지정하면 나머지 Partition
Column이 data column처럼 처리되므로 실행을 중단한다.
"""
partition_columns = [
c.name for c in spark.catalog.listColumns(table_name) if c.isPartition
]
if not partition_columns:
raise RuntimeError(f"Partition Table이 아닙니다: {table_name}")
if {c.lower() for c in partition_columns} != {c.lower() for c in partition_spec}:
raise ValueError(
"모든 Partition Column을 지정해야 합니다. "
f"테이블 Partition Column: {partition_columns}, "
f"입력값: {list(partition_spec.keys())}"
)
return partition_columns
# ------------------------------------------------------------
# Partition Location / HDFS
# ------------------------------------------------------------
def get_partition_location(table_name, partition_spec):
rows = spark.sql(
f"DESCRIBE FORMATTED {table_name} "
f"PARTITION ({make_partition_sql(partition_spec)})"
).collect()
for row in rows:
key = str(row[0]).strip() if row[0] is not None else ""
if key == "Location":
return str(row[1]).strip()
raise RuntimeError("Partition Location을 찾을 수 없습니다.")
def _hdfs(path_string):
jvm = spark.sparkContext._jvm
conf = spark.sparkContext._jsc.hadoopConfiguration()
path = jvm.org.apache.hadoop.fs.Path(path_string)
return path.getFileSystem(conf), path
def get_hdfs_statistics(path_string):
fs, path = _hdfs(path_string)
if not fs.exists(path):
raise RuntimeError(f"HDFS Path가 존재하지 않습니다: {path_string}")
summary = fs.getContentSummary(path)
return {"bytes": summary.getLength(), "files": summary.getFileCount()}
def delete_hdfs_path(path_string):
fs, path = _hdfs(path_string)
if fs.exists(path):
fs.delete(path, True)
def calculate_target_files(partition_size_bytes, target_file_size_mb):
target_bytes = target_file_size_mb * 1024 * 1024
return max(1, math.ceil(partition_size_bytes / target_bytes))
# ------------------------------------------------------------
# Main Compaction
# ------------------------------------------------------------
def compact_partition(
table_name,
partition_spec,
staging_base,
target_file_size_mb=512,
min_file_count=100,
target_files_override=None,
merge_mode="coalesce",
keep_staging=False,
):
print(f"Table : {table_name}")
print(f"Partition : {partition_spec}")
# 1. Table / Partition 검증
validate_table(table_name)
partition_columns = validate_partition_spec(table_name, partition_spec)
# 2. Partition Location 및 HDFS 통계
partition_location = get_partition_location(table_name, partition_spec)
stats = get_hdfs_statistics(partition_location)
partition_size, current_files = stats["bytes"], stats["files"]
print(f"Location : {partition_location}")
print(f"Partition Size : {human_size(partition_size)}")
print(f"Current Files : {current_files}")
# 3. Compaction 필요 여부
if current_files < min_file_count:
print(f"파일 수가 {min_file_count}개 미만이므로 Compaction을 수행하지 않습니다.")
return
# 4. Target File 개수 계산
if target_files_override is not None:
if target_files_override < 1:
raise ValueError("--target-files는 1 이상이어야 합니다.")
target_files = target_files_override
else:
target_files = calculate_target_files(partition_size, target_file_size_mb)
print(f"Target Size : {target_file_size_mb} MB")
print(f"Target Files : {target_files}")
if current_files <= target_files:
print("현재 파일 수가 Target Files 이하이므로 Compaction을 수행하지 않습니다.")
return
# 5. Data Column 목록 (Partition Column 제외)
partition_names = {c.lower() for c in partition_columns}
data_columns = [
c for c in spark.table(table_name).columns if c.lower() not in partition_names
]
# 6. 대상 Partition 읽기
where_sql = make_where_sql(partition_spec)
source_df = spark.sql(f"SELECT * FROM {table_name} WHERE {where_sql}")
source_count = source_df.count()
print(f"Source Rows : {source_count}")
if source_count == 0:
print("Partition 데이터가 없습니다.")
return
# 7. Spark Partition 조절
if merge_mode == "repartition":
merged_df = source_df.repartition(target_files)
else:
merged_df = source_df.coalesce(target_files)
# 8. Staging ORC Write
run_id = datetime.now().strftime("%Y%m%d_%H%M%S") + "_" + uuid.uuid4().hex[:8]
staging_path = (
f"{staging_base.rstrip('/')}/{table_name.replace('.', '_')}/{run_id}"
)
print(f"Staging Path : {staging_path}")
merged_df.write.mode("overwrite").format("orc").save(staging_path)
# 9. Staging 검증 — 실패 시 원본 Partition은 손대지 않는다
staged_df = spark.read.format("orc").load(staging_path)
staged_count = staged_df.count()
print(f"Staging Rows : {staged_count}")
if source_count != staged_count:
raise RuntimeError(
"Source와 Staging Row Count가 일치하지 않습니다. "
"원본 Partition은 변경하지 않습니다."
)
# 10. Staging → 원본 Partition INSERT OVERWRITE
# Staging을 다시 읽으면 Spark가 큰 ORC 파일을
# spark.sql.files.maxPartitionBytes(기본 128MB) 단위로 쪼개 읽으므로
# 그대로 INSERT하면 파일 수가 Target Files의 몇 배가 된다.
# 다시 coalesce하여 최종 파일 수를 Target Files로 맞춘다.
temp_view = "compact_" + uuid.uuid4().hex
staged_df.coalesce(target_files).createOrReplaceTempView(temp_view)
select_columns = ", ".join(f"`{c}`" for c in data_columns)
overwrite_sql = f"""
INSERT OVERWRITE TABLE {table_name}
PARTITION ({make_partition_sql(partition_spec)})
SELECT {select_columns}
FROM {temp_view}
"""
print(overwrite_sql)
spark.sql(overwrite_sql)
# 11. 최종 Row Count 검증 — 실패 시 Staging은 복구용으로 남는다
target_count = spark.sql(
f"SELECT COUNT(*) AS cnt FROM {table_name} WHERE {where_sql}"
).first()["cnt"]
print(f"Before Rows : {source_count}")
print(f"After Rows : {target_count}")
if source_count != target_count:
raise RuntimeError("Compaction 이후 Row Count가 일치하지 않습니다.")
# 12. 결과 확인 및 정리
after = get_hdfs_statistics(partition_location)
print(f"Files : {current_files} → {after['files']}")
print(f"Size : {human_size(partition_size)} → {human_size(after['bytes'])}")
spark.catalog.dropTempView(temp_view)
if not keep_staging:
delete_hdfs_path(staging_path)
else:
print(f"Staging Directory 유지: {staging_path}")
print("Compaction SUCCESS")
# ------------------------------------------------------------
# Command Line
# ------------------------------------------------------------
def main():
parser = argparse.ArgumentParser(description="Hive ORC Partition Small File Compaction")
parser.add_argument("--table", required=True, help="Hive table: db.table")
parser.add_argument("--partition", action="append", required=True,
help="예: --partition dt=2026-10-01 (복수 지정 가능)")
parser.add_argument("--staging-base", default="hdfs:///tmp/hive_orc_compaction")
parser.add_argument("--target-file-size-mb", type=int, default=512)
parser.add_argument("--min-file-count", type=int, default=100)
parser.add_argument("--target-files", type=int, default=None,
help="자동 계산 대신 파일 개수 직접 지정")
parser.add_argument("--merge-mode", choices=["coalesce", "repartition"], default="coalesce")
parser.add_argument("--keep-staging", action="store_true")
args = parser.parse_args()
compact_partition(
table_name=args.table,
partition_spec=parse_partition_args(args.partition),
staging_base=args.staging_base,
target_file_size_mb=args.target_file_size_mb,
min_file_count=args.min_file_count,
target_files_override=args.target_files,
merge_mode=args.merge_mode,
keep_staging=args.keep_staging,
)
if __name__ == "__main__":
main()코드에서 짚어 둘 부분이 두 가지 있습니다.
- ACID/Bucket 검사에
SHOW CREATE TABLE을 쓰지 않습니다. Spark 3.x의SHOW CREATE TABLE은 Hive DDL이 아니라 Spark DDL을 출력하고, Transactional Table에서는 예외를 냅니다. 그래서SHOW TBLPROPERTIES의transactional과DESCRIBE FORMATTED의Num Buckets를 직접 확인합니다. - Staging을 다시 읽은 뒤 한 번 더
coalesce합니다. 512 MB짜리 파일을 다시 읽으면 Spark는spark.sql.files.maxPartitionBytes(기본 128 MB) 단위로 쪼개 읽기 때문에, 그대로INSERT하면 최종 파일 수가 Target Files의 몇 배가 됩니다.
실행 방법
# 기본: 목표 512 MB, 최소 파일 수 100, coalesce
spark-submit \
--master yarn \
--deploy-mode cluster \
hive_orc_compaction.py \
--table dw.sales \
--partition dt=2026-10-01| 옵션 | 기본값 | 설명 |
|---|---|---|
--table | (필수) | db.table |
--partition | (필수) | col=value. 복수 Partition Column이면 반복 지정 |
--target-file-size-mb | 512 | 목표 ORC 파일 크기 |
--target-files | 자동 계산 | 파일 수 직접 지정(1 이상). 지정 시 크기 계산보다 우선 |
--min-file-count | 100 | 이 값 미만이면 Skip |
--merge-mode | coalesce | coalesce 또는 repartition |
--staging-base | hdfs:///tmp/hive_orc_compaction | Staging 루트 경로 |
--keep-staging | off | 성공해도 Staging 삭제하지 않음 |
목표 크기를 256 MB로 바꾸거나, 파일 수를 8개로 고정하려면:
spark-submit hive_orc_compaction.py --table dw.sales --partition dt=2026-10-01 --target-file-size-mb 256
spark-submit hive_orc_compaction.py --table dw.sales --partition dt=2026-10-01 --target-files 8PARTITIONED BY (year STRING, month STRING, day STRING)처럼 Partition Column이 여러 개라면 모두 지정해야 합니다. 일부만 지정하면 프로그램이 실행을 중단합니다.
spark-submit hive_orc_compaction.py \
--table dw.sales \
--partition year=2026 \
--partition month=10 \
--partition day=01coalesce와 repartition
일반적인 Small File 정리는 coalesce로 시작합니다. 결과 파일 크기가 심하게 불균등하면 --merge-mode repartition을 사용합니다. Shuffle 비용이 들지만 파일 크기가 고르게 나옵니다.
실제 처리 예
Table : dw.sales
Partition : dt=2026-10-01
Partition Size : 9.7 GB
ORC Files : 1,243
Target Size : 512 MB → Target Files ≈ 20ORC를 다시 쓰면 Stripe 구성과 압축률이 달라지므로 Compaction 후 전체 용량은 원본과 약간 다를 수 있습니다.
운영 시 주의사항
적재 중인 Partition은 건드리지 않는다
오늘(dt=2026-10-01) Partition에 데이터가 계속 들어오고 있다면 Compaction 대상에서 제외합니다. Spark가 읽은 뒤 들어온 데이터는 INSERT OVERWRITE로 사라집니다. 일반적으로 D-1 이전, 적재가 완전히 끝난 Partition만 처리합니다.
INSERT OVERWRITE는 원자적이지 않다
Non-ACID 테이블의 INSERT OVERWRITE는 기존 파일 삭제와 새 파일 이동으로 이뤄집니다. 실행 중에 다른 사용자가 해당 Partition을 조회하면 비어 있거나 일부만 있는 데이터를 볼 수 있으므로 조회가 적은 시간대에 실행합니다.
Bucket 테이블은 대상이 아니다
CLUSTERED BY (customer_id) INTO 32 BUCKETS 같은 테이블은 파일 수와 배치 자체가 Bucket 규칙의 일부입니다. 임의로 coalesce하면 Bucket 구조가 깨지므로, 프로그램은 Num Buckets > 0이면 실행을 중단합니다.
Impala를 함께 쓴다면
Spark/Hive로 파일을 교체한 뒤 Impala에서 조회한다면 메타데이터를 갱신해야 합니다.
REFRESH dw.sales PARTITION (dt='2026-10-01');실패 시 복구
- Staging 검증 단계에서 실패하면 원본 Partition은 변경되지 않습니다.
INSERT OVERWRITE후 최종 검증에서 실패하면 원본은 이미 교체된 상태지만, 예외로 종료되므로 Staging Directory가 삭제되지 않고 남아 있습니다. 실행 로그의Staging Path로 다시 복구할 수 있습니다.
staged_df = spark.read.format("orc").load(
"hdfs:///tmp/hive_orc_compaction/dw_sales/<run_id>"
)
staged_df.createOrReplaceTempView("compact_recovery")
spark.sql("""
INSERT OVERWRITE TABLE dw.sales
PARTITION (dt='2026-10-01')
SELECT id, customer_id, amount
FROM compact_recovery
""")Compaction이 끝나면 Metastore 통계도 갱신해 두는 것이 좋습니다.
ANALYZE TABLE dw.sales PARTITION (dt='2026-10-01') COMPUTE STATISTICS;운영 자동화 구조
일 배치 ETL 뒤에 D-1 Partition을 자동으로 검사·병합하는 형태를 권장합니다.
권장 운영 값
| 설정 | 권장 시작값 |
|---|---|
| Compaction 단위 | Hive Partition |
| 최소 파일 수 | 100개 |
| 목표 ORC 크기 | 512 MB (256 MB ~ 512 MB) |
| 작은 Partition | 최소 1 File |
| 기본 Merge 방식 | coalesce |
| 데이터 편중 시 | repartition |
| 대상 | 적재 완료 Partition (일반적으로 D-1 이전) |
| 검증 | Row Count (Source·Staging·최종) |
| Staging | HDFS 별도 경로 |
| ACID 테이블 | Hive COMPACT 'MAJOR' |
| Bucket 테이블 | 별도 처리 |
정리
원본 Partition → Spark Read → 파일 수 축소 → Staging ORC
→ 건수 검증 → INSERT OVERWRITE PARTITION → 건수 재검증 → Staging 삭제Non-ACID Hive ORC 테이블은 Compactor가 정리해 주지 않으므로 직접 병합해야 합니다. 이때 가장 중요한 두 가지는 원본 Partition을 읽으면서 같은 경로에 바로 Overwrite하지 않는 것과 데이터가 계속 적재 중인 Partition을 Compaction하지 않는 것입니다. 나머지(목표 파일 수 계산, 이중 Row Count 검증, Staging 보존)는 이 두 원칙을 운영에서 안전하게 지키기 위한 장치입니다.