Cross-Index JOIN
Stop ETL’ing Elasticsearch into your warehouse just to JOIN it.
Elasticsearch has no native cross-index JOIN — and a stock JDBC driver on top of ES can’t add one either. SoftClient4ES gives you two things Elasticsearch can’t do on its own:
- Query-time cross-index JOIN —
INNER/LEFT/RIGHT/FULL OUTERjoins across indices (and across clusters), at query time, on every surface: the REPL, the JDBC driver, the ADBC driver, the Arrow Flight SQL sidecar, and Federation. The drivers are the free delivery channel; the JOIN depth is metered. - Persisted Materialized Views — the other superpower: denormalize cross-index data once into a queryable view. (That has its own page; this page is about query-time JOINs.)
Runtime requirement — Java 11+. The JOIN engine is built on Apache Arrow 18.x, which ships Java-11 bytecode, and the 0.20+ REPL/driver line also bundles logback 1.5.x (Java-11 as well). Every JVM-side surface (the REPL, the JDBC/ADBC drivers) requires Java 11 or newer at runtime (17+ for ES 9); the Flight SQL and Federation Docker images bundle their own Java 17.
The JOIN ladder
Three meters gate how far a JOIN can reach. Think of them as rungs:
- Query-time depth —
maxJoins— how many cross-index JOINs a single query may contain. - Cross-cluster reach —
maxClusters— how many Elasticsearch clusters a Federation deployment may span. - Persisted —
maxMaterializedViews— how many denormalized views you may keep materialized.
Every tier has every feature. The meters gate scale, not on/off switches — see Licensing & the meters.
Which row do I need?
Cross-index JOIN ships in three shapes (“rows”). Pick by where your data lives:
| Row | Name | What it does | Engine | Surfaces |
|---|---|---|---|---|
| Row 1 | Passthrough / same-cluster | Cross-index JOIN of 2+ indices in one ES cluster | Each leg is an ES sub-query; the JOIN runs in-process in embedded DuckDB | REPL, JDBC, ADBC, Flight SQL sidecar |
| Row 2 | Cross-cluster conveyor | JOIN / INSERT / CTAS where source and target live in different clusters | Coordinator runs the source SELECT and conveyors the Arrow stream to the target sidecar for bulk-load | Federation |
| Row 3 | Multi-source coordinator | JOIN across 2+ source clusters | Coordinator stages each leg to Parquet scratch + a per-query DuckDB view, joins coordinator-local | Federation |
One sentence to decide: same cluster → Row 1; copy results between two clusters → Row 2; join across two or more clusters → Row 3.
The rule the engine actually applies: it counts the distinct source catalogs in the rewritten FROM / JOIN clauses versus the target catalog. Same (or no) catalog → Row 1; exactly one source catalog different from the target → Row 2; two or more source catalogs → Row 3.
What does NOT work yet: Cross-index JOINs are first-class in this release, but two things are intentionally not here yet: arbitrary subqueries / CTEs in a JOIN query land in the next release (Quarter 4 2026), and heterogeneous Row-3 sources (joining ES with Postgres, MySQL, Snowflake, …) land in the upcoming release (Quarter 1 2027) — this release’s Row 3 is multi-Elasticsearch only. See Known limitations for the full list.
Row 1 — same-cluster passthrough
A Row-1 JOIN reaches two or more indices in one Elasticsearch cluster. Each table in the query becomes its own ES sub-query (with any single-table WHERE pushed down into it); the cross-index JOIN itself is executed in-process by an embedded DuckDB engine inside the driver, sidecar, or REPL. The first JOIN on a fresh process warms up the DuckDB native library once.
Row 1 — same-cluster passthrough: the embedded DuckDB engine fans each table out to its own ES sub-query in one cluster, then joins the results in-process before returning rows to the client.
All examples below are transcribed verbatim from the SoftClient4ES JDBC integration test suite.
They run against two fixtures: jdbc_join_emp (emp_id, dept_id, name, salary — 6 rows: Alice/1/1/6000, Bob/2/1/4000, Carol/3/1/8000, Dave/4/2/5500, Eve/5/2/3000, Orphan/6/99/4500) and jdbc_join_dept (dept_id, dept_name — Engineering=1, Marketing=2, Empty=9). The orphan employee (dept_id=99) and the empty department (dept_id=9) surface the outer-join NULLs.
A few examples use a separate small fixture jdbc_test (id, name, value — a 3-column table, not the employees table; it has no salary/dept_id).
SELECT — INNER / LEFT / RIGHT / FULL OUTER
INNER JOIN (the canonical first example):
SELECT e.name, e.salary, d.dept_nameFROM jdbc_join_emp eJOIN jdbc_join_dept d ON e.dept_id = d.dept_id;-- 5 rows (orphan employee with dept_id=99 dropped by the INNER JOIN)LEFT JOIN with NULLs — unmatched left rows surface NULL on the right (this one uses the jdbc_test fixture):
SELECT t.id, t.name, d.dept_nameFROM jdbc_test tLEFT JOIN jdbc_join_dept d ON t.id = d.dept_idORDER BY t.id ASC;-- only ids 1→Engineering and 2→Marketing match; the rest get NULL dept_nameRIGHT OUTER and FULL OUTER:
-- RIGHT OUTER: dept side fully preserved; empty dept 9 → NULL e.nameSELECT e.name, d.dept_nameFROM jdbc_join_emp eRIGHT OUTER JOIN jdbc_join_dept d ON e.dept_id = d.dept_idORDER BY d.dept_name ASC;
-- FULL OUTER: NULLs on both sides (empty dept 9 → NULL name; orphan emp dept 99 → NULL dept_name)SELECT e.name, d.dept_nameFROM jdbc_join_emp eFULL OUTER JOIN jdbc_join_dept d ON e.dept_id = d.dept_id;Predicate pushdown
Each single-table WHERE predicate is pushed into its own ES sub-query before the JOIN, so a filtered JOIN returns strictly fewer rows:
SELECT e.name, e.salary, d.dept_nameFROM jdbc_join_emp eJOIN jdbc_join_dept d ON e.dept_id = d.dept_idWHERE e.salary > 5000 AND d.dept_name = 'Engineering';GROUP BY / HAVING / ORDER BY
Aggregation runs in DuckDB after the JOIN:
SELECT d.dept_name, COUNT(*) AS cntFROM jdbc_join_emp eJOIN jdbc_join_dept d ON e.dept_id = d.dept_idGROUP BY d.dept_nameHAVING COUNT(*) > 1ORDER BY COUNT(*) DESC;-- Engineering (3), Marketing (2) survive HAVINGSELECT aliases and ordinals work here — since arrow-extensions 0.2.5 (REPL bundle
0.20.4, JDBC / ADBC / Flight SQL driver0.2.5).ORDER BY cnt,HAVING cnt > 1,GROUP BYon an alias and ordinal forms such asORDER BY 2all resolve against the final SELECT list after the join, and alias matching is case-insensitive. Two limits remain: an alias is not legal inSELECT,ONorWHERE— nothing has been computed at that point — and it must be written bare —ORDER BY d.cntqualifies a name no table owns. Since arrow-extensions 0.3.3 the planner rejects the shapes it can recognise (a table-alias-qualified name that is a computed alias or a renamed column, such asd.cntord.dnameford.dept_name AS dname) with an error naming the bare spelling, and a leg that asks Elasticsearch for a column its index does not have fails with an error naming the column and the table alias — never DuckDB’s internalInvalid Input Error, which is what both shapes produced before 0.3.3. A qualified name that really is a column keeps working even when it spells an alias (o.id AS id … ORDER BY o.id,ROUND(o.amount) AS amount … ORDER BY o.amount), and a column used only in a leg’s own pushed-downWHERE(WHERE o.deleted_at IS NULLon an index no document has carried the field into yet) is still accepted even when the mapping does not have it. Before 0.2.5 all of these were rejected withAmbiguous column.The remaining ORDER BY gotcha: you cannot
ORDER BYa column that exists on both sides of the JOIN — order by a column unique to one side, e.g.d.dept_name, not the shared join keyd.dept_id.
Self-join — one index under two aliases
A table can be joined to itself: give the index two aliases and qualify every column. Each alias becomes its own ES sub-query (with its own pushed-down WHERE), exactly as two distinct indices would:
SELECT a.id, b.amountFROM orders aJOIN orders b ON a.id = b.idWHERE a.id <= 5 AND b.id <= 5;-- 5 rows: each order paired with itself
SELECT a.id AS left_id, b.id AS right_idFROM orders aJOIN orders b ON a.customer_id = b.customer_id;-- every pair of orders sharing a customer (fan-out, see below)Since engine
0.23.0with arrow-extensions0.3.3. Before that, the alias map kept only the last alias of a repeated index, so the first leg’s columns never resolved, theONclause lost its join key and DuckDB answeredParser Error: syntax error at end of input(softclient4es-arrow#144). Four rules follow, each enforced with a message that says why: qualify every column (a bareidis ambiguous between the two legs); alias a column you select from both legs (SELECT a.id AS left_id, b.id AS right_id—SELECT a.id, b.id,SELECT *anda.*over a self-join are rejected, because the result would carry the same column name twice and keep only one value — over two different indices that share a column name the REPL names the colliding columns<alias>.<column>(o.id,c.id) and leaves the others bare, while JDBC / ADBC / Flight SQL keep Arrow’s duplicate labels; andSELECT *over a JOIN omits object and nested columns and their sub-fields — list them explicitly to select sub-fields); give each leg its own alias — aliases compare case-insensitively, soFROM orders a JOIN orders a,FROM orders A JOIN orders aandFROM orders JOIN ordersare all rejected; and write a self-join withJOIN … ON— the comma formFROM orders a, orders bis a multi-index search with no join engine behind it and is rejected with a message pointing at theJOINspelling.
JOIN cardinality — fan-out on a non-unique key
A JOIN on a key that is not unique on the other side multiplies rows. That is standard SQL and the engine is doing it correctly, but it is the easiest way to get plausible-looking wrong numbers, because the row multiplication is invisible in the output.
With one tenant that has 2 EU error rows, 3 US rows and 2 AP rows:
SELECT eu.tenant_id, COUNT(*) AS eu_errors, AVG(us.latency_ms) AS us_avg_latency, AVG(ap.latency_ms) AS ap_avg_latencyFROM eu_events AS euJOIN us_events AS us ON eu.tenant_id = us.tenant_idJOIN ap_events AS ap ON eu.tenant_id = ap.tenant_idWHERE eu.level = 'ERROR'GROUP BY eu.tenant_id;eu_errors comes back as 12 — that is 2 × 3 × 2, one row per combination. The answer meant by the query is 2.
Why this is worse than one wrong column. Within a group whose rows all fan out by the same factor:
| Aggregate | Under fan-out |
|---|---|
COUNT, SUM | inflated by the fan-out factor |
AVG, MIN, MAX | unchanged — uniform duplication preserves them |
So in the example above both averages are exactly right, and only the count is wrong. Three columns out of four corroborate a result that is 6× off, and nothing in the output signals that a fan-out happened. Do not use “the averages look sensible” as a sanity check on a JOIN. (If the factor varies across rows within a group — because you grouped by something coarser than the join key — then AVG is silently weighted too, and it is wrong as well.)
What to do instead:
- Count a key from one side rather than rows of the joined product:
COUNT(DISTINCT eu.event_id). - Or aggregate before joining, so each side contributes one row per key — a materialized view per leg is the durable form of this.
- Sanity-check the row count against the left side alone before adding aggregates.
INNER JOIN also drops rows. A key absent from any joined table disappears from the result entirely — a tenant running in EU and US but not APAC vanishes from a query whose name says “every region”. Use LEFT JOIN when the left side is the population you actually mean.
In Federation specifically: each leg is staged and joined coordinator-local, so a fan-out inflates the staged intermediate and consumes the joined-output row cap (maxQueryResults, Community 10,000 / Pro 1,000,000). Hitting that cap is reported — see Row truncation at the result cap — but a truncated fan-out is still an answer to a question you did not ask.
ORDER BY … LIMIT (top-N)
SELECT e.name, e.salary, d.dept_nameFROM jdbc_join_emp eJOIN jdbc_join_dept d ON e.dept_id = d.dept_idORDER BY e.salary DESC LIMIT 3;-- 8000 / 6000 / 5500UNNEST + cross-index JOIN
Elasticsearch handles JOIN UNNEST on an ARRAY\<STRUCT\> natively inside the sub-query; the cross-index JOIN to the dimension runs in DuckDB. UNNEST does not count toward maxJoins:
SELECT o.id, oi.product, oi.quantity, d.dept_nameFROM jdbc_join_order_items oJOIN UNNEST(o.items) AS oiJOIN jdbc_join_dept d ON o.customer_id = d.dept_id;-- one row per nested item; the UNNESTed nested fields (product, quantity) populateINSERT … SELECT … JOIN
Write the joined output into a target index (Row-1 write; orphan dropped → 5 rows):
INSERT INTO jdbc_row1_insert_join_targetSELECT e.name, e.salary, d.dept_nameFROM jdbc_join_emp eJOIN jdbc_join_dept d ON e.dept_id = d.dept_id;CREATE TABLE … AS SELECT … JOIN (CTAS)
CREATE TABLE jdbc_row1_ctas_join_target ASSELECT e.name, e.salary, d.dept_nameFROM jdbc_join_emp eJOIN jdbc_join_dept d ON e.dept_id = d.dept_id;CREATE OR REPLACE TABLE … AS
Real REPLACE — drops then recreates, so re-running a narrower query shrinks the table:
CREATE OR REPLACE TABLE jdbc_row1_replace_ctas_target ASSELECT e.name, e.salary, d.dept_nameFROM jdbc_join_emp eJOIN jdbc_join_dept d ON e.dept_id = d.dept_id;INSERT … ON CONFLICT (col) DO UPDATE (upsert)
Idempotent upsert — a stable SHA-1 _id is derived from the conflict column, so re-running merges rather than duplicating (count stays 5, not 10):
INSERT INTO jdbc_row1_insert_join_upsert_targetSELECT e.emp_id, e.name, e.salary, d.dept_nameFROM jdbc_join_emp eJOIN jdbc_join_dept d ON e.dept_id = d.dept_idON CONFLICT (emp_id) DO UPDATE;Two ON CONFLICT forms are rejected: CTAS + ON CONFLICT and INSERT + ON CONFLICT DO NOTHING are rejected at the boundary — Elasticsearch has no native skip-on-conflict with stable ids. Use
INSERT … ON CONFLICT (col) DO UPDATEfor upserts.
Prepared statement through a JOIN
A bound parameter is substituted into the sub-query before planning, so the same PreparedStatement returns different rows for different bindings:
SELECT e.name, d.dept_nameFROM jdbc_join_emp eJOIN jdbc_join_dept d ON e.dept_id = d.dept_idWHERE e.salary > ?; -- bind 3500.0 → more rows; bind 6500.0 → fewer (Carol = 8000 survives)Row 2 — cross-cluster conveyor
A Row-2 operation has its target in a different cluster from its source. The Federation coordinator runs the source SELECT on the source cluster’s sidecar, receives the result as an Arrow stream, and conveyors it to the target sidecar for a bulk-load. (The JOIN in 2a/2b is itself single-source — both source tables are in prod_us — but the target prod_eu is a different cluster, and that is what makes it Row 2. A plain catalog-to-catalog INSERT … SELECT * with no JOIN is also a Row-2 conveyor.)
Row 2 — cross-cluster conveyor: the coordinator runs the source SELECT on prod_us, receives the result as an Arrow stream, and bulk-loads it onto the prod_eu target sidecar.
Cross-cluster references use backtick-quoted catalog prefixes — the catalog name is the Federation servers.<name> alias, which Federation strips before forwarding each leg’s SELECT to its source cluster.
Examples are transcribed from the SoftClient4ES Federation integration test suite.
Cross-cluster INSERT-with-JOIN — the source SELECT runs on prod_us, conveyed to prod_eu:
INSERT INTO `prod_eu`.destSELECT o.id, c.nameFROM `prod_us`.orders oJOIN `prod_us`.customers c ON o.id = c.id;Cross-cluster CTAS-with-JOIN — the target DDL is inferred from the source FlightInfo schema, then conveyed:
CREATE TABLE `prod_eu`.dest ASSELECT o.id, c.nameFROM `prod_us`.orders oJOIN `prod_us`.customers c ON o.id = c.id;Row 3 — multi-source coordinator
A Row-3 JOIN reaches two or more source clusters. The coordinator stages each leg to Parquet scratch on disk and exposes it as a per-query DuckDB view (named q_<UUID> for isolation), then runs the JOIN coordinator-local.
Multi-Elasticsearch only in this release: in this release, Row 3 joins across multiple Elasticsearch clusters. Joining Elasticsearch against heterogeneous sources (Postgres, MySQL, Snowflake, …) is coming in the upcoming release (Quarter 1 2027) — it is not promised for this release.
Row 3 — multi-source coordinator: each source leg (prod_us, prod_fr) is staged to Parquet scratch and exposed as a per-query DuckDB view; the JOIN runs coordinator-local before landing on the prod_eu target.
Multi-source SELECT JOIN (read; the headline three-cluster query — two source catalogs, joined coordinator-local):
SELECT o.id, c.nameFROM `prod_us`.orders oJOIN `prod_eu`.customers c ON o.id = c.id;-- two source catalogs (prod_us, prod_eu) → coordinator-local joinThis is the SRE wedge: correlate logs, metrics, and traces across regional ES clusters in one SQL query. (The full SRE story lives in the Federation operator guide and the three-region example topology.)
Multi-source INSERT-with-JOIN (sources prod_us + prod_fr, target prod_eu):
INSERT INTO `prod_eu`.destSELECT o.id, c.nameFROM `prod_us`.orders oJOIN `prod_fr`.customers c ON o.customer_id = c.id;Multi-source CREATE OR REPLACE TABLE … AS:
CREATE OR REPLACE TABLE `prod_eu`.dest ASSELECT o.id, c.nameFROM `prod_us`.orders oJOIN `prod_fr`.customers c ON o.customer_id = c.id;Performance characteristics
What to expect per row, qualitatively. (Hard latency and throughput numbers belong to a follow-up release’s Arrow benchmark — link forward, don’t pre-empt.)
- Row 1 — lowest latency of the three rows; the JOIN is in-process DuckDB over ES sub-query results, so the dominant cost is the ES sub-queries plus the cross product. The first JOIN on a fresh process pays a one-time DuckDB native-library warm-up.
- Row 2 — adds a network hop (coordinator → source SELECT) plus a bulk-load conveyor to the target; throughput is bound by the Arrow stream and target ingest, not by JSON — data crosses the wire as Arrow RecordBatches.
- Row 3 — adds per-leg Parquet staging on the coordinator before the DuckDB join; cost scales with the largest leg plus the join cardinality.
- Arrow wire format — for Flight SQL and ADBC clients specifically, results stream as Arrow RecordBatches (zero JSON on the wire). The JDBC driver surfaces JOIN rows via an in-process Arrow result set and the REPL renders to a console, so the clean zero-JSON-wire property is a Flight-SQL/ADBC trait, not a blanket one. The value here is interoperability and SQL completeness; the full zero-copy speed story is a follow-up release’s benchmark.
Row truncation at the result cap
The joined output is capped at the tier’s maxQueryResults (Community 10,000). Over-cap results are truncated with a warning — never silently dropped: JDBC and ADBC raise SQLWarning 01004 (Result truncated to N rows), and Flight SQL emits an x-result-truncated header.
Join inputs are never capped: the result cap applies to the joined output only, never to a JOIN leg. Capping a leg would truncate a JOIN input and produce silently wrong results. So a wide Community join can hit the 10,000-row output cap even though none of its inputs were truncated.
Licensing & the meters
Every tier has every feature. The meter gates scale, not an on/off switch.
maxJoins— JOIN depth per query. A 2-table JOIN counts as 1 JOIN, a 3-table JOIN as 2, a 4-table JOIN as 3. Community 2 (up to a 3-table JOIN), Pro 5 (up to a 6-table JOIN), Enterprise ∞.UNNESTdoes not count towardmaxJoins.maxClusters— cross-cluster reach. Community 1, Pro 5, Enterprise ∞. Single-cluster Federation is free — the meter, not a feature flag, is the paywall. ES-only in this release (phrased “across N ES clusters”).maxMaterializedViews— persisted views. Community 1, Pro 50, Enterprise ∞.
Community has Federation. It is capped at 1 cluster — the quota is the paywall, not a feature gate. One sidecar / one cluster boots free; a second cluster makes the Federation sidecar fail to start by design.
The same meters, in source-of-truth order (MV / results / clusters / joins):
| Tier | maxMaterializedViews | maxQueryResults | maxClusters | maxJoins |
|---|---|---|---|---|
| Community | 1 | 10,000 | 1 | 2 |
| Pro | 50 | 1,000,000 | 5 | 5 |
| Enterprise | ∞ | ∞ | ∞ | ∞ |
What happens at each cap (verified)
maxJoinsexceeded → the planner rejects the query. A 4-table query (3 cross-index JOINs) under Community is rejected with a message that names the count, the limit, the next tier, and the upgrade URL:
-- 4 tables = 3 cross-index JOINs → exceeds Community maxJoins=2 → rejectedSELECT t.id, d.dept_name, r.region_name, m.team_nameFROM jdbc_test tJOIN jdbc_join_dept d ON t.id = d.dept_idJOIN jdbc_join_region r ON t.id = r.region_idJOIN jdbc_join_team m ON t.id = m.team_id;-- Error: "Query contains 3 cross-index JOINs ... maximum of 2 ... Upgrade to Pro ...-- See: https://portal.softclient4es.com/pricing"The same query runs unchanged on Pro or Enterprise (higher / no cap).
maxClustersexceeded → the Federation sidecar fails to start (by design — a CrashLoop) rather than silently dropping a cluster.maxQueryResultsexceeded → truncate-with-warning on the no-LIMITpath (SQLWarning 01004/ Flightx-result-truncated), or an HTTP 402 on an explicitLIMITover quota.
The JOIN engine ships in the Elastic-License extensions — free to use (not “source-available”).
For the full price matrix and editions, see Licensing & pricing.
What does NOT work yet
- Arbitrary subqueries and CTEs inside a JOIN query — coming in the next release (Quarter 4 2026).
- Heterogeneous Row-3 sources (joining Elasticsearch with Postgres, MySQL, Snowflake, …) — coming in the upcoming release (Quarter 1 2027); this release’s Row 3 is multi-Elasticsearch only.
- JOIN inside a watcher input (
CREATE WATCHER … FROM a JOIN b ON …) — a watcher input is a single Elasticsearchsearchrequest over a list of indices, so the join is rejected at parse time. Pre-join the sources with a materialized view and have the watcher search the view.
Full list: Known limitations.
Try it in 5 minutes
- JDBC quickstart — single-cluster JOIN from any JDBC tool.
- ADBC quickstart — in-process Arrow.
- Arrow Flight SQL quickstart — columnar streaming.
- Federation operator guide — cross-cluster (Rows 2 & 3) with the Helm chart.
- Example topologies live in the
softclient4es-helmrepo:softclient4es-federation/examples/single-cluster/and.../three-region/.
See also
- Materialized Views — superpower #2: persisted, pre-joined data.
- DQL — Queries — non-JOIN SELECT and
JOIN UNNESTdetail. - Licensing & pricing — the full edition / quota matrix.
- Known limitations.
- Federation operator guide.