Partitioning a table by a high-cardinality column like user_id creates an operational disaster known as over-partitioning, generating millions of nested directories that hold tiny files only a few kilobytes in size. Query engines struggle under severe metadata overload because listing millions of partition paths exhausts driver memory and inflates query planning times from milliseconds to several minutes. At the same time, metastores and object storage rate limits buckle under the load, while analytics queries that do not filter by that specific user ID suffer severe performance degradation.
The small files and metadata storm
In cloud object storage and distributed systems, partitioning creates a distinct physical prefix for every unique column value. When applied to high-cardinality keys, severe issues emerge:
- Filesystem listing saturation: S3 and GCS return directory listings in batches of 1,000 keys. If a table has 5,000,000 partition directories, scanning the table requires thousands of recursive HTTP API calls, triggering cloud throttling errors and taking tens of minutes before query execution even begins.
- Metastore exhaustion: Hive Metastore or cloud catalogs must track partition metadata in their backend relational databases. Storing millions of partitions exhausts catalog memory and causes connection timeouts for all team members.
- Driver out-of-memory crashes: The Spark driver must hold file split metadata for every tiny file in its memory. Creating millions of tiny files generates gigabytes of JVM heap allocations, frequently crashing the driver with OutOfMemory errors.
Query planning and driver memory strain
Most analytical queries aggregate across thousands or millions of users rather than filtering on a single user ID:
- When an analytical query runs without a user_id filter, the engine must open every single partition folder.
- Opening millions of tiny files destroys columnar read benefits because vectorized execution loops and dictionary encodings fail on small blocks.
- The query spends 95 percent of its duration reading metadata rather than processing data rows.
High-cardinality alternatives
As a practical rule of thumb, each partition should contain at least 1 GB of data to justify directory separation. For high-cardinality columns that users filter on frequently, avoid Hive directory partitioning entirely:
- Rely on modern table format clustering (such as Delta liquid clustering or Snowflake clustering) to co-locate similar IDs inside larger files.
- In legacy frameworks, use bucketing or Z-ordering within coarse-grained date partitions (such as partitioning by event_date and clustering by user_id).