Live event contract¶
Runtime¶
stream-events.ipynb is rendered from
utility/notebooks/templates/driver-05-stream.py. It is an optional,
long-running Fabric Spark Structured Streaming driver.
Supported sinks:
eventhouse: direct typed writes through the Spark Kusto connector;delta: development/validation landing table.
The current path does not require Kafka, Event Hubs, or a Fabric Eventstream.
Envelope¶
Every event contains:
event_typetrace_idingest_timestampschema_versionsourcecorrelation_idpartition_keysession_idparent_event_id
Payload fields are event-specific and mapped by EVENT_PAYLOADS.
Business event types¶
receipt_createdreceipt_line_addedpayment_processedinventory_updatedstockout_detectedreorder_triggeredcustomer_enteredcustomer_zone_changedble_ping_detectedtruck_arrivedtruck_departedstore_openedstore_closedad_impressionpromotion_appliedonline_order_createdonline_order_pickedonline_order_shipped
KQL defines an additional unknown_event catch-all table. It is not a
nineteenth generated business event. The current notebook rejects unmapped
event types before committing the micro-batch checkpoint. The catch-all is an
operational queued-ingestion boundary, not an active direct-writer DLQ.
Here DLQ means dead-letter queue: a place to preserve messages that could not
be processed.
Eventhouse writes¶
For each micro-batch, the notebook:
- persists the batch;
- finds present mapped event types;
- resolves one notebook runtime token;
- maps JSON to typed envelope/payload columns;
- writes event types concurrently to their same-named KQL tables;
- uses
FailIfNotExist; - sets
flushImmediately=true; - unpersists the batch.
If kusto_uri is blank, the notebook resolves the KQL database
queryServiceUri by display name in the current workspace.
Trigger and checkpoint behavior¶
- Eventhouse uses a 10-second processing trigger for an unbounded run.
- Bounded runs use a 2-second processing trigger.
- Checkpoints are sink-specific under
Files/setup/stream/checkpoint/<sink>. - A logical stream ID is stored under the checkpoint root. Notebook restarts reuse it; deleting the checkpoint root creates a new event-ID namespace.
Duplicate and failure behavior¶
- Generated business IDs and
trace_idinclude the persisted stream ID. - Eventhouse writes use the connector's transactional mode, deterministic request IDs, duplicate-blob protection, and matching ingest-by/check tags for each stream, event type, and Spark micro-batch.
- Unmapped event types and any required table-write failure raise from
foreachBatch, so Spark does not commit that batch's checkpoint. - Silver transforms remove exact duplicate candidates, reject conflicting rows
with the same identity, and use Delta
MERGEbefore advancing watermarks. - Gold tables are rebuilt with overwrite semantics from duplicate-safe Silver inputs, so reruns do not accumulate aggregate rows.
These controls make expected retries idempotent and are covered by injected failure/replay contracts. The readiness adapter checks live ingestion and checkpoint tags. The remaining IMP-013 evidence is a recent bounded run of the intentionally manual stream.
Cross-layer ownership¶
The exact KQL table shape is in 01-create-tables.kql. JSON ingestion mappings
in 02-create-ingestion-mappings.kql support queued raw-JSON ingestion but are
not used by the direct typed live path.
Silver mappings are separately implemented in
fabric/lakehouse/03-streaming-to-silver.ipynb. Their current divergence from
the historical contract is documented in
Fabric analytics.
contracts/retail-demo.json declares one stable path ID per emitted event and
one derived attribution path. It stores keys, UTC event-time semantics,
targets, terminals, and named exceptions, but never copies physical field or
type inventories. python scripts/check_data_contracts.py derives those
inventories from the driver, KQL, schemas.py, notebooks, and active TMDL.
Truck lifecycle¶
truck_arrived and truck_departed share truck_id, dc_id, store_id, and
shipment_id. The departure payload and envelope timestamp are later than the
arrival timestamp. Normal dwell is a deterministic 30-75 minutes; every fifth
eligible logistics value takes the deterministic 120-minute late path, above the
90-minute rule threshold.
Silver joins the two event tables on all four lifecycle keys and writes one
completed fact_truck_moves row. Eventhouse and Gold expose dwell in minutes
with STORE_<id> / DC_<id> site labels.
tests/test_truck_dwell_contract.py verifies the generator, rendered stream
notebook, KQL table/function contract, Silver join, Gold aggregation, queryset,
dashboard template, and 90-minute rule. A live Fabric smoke run remains the
deployment verification gate.
Marketing attribution¶
Attributed store and online scenarios emit two ad_impression touches and one
purchase journey without introducing another event type. The existing envelope
correlation_id carries attribution_journey_id through impressions,
receipt/order creation, promotion, payment, and online status events.
The contract is deterministic last-touch within an inclusive seven-day window,
ordered by touch timestamp and impression ID. Silver writes one
fact_marketing_attribution row after purchase, promotion, and approved
payment reconcile; attributed rows additionally require a valid selected
touch, while purchases without a journey are marked
UNATTRIBUTED_NO_JOURNEY.
Money remains integer cents:
gross_subtotal_cents - discount_cents = net_subtotal_cents
net_subtotal_cents + tax_cents = total_cents
approved payment_cents = total_cents
attributed_revenue_cents = net_subtotal_cents
Eventhouse exposes the same logic through fn_marketing_attribution() and
fn_campaign_performance(). Gold publishes campaign_performance_daily, and
the Direct Lake model exposes both the audit fact and campaign KPIs. Contract
tests cover batch determinism, live payload/KQL mappings, Silver/Gold wiring,
semantic measures, and all cent equations. A live Fabric smoke run remains the
deployment verification gate.
Contract evidence and live boundary¶
Eight fixture scenarios cover all 18 emitted types plus nullable and unknown-event variants. The parameterized matrix validates the nine-field envelope, KQL DDL/mappings, UTC values, business/dedupe keys, Silver/Gold routes, and either an active semantic terminal or the named streaming-only exception.
Repository acceptance for IMP-005 is complete. A live Fabric staging run
through Eventhouse, the optional Silver/Gold projection, and Direct Lake
remains the external evidence boundary; the static check does not claim that
tenant-level execution.