Blog
hiveorcpysparksparkcompactionsmall-fileshdfsoperations

Compacting Hive Non-ACID ORC Tables

The Hive compactor does not run on non-ACID ORC tables. Here is how to safely merge their small files with PySpark using a staging area and row-count validation, plus operational guidelines.

Data DynamicsOctober 1, 202615 min read

When a Hive ORC table receives repeated INSERTs, hundreds or thousands of small ORC files pile up under a single partition. Small files eat NameNode memory and increase the number of files and splits every query has to open, slowing down Hive, Impala and Spark alike.

dt=2026-10-01
  2 MB, 3 MB, 1 MB, 4 MB, 2 MB, ...   → 1,000 files

For ACID (transactional) tables, the Hive compactor cleans this up via ALTER TABLE ... COMPACT 'MAJOR'. But the compactor does not run on non-ACID tables (external or transactional=false). This post shows how to safely merge small files in non-ACID ORC tables with PySpark. The full code lives in the DataDynamics/hive-orc-table-compaction repository.

The goal is not to squash everything into one file but to rebuild the partition into ORC files of a sensible size (256 MB – 512 MB). For roughly 2.4 GB (1,000 files averaging 2.4 MB), 5–10 files is about right.

First: ACID or non-ACID?

SHOW TBLPROPERTIES dw.sales('transactional');
hdfs dfs -ls <partition path>
What you seeTable type
transactional=trueACID
base_*, delta_*, delete_delta_* directories under the partitionACID
Data files named like bucket_00000ACID
None of the above; files such as 000000_0 sit directly in the partitionNon-ACID (this post)
Loading diagram…

Do not use this post's PySpark program on ACID tables. Spark cannot correctly read or write ACID tables without the Hive Warehouse Connector, and on HDP 3.x / CDP, managed tables are ACID by default even if you never asked for it. ACID tables are compacted with Hive's own compaction — see the separate Hive ACID ORC table compaction guide.

What Hive alone can do: CONCATENATE and hive.merge

Hive itself offers two options for non-ACID ORC tables.

1) ALTER TABLE ... CONCATENATE — stitches ORC files together at the stripe level. It is fast because nothing is decompressed and rewritten, but the original small stripes remain and you have little control over the resulting file sizes.

ALTER TABLE dw.sales PARTITION (dt='2026-10-01') CONCATENATE;

2) hive.merge.* settings — add a merge stage at write time so fewer small files are produced in the first place. This is prevention, not a cure for files that already exist.

SET hive.merge.tezfiles=true;              -- Tez engine
SET hive.merge.mapfiles=true;
SET hive.merge.mapredfiles=true;
SET hive.merge.smallfiles.avgsize=134217728;   -- merge if average < 128 MB
SET hive.merge.size.per.task=536870912;        -- target 512 MB per merged file

If you want to rewrite existing small files to a target size and validate before swapping them in, the Spark approach is more flexible. The rest of this post covers it.

Overall flow

Loading diagram…

Example table:

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

ORC files written by Hive usually have no extension (000000_0); after Spark compaction they become part-00000-<uuid>-c000.snappy.orc.

Why not overwrite the source directly

Overwriting the partition you are reading from means Spark deletes the source directory mid-read. Spark usually blocks this with the error below, but on any path where it doesn't, you lose data.

Cannot overwrite a path that is also being read from.
Loading diagram…

Write to a staging area on a separate HDFS path, validate it, then INSERT OVERWRITE the original partition from staging — that is the core of this program.

Target file count and the compaction decision

Instead of hard-coding target_files = 4, derive it from the partition size:

Target Files = ceil(Partition Size / Target File Size)
 
e.g. 10 GB / 512 MB = 10 × 1024 / 512 = 20

Then decide whether to compact:

Loading diagram…

With 1,027 files, 10 GB and a 512 MB target, the partition goes from 1,027 to about 20 files.

The PySpark program

hive_orc_compaction.py provides:

  • Safety: blocks ACID and bucketed tables; verifies that every partition column is specified
  • Decision: looks up the partition location in the Hive Metastore, checks HDFS size and file count, computes the target file count, skips below the minimum file count
  • Merge: coalesce() or repartition(), writes staging ORC
  • Validation: source vs. staging row count → INSERT OVERWRITE PARTITION → final row count
  • Cleanup: deletes staging only on success (kept for recovery on failure)
#!/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 (repeatable) → dict"""
    result = {}
    for item in partition_args:
        if "=" not in item:
            raise ValueError(f"Invalid partition format: {item}")
        key, value = item.split("=", 1)
        key, value = key.strip(), value.strip()
        if not key:
            raise ValueError(f"Empty 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):
    """
    Block tables that must not be rewritten by Spark:
      1. Hive ACID (transactional=true)
      2. Hive Bucketed Table (CLUSTERED BY)
 
    In Spark 3.x, SHOW CREATE TABLE prints Spark DDL (not Hive DDL)
    and throws on transactional tables, so we query Metastore
    information directly.
    """
    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(
            "This is a Hive ACID transactional table. "
            "Use Hive MAJOR COMPACTION instead of a Spark merge."
        )
 
    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(
                    "This is a Hive bucketed table. "
                    "Do not use this program on CLUSTERED BY tables."
                )
    print("Table validation OK")
 
 
def validate_partition_spec(table_name, partition_spec):
    """
    Ensure every partition column is specified.
    On a year/month/day table, passing only year=2026 would make the
    remaining partition columns behave like data columns, so abort.
    """
    partition_columns = [
        c.name for c in spark.catalog.listColumns(table_name) if c.isPartition
    ]
    if not partition_columns:
        raise RuntimeError(f"Not a partitioned table: {table_name}")
    if {c.lower() for c in partition_columns} != {c.lower() for c in partition_spec}:
        raise ValueError(
            "All partition columns must be specified. "
            f"Table partition columns: {partition_columns}, "
            f"given: {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 not found.")
 
 
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 does not exist: {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. Validate table / partition
    validate_table(table_name)
    partition_columns = validate_partition_spec(table_name, partition_spec)
 
    # 2. Partition location and HDFS statistics
    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. Is compaction needed?
    if current_files < min_file_count:
        print(f"Fewer than {min_file_count} files; skipping compaction.")
        return
 
    # 4. Compute target file count
    if target_files_override is not None:
        if target_files_override < 1:
            raise ValueError("--target-files must be >= 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("Current file count <= target files; skipping compaction.")
        return
 
    # 5. Data columns (excluding partition columns)
    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. Read the target 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 is empty.")
        return
 
    # 7. Adjust Spark partitions
    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. Validate staging — on failure the original partition is untouched
    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 and staging row counts differ. "
            "The original partition is left unchanged."
        )
 
    # 10. Staging → original partition via INSERT OVERWRITE
    # Re-reading staging splits large ORC files by
    # spark.sql.files.maxPartitionBytes (default 128MB), so inserting
    # as-is would yield several times more files than target_files.
    # coalesce again to pin the final file count to 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. Final row count check — on failure staging is kept for recovery
    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("Row count mismatch after compaction.")
 
    # 12. Report and clean up
    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"Keeping 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="e.g. --partition dt=2026-10-01 (repeatable)")
    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="Set file count explicitly instead of auto-calculating")
    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()

Two details in the code are worth calling out:

  • It does not use SHOW CREATE TABLE for the ACID/bucket checks. In Spark 3.x, SHOW CREATE TABLE prints Spark DDL rather than Hive DDL and throws on transactional tables, so the program reads transactional from SHOW TBLPROPERTIES and Num Buckets from DESCRIBE FORMATTED directly.
  • It calls coalesce again after re-reading staging. When Spark re-reads 512 MB files it splits them by spark.sql.files.maxPartitionBytes (128 MB by default), so inserting as-is would produce several times more files than the target.

Running it

# Defaults: 512 MB target, min 100 files, coalesce
spark-submit \
  --master yarn \
  --deploy-mode cluster \
  hive_orc_compaction.py \
  --table dw.sales \
  --partition dt=2026-10-01
OptionDefaultDescription
--table(required)db.table
--partition(required)col=value; repeat for multiple partition columns
--target-file-size-mb512Target ORC file size
--target-filesautoExplicit file count (>= 1); overrides the size calculation
--min-file-count100Skip below this file count
--merge-modecoalescecoalesce or repartition
--staging-basehdfs:///tmp/hive_orc_compactionStaging root path
--keep-stagingoffKeep staging even on success

To target 256 MB, or to force exactly 8 files:

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

With multiple partition columns such as PARTITIONED BY (year STRING, month STRING, day STRING), you must specify all of them; the program aborts if any are missing.

spark-submit hive_orc_compaction.py \
  --table dw.sales \
  --partition year=2026 \
  --partition month=10 \
  --partition day=01

coalesce vs. repartition

Loading diagram…

Start with coalesce for ordinary small-file cleanup. If the resulting files are badly uneven, use --merge-mode repartition: it costs a shuffle but produces evenly sized files.

A real example

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…

Rewriting ORC changes stripe layout and compression ratio, so the total size after compaction may differ slightly from the original.

Operational caveats

Never touch a partition that is still loading

If today's partition (dt=2026-10-01) is still receiving data, leave it out. Anything written after Spark reads the partition is wiped by INSERT OVERWRITE. As a rule, only compact partitions whose loading is fully finished, typically D-1 or older.

Loading diagram…

INSERT OVERWRITE is not atomic

On a non-ACID table, INSERT OVERWRITE deletes the old files and moves the new ones in. A user querying the partition in the meantime may see it empty or partially populated, so run compaction during low-query hours.

Bucketed tables are out of scope

For tables like CLUSTERED BY (customer_id) INTO 32 BUCKETS, the file count and layout are part of the bucketing contract. An ad-hoc coalesce breaks it, so the program aborts when Num Buckets > 0.

If you also use Impala

After Spark/Hive swaps the files, refresh Impala's metadata before querying:

REFRESH dw.sales PARTITION (dt='2026-10-01');

Recovering from failure

  • If staging validation fails, the original partition is untouched.
  • If the final check fails after INSERT OVERWRITE, the original has already been replaced — but the program exits with an exception, so the staging directory is not deleted. Use the Staging Path from the log to restore:
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
""")

After compaction it is also a good idea to refresh Metastore statistics:

ANALYZE TABLE dw.sales PARTITION (dt='2026-10-01') COMPUTE STATISTICS;

Automating it

We recommend checking and compacting the D-1 partition automatically after the daily ETL:

Loading diagram…
SettingStarting value
Compaction unitHive partition
Minimum file count100
Target ORC size512 MB (256 MB – 512 MB)
Small partitionsAt least 1 file
Default merge modecoalesce
Heavy skewrepartition
TargetsFully loaded partitions (typically D-1 or older)
ValidationRow count (source, staging, final)
StagingSeparate HDFS path
ACID tablesHive COMPACT 'MAJOR'
Bucketed tablesHandle separately

Summary

Source partition → Spark read → reduce file count → staging ORC
  → validate count → INSERT OVERWRITE PARTITION → re-validate count → delete staging

Non-ACID Hive ORC tables are not cleaned up by the compactor, so you have to merge them yourself. The two rules that matter most are never overwrite the partition you are reading from in place and never compact a partition that is still being loaded. Everything else — target file calculation, double row-count validation, keeping staging on failure — exists to uphold those two rules safely in production.