Optimizing Athena Queries on Large Transactional Tables
When querying very large transactional datasets in Amazon Athena, performance challenges often arise as soon as you introduce joins — especially when joining to a secondary table that may contain zero to many related records. Even if the secondary table is small or has very few columns, query time and data scanned can increase dramatically.
This post walks through why that happens, what doesn’t work, and which strategies do improve performance in Athena — including when to use partitioning, when to avoid it, and when Iceberg tables are the right solution.
The Core Problem
Athena is a distributed, columnar query engine built on Presto/Trino. It is optimized for scanning large amounts of data efficiently — not for row-by-row lookups or correlated subqueries.
A common pattern that causes trouble looks like this:
- A very large transactional table (hundreds of millions of rows)
- A secondary table used only to check whether a related record exists
- A join key with very high cardinality
- No additional filtering available for the secondary table
Even if the join is logically simple (“does a related record exist?”), Athena may still scan all files from both tables, resulting in high data scanned and slow execution.
Why Correlated Subqueries Hurt Performance
Using a correlated subquery to check for existence seems efficient, but in Athena it often forces a full scan of the secondary table:
Athena cannot use indexes or point lookups. Instead, it must evaluate the subquery in a distributed context, which often results in:
- Full table scans
- Distributed data shuffles
- High latency and cost
Use Joins or Semi-Joins Instead
A better approach is to replace correlated subqueries with:
- A LEFT JOIN to a deduplicated lookup table, or
- A SEMI JOIN / IN clause when you only need existence
Example pattern:
This gives Athena more freedom to optimize execution.
Partitioning: What Helps and What Doesn’t
❌ Don’t Partition on High-Cardinality Keys
Partitioning by a join key with millions of distinct values (like an object or transaction ID) is a common mistake.
Why this fails:
- Athena struggles with very large numbers of partitions
- Glue metadata overhead becomes extreme
- Query planning time increases
- Overall performance gets worse
As a rule of thumb: never partition on a key with millions of distinct values.
✅ Partition on Fields You Always Filter By
Partitioning is extremely effective when done correctly. If every query includes filters like:
- Organization or company
- Fiscal year
- Accounting period or month
Then partitioning on all of those fields together is ideal.
This allows Athena to:
- Skip entire partitions during query planning
- Read only a small subset of data
- Dramatically reduce data scanned
For example:
PARTITIONED BY (company, fiscal_year, accounting_period)
As long as the total number of partition combinations stays reasonable (typically under tens of thousands), this approach works very well.
Precomputing Flags Instead of Re-Joining
If you frequently need to know whether a related record exists, the most efficient approach is often to materialize that information once.
Instead of joining every time:
- Build a new table that includes the original transactional data
- Add a derived boolean or integer flag indicating existence
- Partition that table using the same required filters
Once materialized, future queries:
- Require no joins
- Scan fewer columns
- Run faster and cost less
Compression Matters (Yes, Use Snappy)
For Athena and Parquet:
- Snappy compression is the best default
- It provides fast decompression
- It’s fully supported by Athena, Glue, and Iceberg
- It balances performance and storage cost
Avoid GZIP for interactive queries — it compresses well but slows down reads significantly.
When Iceberg Is the Better Choice
Traditional Hive-style tables work well for static data, but when your data evolves or you want more flexibility, Iceberg tables are often the better option.
Iceberg advantages:
- ACID transactions
- Built-in metadata management
- No MSCK REPAIR TABLE
- Schema evolution
- Partition evolution
- Efficient UPDATE, DELETE, and MERGE
Creating an Iceberg table with Athena typically looks like:
For large transactional datasets that are queried repeatedly and enriched over time, Iceberg often delivers better long-term performance and maintainability.
Key Takeaways
- Athena scans data, it doesn’t “look up” rows
- Correlated subqueries are expensive — avoid them
- Never partition on high-cardinality IDs
- Partition on fields that appear in every query filter
- Materialize frequently used existence checks
- Use Parquet + Snappy
- Consider Iceberg for large, evolving datasets
Final Thought
Athena performance is less about SQL cleverness and more about data layout. Once your tables are partitioned correctly, compressed properly, and structured for how they’re queried, performance issues often disappear — even at massive scale.
If you design your storage with the query patterns in mind, Athena can comfortably handle hundreds of millions (or billions) of rows without breaking a sweat.
Retiring Lawson, PeopleSoft, or Oracle? APIX archives the entire application — every table, every year, attachments and security included — into your own AWS account in about 30 days, so you can decommission the legacy system and keep full access to the history.





