After Indexing is Correct: Hash Join Spills, Splitting Joins Without N+1 & 4 UPSERT Traps
When a query runs slowly, the standard engineering instinct is predictable: run EXPLAIN, check if the type column indicates a full table scan (ALL), and add an index to the filtered column.
Yet in high-scale production systems, situations regularly occur where the index is correctly designed, the execution plan confirms the index is chosen, and yet execution latency stalls between 3 and 8 seconds.
Why do queries remain slow when indexing is correct? When should complex multi-table SQL queries be deconstructed into application-tier joins rather than forced onto the database engine? And why does the standard INSERT ... ON DUPLICATE KEY UPDATE (UPSERT) pattern secretly burn through primary key sequences and produce baffling deadlocks?
This article explores high-level query optimization techniques once baseline indexing has already been perfected.
-- Table has 1 row (id=1, email='[email protected]') INSERT INTO users (email, name) VALUES ('[email protected]', 'Alice') ON DUPLICATE KEY UPDATE name = VALUES(name); -- Next inserted row receives id=3 (id=2 was burned!)
1. When Indexing is Correct but Queries Crawl: The Hash Join Disk Spill
In modern database engines (MySQL 8.0.18+ and PostgreSQL), when joining large datasets or when indexes cannot be utilized, the optimizer selects the Hash Join algorithm:
The Two Phases of a Hash Join:
Phase 1: Build Phase (Construct In-Memory Hash Table)
βββ Reads smaller input (Build Table) ββ> Populates Hash Table in memory:
βββ If size <= join_buffer_size: Fits entirely in RAM (O(1) memory lookup)
βββ If size > join_buffer_size: π₯ DISK SPILL (Partitions written to temporary disk files)!
Phase 2: Probe Phase (Scan Larger Table)
βββ Reads rows from larger input (Probe Table) ββ> Probes hash partitions to return matches.
Starting in MySQL 8.0.18 and standard in PostgreSQL, the optimizer uses Hash Joins for equi-joins when no index is available.
Phase 1 (Build): Reads the smaller table and constructs an in-memory hash table keyed by the join attribute.
Phase 2 (Probe): Reads the larger table and probes each row against the hash table in memory in $O(1)$ time.
The Disk Spill Trap: If the build table exceeds join_buffer_size (MySQL) or work_mem (PostgreSQL), the hash table spills into temporary disk files, degrading latency from 20ms to 8 seconds!
-- Look for hash join spills in EXPLAIN ANALYZE: -> Inner hash join (orders.user_id = users.id) (actual time=0.45..12.3 rows=50000) -- Warning if spilling to disk: Batches: 5, Memory Usage: 32768kB (Disk Spill: Yes)
The Memory Threshold: join_buffer_size & work_mem
- When the build dataset exceeds the allocated join memory buffer (
join_buffer_sizein MySQL orwork_memin PostgreSQL): - The engine cannot retain the hash table in memory. It partitions the dataset into chunks and spills them to temporary disk files.
- The algorithm transforms from an in-memory lookup into multi-pass random disk I/O.
- The Result:
EXPLAINreports a modern Hash Join, but the query spends 8 seconds waiting on disk I/O.
Production Remediation
- Inspect with
EXPLAIN ANALYZE: Look for metrics indicating disk spills (e.g.,Batches: 5, Memory Usage: ... (Disk Spill: Yes)). - Evaluate Index Nested Loop Join: Ensure the join column on the probed table has a covering index so the engine can utilize an Index Nested Loop Join () instead of scanning and hashing entire tables.
- Elevate Session Buffers: For dedicated reporting queries, elevate memory for that session:
SET SESSION join_buffer_size = 268435456; -- 256 MB
2. The Art of Splitting Multi-Table JOINs Without Creating N+1
Developers frequently write monolithic queries joining 4 to 6 tables (JOIN users JOIN orders JOIN order_items JOIN products JOIN payments...) believing that "letting the database do everything in one round-trip is always fastest."
In high-throughput microservices architectures, massive multi-table joins are actively discouraged on operational paths. Instead, high-scale platforms employ Application-Side Joins:
// Step 1: Fetch top 50 orders (1 roundtrip)
List<Order> orders = orderRepo.findByUserId(userId, limit(50));
List<Long> orderIds = orders.stream().map(Order::getId).toList();
// Step 2: Batch fetch all related items (1 roundtrip)
List<OrderItem> items = itemRepo.findByOrderIdIn(orderIds);
// Step 3: In-memory hash mapping (0 roundtrips)
Map<Long, List<OrderItem>> itemMap = items.stream()
.collect(groupingBy(OrderItem::getOrderId));1. No Cartesian Products: Prevents transferring redundant order columns for every child item row over the wire.
2. Cache-Friendly: Each entity (Order, Item) can be independently cached in Redis with high cache hit ratios.
3. Microservice-Ready: Works seamlessly when `orders` and `order_items` are partitioned across different database shards.
Vulnerabilities of Complex Multi-Table JOINs
- Network Bloat: Joining
orderswithorder_itemsduplicates order-level metadata (shipping addresses, customer names) across every single line item sent over the network. - Cache Invalidation Coupling: Caching the unified result of a 5-table join means that an update to any single product description invalidates the cached representation of all associated orders.
- Sharding Incompatibility: When tables are partitioned across separate physical databases (
ordersin Shard-1,usersin Shard-2), cross-database SQL joins become physically impossible.
Enterprise Pattern: Batch Fetching (2 Queries, Zero N+1)
Avoid issuing single-record queries in loops ( queries). Use bounded Batch Ingestion:
// β
Step 1: Fetch 50 orders (1 Network Roundtrip)
List<Order> orders = orderRepository.findByUserId(userId, PageRequest.of(0, 50));
List<Long> orderIds = orders.stream().map(Order::getId).toList();
// β
Step 2: Fetch all child items for all 50 orders at once (1 Network Roundtrip)
List<OrderItem> allItems = orderItemRepository.findByOrderIdIn(orderIds);
// β
Step 3: Stitch together in application memory using a HashMap (O(1) CPU, 0 Disk I/O)
Map<Long, List<OrderItem>> itemsByOrderId = allItems.stream()
.collect(Collectors.groupingBy(OrderItem::getOrderId));
orders.forEach(order -> order.setItems(itemsByOrderId.getOrDefault(order.getId(), List.of())));
- Query Overhead: Exactly 2 simple queries.
- Each query hits primary keys or indexed foreign keys, completing in .
- Each entity can be cached independently in Redis with a cache hit ratio exceeding 95%.
3. The Deterministic Pattern: SELECT id Then Write by ID
When executing bulk updates on expired or pending records, developers often write:
-- β DANGEROUS: Causes unconstrained gap locks and frequent deadlocks
UPDATE orders
SET status = 'EXPIRED'
WHERE status = 'PENDING' AND created_at < NOW() - INTERVAL 1 DAY
LIMIT 1000;
Naive batch updates using range conditions (e.g. UPDATE orders SET status = 'EXPIRED' WHERE created_at < NOW() - INTERVAL 7 DAY LIMIT 1000) acquire arbitrary, non-deterministic gap locks that routinely deadlock against concurrent user writes.
SELECT id FROM orders WHERE created_at < ? LIMIT 1000;Zero write locks acquired.
orderIds.sort(naturalOrder());Guarantees every thread acquires locks in identical ascending order.
UPDATE orders SET status = 'EXPIRED' WHERE id IN (:sortedIds);Locks exact rows, 0 gap locks, 0 deadlocks.
Why This Statement Causes Deadlocks
- Evaluating a range condition on
created_atcombined with aLIMITclause causes InnoDB to acquire Next-Key Locks across unpredictable index gaps. - When multiple background workers or user checkout transactions touch adjacent rows, these gap locks conflict, causing immediate Deadlock 1213 errors.
The Safe 3-Step Pattern: Deterministic ID Mutation
- Read IDs in Read-Only Mode:
SELECT id FROM ordersWHERE status = 'PENDING' AND created_at < :expiredTimeLIMIT 1000;
- Sort IDs in Application Memory:
Sort the retrieved
orderIdsin strictly ascending numerical order (Collections.sort(ids)). This guarantees every worker thread requests row locks in identical physical sequence. - Execute Update by Primary Key:
InnoDB acquires Record Locks only on Primary Key records. Zero Gap Locks are generated, completely eliminating deadlock risk.UPDATE orders SET status = 'EXPIRED' WHERE id IN (:sortedIds);
4. The Four Deadliest Traps of ON DUPLICATE KEY UPDATE
The INSERT ... ON DUPLICATE KEY UPDATE (UPSERT) construct is widely used to handle insert-or-update flows. However, it introduces four critical production traps:
-- Table has 1 row (id=1, email='[email protected]') INSERT INTO users (email, name) VALUES ('[email protected]', 'Alice') ON DUPLICATE KEY UPDATE name = VALUES(name); -- Next inserted row receives id=3 (id=2 was burned!)
Trap 1: Auto-Increment ID Burn & Sequence Exhaustion
Consider a table tracking user page counters:
CREATE TABLE user_counters (
id INT AUTO_INCREMENT PRIMARY KEY,
user_id INT NOT NULL UNIQUE,
visit_count INT DEFAULT 1
);
Executing:
INSERT INTO user_counters (user_id, visit_count) VALUES (42, 1)
ON DUPLICATE KEY UPDATE visit_count = visit_count + 1;
- Physical Engine Behavior: Before checking if
user_id = 42exists, InnoDB always allocates the next Auto-Increment ID from the global table counter! - When the duplicate is detected and an
UPDATEoccurs, the allocated ID is discarded and permanently lost. - On a high-throughput counter updating millions of times daily, an
INTprimary key (maximum value 2.1 billion) will be completely exhausted within months, triggering fatalDuplicate entry for key 'PRIMARY'errors that take down the service!
Trap 2: Insert Intention vs Gap Lock Deadlocks
When an UPSERT detects a conflict on a unique secondary index:
- InnoDB must acquire an exclusive X Next-Key Lock on the duplicate record.
- When concurrent transactions perform upserts into adjacent ranges, their Gap Locks and Insert Intention Locks create cyclic wait dependencies, producing immediate deadlocks.
Trap 3: Replication Data Drift Under Statement-Based Replication
If using binlog_format = STATEMENT:
- Non-deterministic functions (e.g.,
NOW()) or statements evaluated across multiple unique keys can be applied in different physical index orders on primary vs replica servers. - Data drifts silently between primary and replica nodes.
- Mandate: Always enforce
binlog_format = ROWin production MySQL deployments.
Trap 4: Ambiguity with Multiple Unique Indexes
Consider a table with both a Primary Key (id) and a Unique Key (email):
-- Existing rows:
-- Row 1: (id = 1, email = '[email protected]')
-- Row 2: (id = 2, email = '[email protected]')
ON DUPLICATE KEY UPDATE name = 'Hacker';
- The statement conflicts on
id = 1with Row 1, but also conflicts onemail = '[email protected]'with Row 2! - MySQL cannot update both rows. It updates ONE arbitrary row depending on which index the engine scans first!
- This causes silent, undetected data corruption and security breaches.
