Skip to content
AITroveRead. Build. Understand.
Make this comfortable

Partition planning, pruning and skew

Last updated: 6 Oct 20265 min read
tutorial
AdvancedBy AITrove Editorial

A partition spec groups rows by a transform so readers can skip unrelated files; its value depends on query predicates and the distribution of rows.

Start with the dominant read

A support dashboard usually filters orders by event day and sometimes region. Day partitions can bound the files considered for a seven-day report. Partitioning by individual order ID would create nearly one partition per order, making metadata and file counts unmanageable. A partition key must reduce scans without exploding the number of small groups.

Use event time deliberately

Landing date and business event date are different. Late events received today may belong to an event-day partition from last month. If consumers ask for historical order date, partition on that date and include old partitions in the correction plan. Late-arrival policy should define how long the system rewrites earlier partitions.

Inspect skew before scaling

A region named 'unknown' may contain 68% of all orders. Giving it one processing task while smaller regions each get one task leaves most workers idle. Calculate counts and bytes by partition and by join key. Salted keys, repartitioning, or a separate heavy-key path may help; each changes shuffle cost and must preserve the output key and count.

Keep partition evolution readable

A table can move from day partitions to month plus bucket partitions as data and access patterns change. Readers must still find older data laid out under the earlier spec. A transaction-aware table catalog can track which spec produced each file. Without one, a migration needs a versioned manifest and a read path aware of both layouts. Snapshot commits isolate the switch.

Measure the trade

Compare the fraction of files pruned for real queries, total files, worst partition size, write amplification from corrections, and p95 query latency. A spec that reduces a full scan but creates 190,000 tiny files may be a net loss. Compaction can reduce file count, but cannot fix a partition key that does not match filters.

Implementation

python
from collections import Counter

orders = [
    {"day": 47, "region": "north"}, {"day": 47, "region": "unknown"},
    {"day": 47, "region": "unknown"}, {"day": 48, "region": "south"},
    {"day": 49, "region": "unknown"},
]

def partition_counts(records):
    counts = Counter((record["day"], record["region"]) for record in records)
    return dict(counts)

counts = partition_counts(orders)
assert counts[(47, "unknown")] == 2
assert sum(counts.values()) == len(orders)

Performance and operating cost

Counting N rows is O(N) time and O(P) memory for P observed partitions. A finer partition spec can reduce read bytes but raises metadata, file-opening, and correction costs. Benchmark the full workload rather than a single selective query.

Common Mistakes

  • Do not partition by a near-unique identifier.
  • Do not use ingest day when users filter by event day without accounting for late data.
  • Do not report only average partition size when one key dominates.

Read next

Continue the workflow: Shuffle boundaries and local combiners.

ai-data
data-engineering
Storage details