Stage 3 of 6

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

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

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

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 PackageConsumer
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
Real World Example

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.