Enrichment
Broadcast hash join against small reference tables (internal/enrich).
The reference is read whole into worker RAM; every event matches in O(1)
against the map; a periodic full re-read swaps the map atomically (an
in-flight event finishes against the old image). No lookup per event, no
shuffle, no windowed state.
Rules
- Projection is explicit.
selectis required and lists exactly the reference columns the event receives — nothing unselected ever lands in the sink.select: ["*"]is the wildcard sugar for "take everything"; it must be the only entry. - Column namespacing follows Spark DataFrame semantics. An unrenamed
column is auto-prefixed
{table}.{column}, so two references injecting a same-named column never silently collide.asrenames at load time and overrides the prefix. joinTypeis required and accepts five values:left(aliasleft outer),inner,left semi,left anti.left/left outerpass a miss with the reference columns NULL;innerdrops a miss;left semikeeps a hit without the reference columns in the output;left antikeeps a miss without the reference columns.selectis rejected as a spec error onleft semi/left anti— neither emits reference columns. A delete always survives every join type, bypassing the lookup entirely.onColdStartis decided per BATCH, not per event. Before the reference's first load completes:dropdrops the whole batch;bufferandpassare currently identical — every row in the batch passes as a miss (reference columns NULL, following the join type). There is no per-event buffer or drain queue;bufferLimitsis still validated as spec grammar but tunes nothing today.- A wildcard reference's real columns are known before boot completes, not just after the first refresh — see Enrich internals.
- Point-in-time, and it says so. Enriched columns are not reproducible by replay (the reference is a snapshot, not CDC); source columns and position stay deterministic either way.
Full grammar detail, the boot-time join-key type check, duplicate-key rejection, and cold-start internals: Enrich internals.
Examples
Every example below is a complete, valid pipeline spec — source/sink
are required on the reference too (spec.EnrichSource), and event columns
come from the source table's own schema (introspected for SQL sources),
not a columns: list on tables[].
Single reference:
pipeline: orders-enriched
source:
kind: mysql
uri: mysql://repl:replpass@127.0.0.1:3306/shop
sink:
type: iceberg+rest
uri: http://localhost:8181/api/catalog
warehouse: quickstart_catalog
namespace: bronze
clientId: root
clientSecret: s3cr3t
scope: PRINCIPAL_ROLE:ALL
tables:
- source: shop.orders
target: bronze.orders
primaryKey: [id]
writeMode: upsert
createIfNotExists: true
enrich:
- table: customers
source:
uri: mysql://repl:replpass@127.0.0.1:3306/shop
query: SELECT id, name, tier FROM customers
on: {customer_id: id}
select: [name, tier]
joinType: left
An order row {id: 1, customer_id: 42, amount: 100} lands as
{id: 1, customer_id: 42, amount: 100, customers.name: "Ana", customers.tier: "gold"}.
A customer_id with no matching row lands with customers.name and
customers.tier both NULL (left join, not dropped).
Multiple references (no collision — each gets its own prefix):
enrich:
- table: customers
source: {uri: "mysql://repl:replpass@127.0.0.1:3306/shop", query: "SELECT id, name FROM customers"}
on: {customer_id: id}
select: [name]
joinType: left
- table: products
source: {uri: "mysql://repl:replpass@127.0.0.1:3306/shop", query: "SELECT id, name, category FROM products"}
on: {product_id: id}
select: [name, category]
joinType: left
→ {..., customers.name: "Ana", products.name: "Laptop", products.category: "electronics"}
Renaming with as:
enrich:
- table: customers
source: {uri: "mysql://repl:replpass@127.0.0.1:3306/shop", query: "SELECT id, name, tier FROM customers"}
on: {customer_id: id}
select: [name, tier]
joinType: left
as: {"customers.name": "client_name"}
→ {..., client_name: "Ana", customers.tier: "gold"} — name renamed,
tier keeps its prefix. as keys must be the table-prefixed column
name (customers.name, not name) and must name a column in select.
spec.Validate() and enrich.New() agree on this rule (fixed by #78);
the example above boots and runs as shown.
joinType is required on every reference — there is no default miss
policy; see Rules above for the five accepted values.