Ingestion Layer → Streaming Ingestion

Redpanda (Kafka Compatible)

The Kafka-API-compatible streaming backbone CXOS uses in place of Apache Kafka itself, chosen for lower operational overhead at equivalent throughput.

High-Level Design

Redpanda is the durable log every ingested event flows through.

Data Source
Ingestion API
Producer — publishes validated events
→
Ingestion
Redpanda Cluster
Kafka-compatible durable log
→
Processing
.NET Core Stream Worker
Consumer — Confluent.Kafka .NET client
→
Foundation
Data Lakehouse
Where consumed events are ultimately written
→
Intelligence
Schema Registry
Enforces the contract at the broker level
→
Activation
Every Real-time Consumer
The backbone every real-time stage depends on

💼 Business Context

  • The streaming backbone is invisible to business stakeholders when it works, and existential when it doesn't — it's the platform's real-time nervous system
  • Kafka-API compatibility means the ecosystem of Kafka connectors, client libraries, and tooling all work unmodified
  • Owned by Platform / Infrastructure Engineering

🔌 Technical Overview

Where CXOS is deployed with Redpanda directly (rather than Azure Event Hubs' native Kafka-compatible endpoint), it runs as a managed cluster fronting the same producer/consumer contracts used elsewhere — the .NET Core Ingestion API and Stream Worker use the standard Confluent.Kafka .NET client against Redpanda exactly as they would against Event Hubs, so the rest of the platform's code is broker-agnostic. Topic partitioning follows the event's user_id/anonymous_id to preserve per-customer ordering.

Compatibility

Kafka producer/consumer API Schema Registry compatible Confluent.Kafka .NET client

💾 Topic Configuration

{
  "topic": "cxos.events.raw",
  "partitions": 24,
  "partition_key": "user_id | anonymous_id",
  "retention_hours": 168,
  "compaction": false
}

🔗 Integration Points

  • Redpanda cluster (or Azure Event Hubs' Kafka-compatible endpoint, interchangeably)
  • Confluent.Kafka .NET client — used by the Ingestion API (producer) and Stream Worker (consumer)
  • Schema Registry — enforces the Cxos.Ingestion.Client contract at the broker level
  • Azure Monitor — broker-level throughput, consumer lag, and partition-skew dashboards

🧰 Services Consumed

  • Owning microservice — Cxos.Ingestion.Infrastructure (see the Full Application Service Map)
  • No dedicated database — Azure Event Hubs is a message stream, not a store

⚠️ Non-Functional Considerations

  • Scale: partition count is sized for peak expected throughput with headroom; partitioning by customer ID keeps load reasonably balanced
  • Latency: sub-second producer-to-consumer latency under normal load
  • Reliability: replication factor 3 across availability zones; offset commits are transactional with downstream writes where exactly-once semantics matter
  • Security/Privacy: mTLS between all producers/consumers and the broker; topic-level ACLs restrict which services can produce or consume which event types

🎯 Enterprise Example

The platform migrates its highest-volume topic from a self-managed Kafka cluster to Redpanda without changing a single line of the Stream Worker's consumer code, since both speak the same Kafka wire protocol — cutting operational overhead while preserving every existing integration.

← Back to Streaming Ingestion