Choose join strategies from cardinality and skew evidence
Join Strategies and Dimensional Modeling Benefits: Broadcast/Hash/Merge Concepts, Cardinality, and Skew
Reason about join execution from cardinality, selectivity, uniqueness, and skew while preserving dimensional grain.
Learning outcomes
Distinguish broadcast, hash, and merge join concepts from the specific nested-loop/index-search behavior observed in SQLite.
Use cardinality, uniqueness, selectivity, and key distribution to predict join work and fanout risk.
Quantify skew and explain why dimensional models often make dimension-side joins easier without making them automatically cheap.
Reject forced join strategies or hints when the engine/data evidence does not support them.
Preserve fact grain and metric totals while changing physical join execution.
Chapter 26 begins from the accepted AtlasMart state established through Chapters 01–25: 10 current paid lines, 8 orders, 12 units, 820 USD gross revenue, 495 USD cost, and 325 USD gross profit. Source progress remains committed through sequence 208; Chapter 20 metric contracts, Chapter 21 certified marts, Chapter 22 security controls, Chapter 23 lineage/ownership, Chapter 24 tests, and Chapter 25 SLO/incident practices remain authoritative. The 360,000-row benchmark below is a clearly labeled scale fixture for performance mechanics only; it never replaces these business control totals.
Runtime: Python 3.13.5 + SQLite 3.46.1. Benchmark grain: one synthetic paid sales line per row. Scale fixture: 360,000 rows across 180 dates, 100 products, and 5,000 customers; product 1 intentionally receives about 30% of rows and North receives about 50% of customer-linked rows. Storage: a local SQLite database, WAL journal mode, temporary structures configured in memory. Timing: 3 warm-up executions then 17 measured executions; reported cache state is warm. The operating-system cold cache is not forcibly cleared, so no “cold-cache” claim is made. Security: synthetic identifiers only. Distributed limitations: SQLite does not expose distributed shuffle bytes, warehouse queues, or hardware SIMD counters; those mechanisms are explained conceptually and, where useful, modeled explicitly rather than mislabeled as SQLite measurements.
1. Realistic problem: one hot product owns a worker
AtlasMart’s distributed-cloud migration team notices that product-level joins would likely repartition by product key. The local scale fixture intentionally makes product 1 account for 30.08% of fact rows. In a hash-distributed system, a naive hash on product can send much more work to one partition. The business risk is high p99 latency and underused parallel workers even when total CPU capacity looks ample.
2. Broadcast, hash, and merge are execution mechanisms, not modeling patterns
| Join concept | Mechanism | Good evidence | Boundary/failure |
|---|---|---|---|
| Broadcast | Replicate a small input to workers holding partitions of a large input. | Dimension is demonstrably small enough; avoids repartitioning the fact. | Large or rapidly growing dimension can exceed memory/network budget. |
| Hash | Hash rows by join key, then match equal-key buckets. | Equality join, adequate memory, balanced key distribution. | Hot keys create skew; poor cardinality estimates can size/build the wrong side. |
| Merge | Read both inputs in compatible key order and advance through them. | Inputs are already sorted/ordered or sorting cost is justified. | Sorting can dominate; duplicate keys can still multiply rows. |
| SQLite local plan | Nested scans/searches selected by SQLite planner. | EXPLAIN QUERY PLAN shows SCAN/SEARCH and index use. | Do not describe this result as a distributed broadcast/hash/merge join. |
3. Cardinality comes before strategy
Cardinality is the number of rows or distinct key values relevant to an operation. Before choosing a join strategy, establish: fact rows; dimension rows; uniqueness of dimension keys; selectivity of filters; null/orphan rates; and key-frequency distribution. A dimension with one row per product can safely join many fact rows because its key is unique. A “dimension” with duplicate product keys silently multiplies fact rows and corrupts revenue no matter how fast the join algorithm is.
SELECT product_sk, COUNT(*) AS nFROM dim_productGROUP BY product_skHAVING COUNT(*) <> 1;-- expected: no rowsSELECT COUNT(*) AS orphan_factsFROM fact_sales fLEFT JOIN dim_product p ON p.product_sk = f.product_skWHERE p.product_sk IS NULL;-- expected: 0
4. Executed skew evidence
| Distribution evidence | Observed value |
|---|---|
| Benchmark fact rows | 360,000 |
| Top product rows | 108,288 |
| Top product share | 30.08% |
| Naive 8-way product-hash partition rows | 30677, 138663, 33019, 33180, 32824, 30646, 30661, 30330 |
| Naive max/mean skew ratio | 3.081 |
| Simulated salted-hot-key rows | 44213, 43911, 46555, 46716, 46360, 44182, 44197, 43866 |
| Simulated salted max/mean ratio | 1.038 |
The salting calculation is a simulation, not an SQLite feature. It demonstrates the mechanism: splitting one hot key across several execution buckets can reduce processing skew, but doing so complicates grouping/join semantics because partial results must be recombined. Never salt business keys in the logical dimensional model merely to satisfy an execution engine.
5. Dimensional modeling helps, but does not guarantee a fast join
A star schema often provides compact, unique-key dimensions and a large fact at one explicit grain. That shape gives planners useful options: a small dimension may be broadcast in a distributed engine, and primary-key lookups are cheap in row stores. But dimensions can become wide, SCD Type 2 history can grow, many-to-many bridges can multiply rows, and skew can dominate. “Star schema = broadcast join” is folklore, not a rule.
6. Local plan evidence and its limits
SEARCH f USING INDEX ix_fact_sales_date_customer (order_date>? AND order_date<?)SEARCH c USING INTEGER PRIMARY KEY (rowid=?)SEARCH p USING INTEGER PRIMARY KEY (rowid=?)USE TEMP B-TREE FOR GROUP BYUSE TEMP B-TREE FOR ORDER BY
The plan shows an indexed date-range search into the fact followed by primary-key searches into customer and product dimensions. SQLite is single-node, so there is no network shuffle. In a distributed engine, the same logical SQL might become broadcast, hash, merge, or another engine-specific strategy depending on table statistics and resource policy.
7. Controlled failure: force a strategy because it was fast elsewhere
Copying a “broadcast every dimension” recommendation from another engine ignores dimension size, memory, concurrency, statistics, and implementation details. Likewise, forcing join order or hints from one benchmark can regress another workload after data grows. The safe repair is evidence: validate uniqueness/cardinality, profile key frequencies, examine the current engine plan, measure representative workload percentiles, then keep a rollback path.
What if the smallest table is not the best build/broadcast side?
Show answer
A highly selective predicate on the fact can make the filtered fact smaller than a nominally small dimension. Runtime cardinality after filters, not table labels alone, should drive the decision.
8. Boundary: pre-aggregate only when grain remains explicit
Pre-aggregating fact rows before a dimension join can reduce rows, but only when the requested metric is additive across the removed detail and the grouping keys preserve every required slice. Distinct customers, weighted bridge allocations, and non-additive ratios can change if you aggregate too early. Performance rewrites inherit all grain/metric rules from Chapters 02–09 and Chapter 20.
9. Production judgment and bridge
Accept a performance change only after result semantics, security policy, freshness/history behavior, retry/replay behavior, and reconciliation still pass. Performance evidence must name the workload, dataset, engine/version, cache state, concurrency, physical layout, and percentile—not just “faster.” Treat any cost or latency number in this chapter as fixture-specific, not a universal target. Keep rollback simple: indexes/materializations can be removed, query rewrites reverted, and workload-pool assignments restored while the logical model and governed metric contracts remain unchanged.
Next: Partition/Cluster Pruning, Predicate Pushdown, Materialization, Approximation, and Query Rewrites.
10. Verification checklist
- Dimension join keys remain unique at the declared join grain.
- Fact/dimension orphan controls remain zero or use an explicit unknown-member policy.
- Hot-key distribution is measured, not guessed.
- SQLite plans are described as SQLite plans, not distributed join algorithms.
- Any skew-mitigation technique is kept in the physical execution layer and reconciled to unsalted business keys.
Knowledge check
Check your understanding
- What is the central mechanism in “Join Strategies and Dimensional Modeling Benefits: Broadcast/Hash/Merge Concepts, Cardinality, and Skew”, and which AtlasMart grain or metric contract must remain unchanged?
- Which observable evidence in this lesson distinguishes the correct design from the controlled failure?
- Which assumptions are local or engine-specific, and what must be re-checked before production use?
Review the answers
1. Preserve the lesson’s declared business grain, history semantics, governed metric definitions, and reconciliation controls while changing only the mechanism under study.
2. Use the lesson’s counts, sums, checksums, plans, traces, timing/cost calculations, or failure-state evidence—not a green task status or naming convention alone.
3. Re-check runtime/version, data scale and distribution, cache/concurrency, storage layout, security context, pricing/region where relevant, and the exact product guarantees before production adoption.
Authoritative references
- SQLite — EXPLAIN QUERY PLANOfficial interpretation of scan/search operators, index use, temporary B-trees, and the warning that plan-output format is not a stable application API.
- SQLite — Query PlanningOfficial background on table scans, multi-column indexes, sorting, and why the planner chooses among semantically equivalent algorithms.
- SQLite — EXPLAINOfficial semantics and limitations of EXPLAIN/EXPLAIN QUERY PLAN used by the local lab.
- Python — statisticsUsed for deterministic latency summaries in the local benchmark harness.
- BigQuery — Understand reservationsNon-prerequisite vendor example showing how a cloud warehouse can isolate workloads with resource pools; the Chapter 26 concepts do not require BigQuery.
11. Lab cleanup/reset
The mandatory lab is local and synthetic. Delete
ch26_lab/atlasmart_perf.sqlite and the generated
benchmark JSON, then rerun the setup script to restore the
deterministic 360,000-row fixture. No cloud resources, accounts,
paid services, or production credentials are created.