Blog
hiveorcpysparksparkcompactionsmall-fileshdfsoperations

Hive Non-ACID ORC 테이블 Compaction 가이드

Compactor가 동작하지 않는 Hive Non-ACID ORC 테이블의 Small File을 PySpark로 Staging·Row Count 검증을 거쳐 안전하게 병합하는 방법과 운영 기준을 정리했습니다.

Data Dynamics2026年10月1日21 min read
This post is not yet translated. The original Korean version is shown below.

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 files

ACID(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=trueACID
Partition 경로에 base_*, delta_*, delete_delta_* 디렉터리 존재ACID
데이터 파일 이름이 bucket_00000 형태ACID
위에 해당하지 않고 000000_0 등의 파일이 Partition 경로에 바로 존재Non-ACID (이 글의 대상)
Loading diagram…

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 방식이 더 유연합니다. 이하 내용은 이 방식입니다.

전체 처리 구조

Loading diagram…

예시 테이블은 다음과 같습니다.

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_0

Hive가 만든 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.
Loading diagram…

별도 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 여부를 판단합니다.

Loading diagram…

현재 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-mb512목표 ORC 파일 크기
--target-files자동 계산파일 수 직접 지정(1 이상). 지정 시 크기 계산보다 우선
--min-file-count100이 값 미만이면 Skip
--merge-modecoalescecoalesce 또는 repartition
--staging-basehdfs:///tmp/hive_orc_compactionStaging 루트 경로
--keep-stagingoff성공해도 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 8

PARTITIONED 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=01

coalesce와 repartition

Loading diagram…

일반적인 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 ≈ 20
Loading diagram…

ORC를 다시 쓰면 Stripe 구성과 압축률이 달라지므로 Compaction 후 전체 용량은 원본과 약간 다를 수 있습니다.

운영 시 주의사항

적재 중인 Partition은 건드리지 않는다

오늘(dt=2026-10-01) Partition에 데이터가 계속 들어오고 있다면 Compaction 대상에서 제외합니다. Spark가 읽은 뒤 들어온 데이터는 INSERT OVERWRITE로 사라집니다. 일반적으로 D-1 이전, 적재가 완전히 끝난 Partition만 처리합니다.

Loading diagram…

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을 자동으로 검사·병합하는 형태를 권장합니다.

Loading diagram…

권장 운영 값

설정권장 시작값
Compaction 단위Hive Partition
최소 파일 수100개
목표 ORC 크기512 MB (256 MB ~ 512 MB)
작은 Partition최소 1 File
기본 Merge 방식coalesce
데이터 편중 시repartition
대상적재 완료 Partition (일반적으로 D-1 이전)
검증Row Count (Source·Staging·최종)
StagingHDFS 별도 경로
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 보존)는 이 두 원칙을 운영에서 안전하게 지키기 위한 장치입니다.