Transformation & Processing
Raw events are deduplicated, stitched to an identity, sessionized, and modeled — in real time for immediate use, and in batch for deeper aggregation.
⚡ Stream Processing
Real-time — Bytewax / Quix Streams
Event Deduplication
Removes duplicates so one action is never counted twice.
Sessionization
Groups raw events into the logical sessions analytics runs on.
Identity Resolution (Stitching)
Links events across devices back to the same real customer.
Data Quality Checks
Continuous checks that catch regressions before they reach a dashboard.
Real-time Enrichment
Adds computed, business-meaningful context to events as they stream.
Derived Event Generation
Synthesizes higher-level events like cart_abandoned from raw patterns.
How to Consume
Cxos.Processing.Api · Query (job status) / Command (backfill trigger)
GET /v1/processing/stream-worker/status
{ "consumer_group": "cxos-processing-stream-worker", "lag_seconds": 2.1, "healthy": true }
POST /v1/processing/backfill { "job": "identity-resolution", "from": "2026-07-01" }
📋 Batch Processing
dbt Core — SQL & Python Models
Rollups & Aggregations
Pre-computed summaries that make dashboards fast at any scale.
Data Modeling
Transforms raw events into the dimensional models the business thinks in.
Backfills
Reprocesses history when a source, bug fix, or rule changes.
Machine Learning Feature Prep
Builds the feature tables ML models train and score against.
How to Consume
Cxos.Processing.Api · Query
// Not called directly for output — read the marts tables it produces
// via Cxos.Intelligence.Api's Query & Analytics Engine. To check a job:
GET /v1/processing/batch-worker/runs?job=rollups_aggregations&status=latest
{ "job": "rollups_aggregations", "status": "success", "rows_written": 2847193 }
📜 Event Schema & Registry
Schema Management
The authoritative definition of every event type's shape.
Versioning
Lets a schema evolve without breaking existing consumers.
Compatibility
Automated checks that a schema change won't silently break anything.
Governance Rules
Policy layer governing naming, PII tags, and ownership.
How to Consume
Cxos.Processing.SchemaGovernance.Api · Query
GET /v1/schema-registry/events/product_viewed?version=latest
{ "event": "product_viewed", "version": 3, "owner_team": "commerce-platform",
"compatibility": "backward", "pii_fields": [] }
🔌 Platform Connectors
The engine adapters and inbound connector family this stage reads from and writes through.
| Connector Package | Consumer |
|---|---|
Cxos.Connectors.* (inbound family) |
Stream Worker · Batch Worker (every normalized event lands here first) |
| dbt-snowflake adapter | Rollups & Aggregations · Batch Processing |
| dbt-spark adapter | Batch Processing (large-scale rollups, backfills) |
Cxos.Ingestion.Client (schema source) |
Stream Worker validation against the shared contract |
Events are deduplicated, stitched to a single customer profile, sessions are built, and purchase events update lifetime value in real-time. Nightly jobs aggregate data for faster analytics.