A Hive dynamic-partition INSERT produced thousands of tiny files — what caused it and what do you tune?
answer
- Count the writers, not the partitions
- Every task opens its own output file
- Route a partition to one reducer
- There is a post-job merge step
- Maybe the partition column is too fine
basics
~20 sEach writing task creates its own file in every partition it touches, so files multiply as tasks times partitions. Funnel each partition to one writer with DISTRIBUTE BY the partition columns, enable Hive's merge settings, and use coarser partition granularity.
solid answer
~50 sWith dynamic partitioning, every task that happens to hold rows for a partition opens its own output file there. If 200 reducers each see rows for 300 partitions, you can get up to 60,000 files instead of 300 — the file count is *tasks × partitions*, not partitions. The primary fix is to control which task writes which partition: add `DISTRIBUTE BY` on the partition columns (or `CLUSTER BY` if you also want ordering) so all rows of a partition are routed to a single reducer, which then writes one file per partition. Second, let Hive merge afterwards with `hive.merge.mapfiles`, `hive.merge.mapredfiles` or `hive.merge.tezfiles`, plus `hive.merge.smallfiles.avgsize` and `hive.merge.size.per.task` to set the trigger and target size. Third, revisit the partition scheme itself: partitioning by hour or by a high-cardinality column guarantees small files no matter how you tune the writers.
code
sql · 8 lines-- before: every reducer may write into every partition
INSERT OVERWRITE TABLE sales PARTITION (dt, country)
SELECT id, amount, dt, country FROM sales_raw;
-- after: one writer per partition
INSERT OVERWRITE TABLE sales PARTITION (dt, country)
SELECT id, amount, dt, country FROM sales_raw
DISTRIBUTE BY dt, country;go deeper
Know what dynamic partitioning is and that lots of small files are bad. Be able to say that each writing task can create a file in every partition it touches.
Explain the tasks-times-partitions arithmetic, what DISTRIBUTE BY changes about the shuffle, and how Hive's merge settings act as a post-job cleanup pass.
Diagnose it end to end: measure average file size, fix routing, watch for the skew you just created, and push back on a partition scheme that guarantees small files whatever you tune.
Own the layout standard — partition granularity, target file size and compaction policy across the platform — so this diagnosis stops recurring in every new pipeline.
## Why the files multiply Dynamic partitioning lets the partition value come from the data rather than from the statement: ```sql SET hive.exec.dynamic.partition.mode=nonstrict; INSERT OVERWRITE TABLE sales PARTITION (dt, country) SELECT id, amount, dt, country FROM sales_raw; ``` Each writing task processes whatever rows it was given. When a task encounters a row for `(dt='2026-01-01', country='DE')` it opens a writer for that partition directory and keeps it open. If rows for a partition are scattered across all tasks, then every task opens a writer in that partition. The output count is therefore bounded by **tasks × distinct partitions**, and with a shuffle-free map-only insert it is bounded by input splits × partitions, which can be worse. That is the whole mechanism, and stating it plainly is most of the answer. ## Fix 1 — route each partition to one writer Add a `DISTRIBUTE BY` on the partition columns so the shuffle sends all rows of a partition to a single reducer: ```sql INSERT OVERWRITE TABLE sales PARTITION (dt, country) SELECT id, amount, dt, country FROM sales_raw DISTRIBUTE BY dt, country; ``` Now each partition has one writer and one output file. Two caveats worth volunteering: - **This creates skew by construction.** A partition holding most of the data becomes one reducer's job and one long-running task. For a badly-skewed partition column, distributing by `dt, country, something_random_bucketed` gives a controlled number of files per partition instead of one, trading a few files for parallelism. - **Do not reach for `ORDER BY` instead.** In Hive, `ORDER BY` forces a *single* reducer for the entire result to produce a global ordering. It certainly reduces the file count, by destroying all parallelism. `SORT BY` orders within a reducer; `CLUSTER BY x` is shorthand for `DISTRIBUTE BY x SORT BY x`. ## Fix 2 — let Hive merge the output Hive can append a merge job that rewrites small files into larger ones: - `hive.merge.mapfiles` — merge the output of map-only jobs (enabled by default). - `hive.merge.mapredfiles` — merge the output of map-reduce jobs (off by default, and this is the one people forget). - `hive.merge.tezfiles` — the Tez equivalent, off by default. - `hive.merge.smallfiles.avgsize` — if the job's average output file is smaller than this, trigger the merge step. - `hive.merge.size.per.task` — the target size each merge task aims to produce. Merging costs an extra pass over the freshly-written data, which is usually a fine trade against thousands of downstream task launches. It is a cleanup, not a substitute for fix 1: it still has to read every tiny file once. ## Fix 3 — question the partition scheme A partition should be large enough to be worth a directory. If partitioning by `dt, country, device_type` yields 30 MB per partition, the layout is the problem. Coarsen it — partition by day only, and let columnar min/max indexes and predicate pushdown inside ORC or Parquet handle the finer filtering. Bucketing is the tool for high-cardinality columns you join on; partitioning is for a handful of coarse, frequently-filtered dimensions. Hive also enforces guardrails that surface this: `hive.exec.max.dynamic.partitions` caps the total number of dynamic partitions a statement may create, and `hive.exec.max.dynamic.partitions.pernode` caps how many a single node may create (100 by default). Hitting the per-node limit produces a fatal error about creating too many dynamic partitions. Raising those limits is occasionally right and much more often a signal that the partition column is too fine-grained — the interviewer is usually checking whether you raise the limit reflexively or ask why it was hit. ## Why anyone cares Three costs compound. Each file is an object in the NameNode's in-memory namespace, so millions of small files pressure the master's heap. Each file (or split) becomes a task, so a query over 60,000 files spends its life launching containers rather than reading data. And columnar formats lose their advantage: an ORC file smaller than a stripe cannot amortize its footer, indexes and dictionary, so compression and pushdown both degrade. ## The complete answer Name the mechanism (tasks × partitions), fix the routing with `DISTRIBUTE BY`, enable merging as cleanup, then challenge the partition granularity — and mention that you would check the resulting average file size afterwards rather than assume the tuning worked.
- Why not just add ORDER BY to reduce the number of output files?In Hive, `ORDER BY` demands a global ordering and so forces a single reducer for the whole result. File count drops to one, and the job loses all parallelism — a scan-sized dataset now streams through one task. `DISTRIBUTE BY` on the partition columns achieves one writer *per partition* while keeping partitions parallel, which is what you actually want.
- DISTRIBUTE BY fixed the file count but now one reducer runs for an hour. What next?You converted a file-count problem into a skew problem: one partition holds most of the rows and now has exactly one writer. Distribute by the partition columns plus a bucketing expression such as a hash of the id modulo N, so the hot partition gets N writers and N files instead of one, and small partitions still get one each.
- How do you decide whether a partition column is too fine-grained?Look at bytes per partition against your target file size. If typical partitions land well under a few hundred megabytes, the directory is not earning its keep. Coarsen the grain — day instead of hour — and rely on columnar min/max indexes and predicate pushdown for the finer filter, or use bucketing for high-cardinality join keys.
saying these in an interview costs you the question
- Assumes one output file per partition regardless of task count
- Raises max.dynamic.partitions.pernode without asking why it was hit
- Uses ORDER BY to cut file count, killing parallelism
- Thinks enabling merge alone removes the need to control writers
- Partitions by a high-cardinality column and blames the cluster