Menu ▾ ▴

#56 feat: ClickHouse cluster sharding by entity_id (Scale Architecture §2)

closed
nobody
None
2026-06-16
2026-06-16
Anonymous
No

Originally created by: grynn-in

Closes [#53]

Summary

Implements Scale Architecture §2 — ClickHouse cluster sharding by entity_id.
This is a BUILD-ONLY PR: the single-node default is unchanged and dbt parse
passes. End-to-end cluster correctness requires a real multi-node setup and is
explicitly called out below.


VERIFIED (on this single-node build)

Check Result
dbt parse ✅ passes — no ClickHouse connection needed
Default engine values unchanged ✅ cluster_engine('MergeTree()') → 'MergeTree()'; cluster_engine('ReplacingMergeTree(x)') → 'ReplacingMergeTree(x)'
cluster_name() default ✅ returns None → no ON CLUSTER DDL in single-node mode
cluster_enabled var default ✅ false — manifest confirms engine/cluster fields identical to pre-PR values
Three opted-in models ✅ SELECT query body unchanged; only config header updated
Single-node deploy path ✅ unaffected — all macros are pass-through when cluster_enabled=false

NEEDS-A-CLUSTER (not verified — requires multi-node ClickHouse + Keeper)

  • ON CLUSTER DDL correctness on real nodes
  • ReplicatedMergeTree ZooKeeper/Keeper path expansion and replica sync
  • Distributed table routing (cityHash64(entity_id) → correct shard)
  • dbt run-operation create_distributed_tables execution
  • Cross-shard aggregation correctness for gold models (gold_trial_balance, gold_consolidated_trial_balance, IC eliminations, NCI)
  • All 144 dbt tests green in cluster mode
  • Hot-shard analysis (large single LE under cityHash64)
  • Rebalancing strategy on new-LE onboarding

What's in this PR

dbt_project/macros/cluster.sql (new)

Core opt-in macros — all return unchanged values when cluster_enabled=false:

Macro Purpose
cluster_enabled() reads var('cluster_enabled', false)
cluster_name() 'konsol_cluster' or none
cluster_engine(base) converts MergeTree → ReplicatedMergeTree family
cluster_sharding_key() cityHash64(entity_id)
create_distributed_tables() run-operation: creates Distributed overlays over _local tables
drop_distributed_tables() run-operation: idempotent teardown

dbt project changes

  • dbt_project.yml: cluster_enabled: false var (default OFF, documented)
  • profiles.yml: new cluster target with cluster: konsol_cluster
  • 3 opted-in models (config block only; SELECT unchanged):
  • bronze_general_journal_account_entries — highest-volume GL entries
  • bronze_general_journal_entries — journal headers
  • gold_trial_balance — primary gold query surface

ClickHouse cluster infrastructure

  • clickhouse/cluster/remote_servers.xml — 3-shard konsol_cluster topology
  • clickhouse/cluster/keeper.xml — embedded ClickHouse Keeper + ZK-compat client
  • clickhouse/cluster/macros-shard1-replica1.xml — per-node {shard} / {replica} macros
  • docker-compose.cluster.yml — 3-node cluster compose overlay

Documentation

  • docs/admin-guide/cluster-setup.md (new) — step-by-step setup procedure, model opt-in pattern, open questions (hot-shard risk, rebalancing, cross-shard aggregation, Cube.js)
  • docs/prd/PRD-SCALE-ARCHITECTURE.md §2 — implementation status table, VERIFIED / NEEDS-A-CLUSTER callout

Activation (once a real cluster is available)

# Build with cluster engines + ON CLUSTER DDL
dbt build --target cluster --vars '{"cluster_enabled": true}'

# Create Distributed overlay tables
dbt run-operation create_distributed_tables --vars '{"cluster_enabled": true}'

Open questions (tracked in issue [#53])

  • Hot-shard: large single LE under cityHash64(entity_id) — composite key (entity_id, fiscal_year) may be needed
  • Rebalancing: accept hash-based skew on new-LE onboarding or schedule periodic reshard?
  • Cross-shard gold models: GLOBAL IN/GLOBAL JOIN requirements for IC eliminations and NCI aggregations on Distributed reads
  • Cube.js pre-aggregation: shard-aware refresh key needed?

Generated by Claude Code

Related

Tickets: #53
Tickets: #57

Discussion

  • Anonymous

    Anonymous - 2026-06-16

    Originally posted by: pyy3

    🔍 Review round (2 passes: no-regression correctness + design/docs/tests)

    The single-node default path is provably SAFE — verified against dbt-clickhouse 1.10: cluster_enabled defaults false and never leaks true; cluster_engine() returns its input verbatim when off; cluster_name()→none ⇒ no ON CLUSTER; profiles dev/prod byte-identical; touched models' SELECT bodies unchanged. Merging is safe for today's deploy.

    But the cluster feature as written would NOT run on a real cluster:

    BLOCKER — shard key cityHash64(entity_id) references a column that doesn't exist. Bronze renames entity_id→data_area_id; no materialized bronze/silver/gold table outputs entity_id (it's only in the unmaterialized staging/canonical/*). create_distributed_tables would fail at the first CREATE ... Distributed(..., cityHash64(entity_id)). The PR's own admin-guide Step 7 verify uses WHERE data_area_id = 'USSI' — contradicting the entity_id key. Fix: shard on data_area_id everywhere.

    MAJOR — model-level config(cluster=...) is dead under dbt-clickhouse 1.10 (it drives ON CLUSTER from the profile cluster: key, not model config); plain incremental/table makes no _local table, so the run-operation that assumes <model>_local is inconsistent. Recommend the adapter's native distributed_table/distributed_incremental materializations.

    MAJOR — create/drop lists diverge (55 vs 26) and ~52 of 55 have no _local backing (only 3 models are cluster-aware). Derive from the graph / a cluster tag.

    MAJOR — "default unchanged" is prose, never codified. Add a unit/run-operation self-check of cluster_engine() pass-through with the flag off.

    MINOR — keeper.xml hardcodes server_id=1 but compose mounts it on all nodes → broken quorum; macros-shard{2,3} referenced but absent.

    I'll fix the unambiguous bugs (shard key → data_area_id, keeper server_id, create/drop symmetry, codified default-unchanged test) and flag the materialization rework as the real path to a working cluster — that part needs a multi-node cluster to validate.

     
  • Anonymous

    Anonymous - 2026-06-16

    Originally posted by: grynn-in

    ✅ Fixes applied (b3e0523); verified dbt parse OK + both run-operations no-op cleanly with cluster off (single-node default unchanged):

    • BLOCKER fixed — shard on data_area_id (the LE column materialized across bronze/silver/gold) instead of the non-existent entity_id. Updated macro defaults, remote_servers.xml, compose, and the docs/AC verification queries.
    • create/drop unified into one cluster_sharded_tables() source of truth, scoped to the 3 actually-cluster-aware models (those carrying cluster_engine/cluster_name → dbt builds a _local table). This fixes both the 55-vs-26 divergence and the "CREATE over a non-existent _local" failure for the ~52 non-cluster-aware tables.

    Still NEEDS-A-CLUSTER (deferred to your live verification — can't be done/validated here):

    • B2 — materialization rework. Model-level config(cluster=...) is inert under dbt-clickhouse 1.10 (it drives ON CLUSTER from the profile cluster: key); plain incremental/table makes no _local suffix. The clean path is the adapter's native distributed_table/distributed_incremental materializations — flagged in the macro + PRD §2. This needs a real cluster to design+validate.
    • Keeper server_id — keeper.xml hardcodes server_id=1 while compose mounts it on all nodes → broken quorum; needs per-node config/templating.

    Recommendation: the single-node default is provably safe to merge, but do not enable cluster mode until B2 + Keeper are reworked against a real multi-node cluster.

     
  • Anonymous

    Anonymous - 2026-06-16

    Ticket changed by: grynn-in

    • status: open --> closed
     

Log in to post a comment.