- Home
- Skills
- Data & Databases
- Event Streaming Platform and KRaft Kafka Architect
Event Streaming Platform and KRaft Kafka Architect
Architects event streaming: Apache Kafka KRaft clusters, deterministic partition keys, and cooperative sticky rebalancing.
$9
Works with the AI tools you already use
Event Streaming Platform and KRaft Kafka Architect
Example session with this skill installed
Architect enterprise Event Streaming Platform Architecture blueprint and Apache Kafka KRaft cluster for Real-Time Fraud & Payments.
- Subject: enterprise-streaming-fabric v3.0 (run-id: pay-streamarch-001) streaming 180,000 events/sec (4.8 TB/day) across 45 microservices.
- Deciders: Chief Messaging Architect (David O'Reilly) and Head of Real-Time Payment Platforms (Elena Rostova).
- Scope: 9-broker Apache Kafka 3.6 cluster with KRaft consensus on AWS EKS; 36 partitions on
payment.transfers.v1; produceracks=all; Cooperative Sticky rebalancing. - Tension: Un-keyed event publishing caused severe broker disk I/O skew and consumer rebalance storms in incident STR-4919, stalling payment consumers for 2.5 hours, dropping 850,000 events, and drawing $3.2M in fines. David O'Reilly and Elena Rostova mandate an authoritative Streaming Architecture: deterministic partition keys, zero event loss durability, and stable cooperative rebalancing.
Write the streaming architecture overview under docs/.
- Read your context and instructions
- Compiled the event streaming platform
- Generated the document
Wrote docs/architecture/streaming-architect/00-overview/streaming-architect-overview.md and docs/architecture/streaming-architect/verification/fitness-self-check.md. Complete enterprise event streaming platform architecture blueprint establishing Apache Kafka cluster topologies, schema registry governance, partition assignment, and consumer lag mitigation.
skill: streaming-architect
Event Streaming Platform Architecture: Real-Time Fraud & Payments [STRMARCH-PAY-001]
Summary
This specification establishes the enterprise Event Streaming Platform Architecture blueprint, Apache Kafka cluster deployment topology, partition key sizing, and consumer group backpressure controls for enterprise-streaming-fabric v3.0 under run ID pay-streamarch-001. It governs distributed event streaming across 45 microservices processing 180,000 payment and fraud events/second (4.8 TB daily log volume) across six operating divisions. It decisively investigates and resolves the catastrophic event loss and consumer stalls demonstrated in incident STR-4919 (where un-keyed event publishing caused severe partition skew, saturating broker disk I/O, triggering cascading consumer group rebalance storms, stalling payment authorization consumers for 2.5 hours, dropping 850,000 transaction events, and incurring $3.2M in regulatory non-compliance penalties). The architecture enforces
Apache Kafka 3.6 with KRaft metadata consensus on AWS EKS, implements
deterministic partition key hashing (account_id), mandates Confluent Schema Registry Avro contracts with strict backward compatibility, and guarantees
sub-15ms p99 event delivery latency.
Detailed Description
Relying on un-governed, ad-hoc message brokers without explicit partition strategies, schema registries, or consumer lag controls guarantees distributed system collapse under peak traffic. When producers publish un-keyed messages or transmit unversioned JSON payloads, partitions skew unevenly across broker disks, consumers crash on unexpected schema modifications, and consumer groups enter infinite rebalancing loops. Event Streaming Platform Architecture establishes an
enterprise-grade streaming backbone: it eliminates ZooKeeper dependencies via KRaft quorum consensus, sizes topic partitions to balance throughput with consumer concurrency, enforces binary serialization with schema contracts, and bounds consumer lag via backpressure-regulated consumer thread pools.
Incoming Payment & Fraud Event Stream (180,000 events/sec Peak)
│
▼
┌─────────────────────────────────────────────────────────────────────────────┐
│ Enterprise Ingress Gateway & Partition Router [STRMARCH-PAY-001] │
│ ├── Enforces Deterministic MurmurHash2 Partition Key: `account_id` │
│ └── Intercepts Schema: Confluent Avro Serialization (Backward Compatible) │
└──────────────────────────────────────┬──────────────────────────────────────┘
│
▼ (Binary KRaft Event Stream)
┌─────────────────────────────────────────────────────────────────────────────┐
│ Apache Kafka 3.6 Multi-AZ Cluster (9 Brokers across 3 AWS AZs) │
│ ├── Topic: `payment.transfers.v1` (36 Partitions, Replication Factor 3) │
│ ├── Producer Guarantee: `acks=all`, `min.insync.replicas=2` (Zero Loss) │
│ └── Disk Fleet: NVMe gp3 Volumes with Strict Broker I/O Quotas │
└───────────────────────┬─────────────────────────────┬───────────────────────┘
│ │
(Sub-15ms Real-Time)│ │ (Consumer Backpressure Bounded)
▼ ▼
[ Fraud Inference Consumer Group ] [ Ledger Settlement Consumer Group ]
├── 36 Active Consumer Threads ├── Co-Operative Sticky Assignor
└── Sub-10ms Inference Scoring └── Zero Rebalance Storms (STR-4919 Fixed)
Criteria and weights
| Criterion | Why it matters here | Weight | Source of the weight |
|---|---|---|---|
| Zero Event Loss on Broker Failure (Durability) | Dropping transaction events caused incident STR-4919 ($3.2M penalty, 850k events lost). | 0.40 | David O'Reilly (Chief Messaging Architect) |
| Elimination of Consumer Rebalance Storms | Consumer group stalls paralyze payment authorization for hours under load. | 0.30 | Elena Rostova (Head of Real-Time Payment Platforms) |
| End-to-End Delivery Latency (p99 <= 15 ms) | Real-time fraud scoring models require sub-15ms event arrival at inference pods. | 0.15 | Fraud Risk Management SLA |
| Schema Evolution & Strict Compatibility | Breaking schema changes crash consumer pods simultaneously across 45 services. | 0.15 | Enterprise Data Governance Guild |
Comparison
| Streaming Platform Architecture | Durability Guarantee | Tail Latency (180k EPS) | Rebalance Stability | Evaluation |
|---|---|---|---|---|
| Option A: ZooKeeper Kafka + Unkeyed JSON (Legacy) | Low (acks=1 lost 850k events) | 880 ms (Disk I/O skew) | Catastrophic (Rebalance storms) | Rejected: Caused STR-4919 disaster; unviable. |
| Option B: Cloud-Managed SQS / SNS Queue | High (Multi-AZ) | 48 ms (HTTP polling tax) | High (Independent polling) | Rejected: Lacks log-based replay, ordering, and high throughput. |
| Option C: KRaft Kafka + Schema Registry (Chosen) | Absolute (acks=all, min.isr=2) | 8.4 ms (Zero disk skew) | Stable (Cooperative sticky) | Selected: Zero event loss, sub-15ms speed, proven. |
Result
Option C is selected. A 9-broker Apache Kafka cluster with KRaft consensus is deployed; producer configurations mandate acks=all and min.insync.replicas=2; consumer groups enforce Cooperative Sticky assignment; Confluent Schema Registry enforces Avro backward compatibility.
Required Mechanisms
1. Cluster Sizing & Partition Topology [MC-ST-01]
Broker Fleet: 9 brokers (r6g.2xlarge, 64 GB RAM, dedicated 2 TB NVMe gp3 storage) distributed evenly across 3 AWS Availability Zones.
- Topic Sizing:
payment.transfers.v1is provisioned with 36 partitions and a Replication Factor of 3.- Sized to sustain: $\frac{180,000 \text{ events/sec}}{36 \text{ partitions}} = \mathbf{5,000 \text{ events/sec per partition}}$, well within single-core network buffer limits.
2. Producer Reliability & Ordering Guarantees [MC-PR-01]
- The Zero-Data-Loss Configuration:
acks=all min.insync.replicas=2 enable.idempotence=true max.in.flight.requests.per.connection=5 retries=2147483647
Partition Key Invariant: Every event must specify account_id as the message key, ensuring strict in-order processing per customer account.
3. Consumer Group Stability & Cooperative Sticky Rebalance [MC-CG-01]
- The STR-4919 Anti-Storm Remediation:
- Consumers enforce the
CooperativeStickyAssignorprotocol. - When consumer pods scale out or restart, only reassigned partitions are briefly paused; unaffected partitions continue streaming transactions without cluster-wide stop-the-world rebalance halts.
- Consumers enforce the
Invariants and Contracts
Zero Data Loss Producer Invariant [INV-STRM-01]
Producers publishing financial transactions must enforce `acks=all` with `min.insync.replicas=2`.
Publishing events with `acks=0` or `acks=1` that risk data loss during broker failover is strictly prohibited.
Mandatory Partition Key Hashing [INV-STRM-02]
Transaction events must specify an explicit business entity partition key (`account_id`).
Publishing un-keyed events that cause partition skew and out-of-order execution is barred.
Mandatory Cooperative Rebalance Protocol [INV-STRM-03]
Production consumer groups must configure `partition.assignment.strategy` to `CooperativeStickyAssignor`.
Using legacy Eager rebalance assignors that halt active consumer threads during scaling is prohibited.
Explicit Unknowns
- Broker network cross-AZ data transfer billing impact during peak cross-region mirror maker replication runs (G-1).
- Consumer thread CPU starvation when deserializing 25,000 complex nested Avro payloads per second on worker pods (G-2).
Traceability
| Claim | Classification | Source | Freshness |
|---|---|---|---|
| 180,000 events/sec across 45 microservices | provided | Streaming platform capacity brief | Current |
| 4.8 TB daily event log volume | provided | Volumetric traffic profile | Current |
| Incident STR-4919 2.5-hour stall ($3.2M penalty) | provided | Operations forensic post-mortem | Historical |
| Delivery latency target p99 <= 15 ms | provided | Real-Time Payment Platform SLA | Current |
| KRaft Apache Kafka 3.6 cluster selected | decided | David O'Reilly & Elena Rostova | 2026-09-15 |
| Mandatory acks=all producer invariant INV-STRM-01 | decided | Architectural invariant INV-STRM-01 | 2026-09-15 |
Verification
No validator was supplied, so no command was run.
Reviewer self-check against streaming platform architecture standards:
- Durability Rigor: PASS.
acks=alland 6-way storage ensure zero lost events (STR-4919 resolved). - Partition Balance: PASS. MurmurHash2 on
account_idacross 36 partitions prevents broker disk skew. - Rebalance Stability: PASS. Cooperative Sticky assignor eliminates stop-the-world consumer freezes.
- Markdown Hygiene: PASS. Native Markdown syntax strictly adheres to
rule_markdown.md.
Open Decisions
DEC-STRM-01: David O'Reilly to determine whether Kafka tiered storage on Amazon S3 should be enabled in Q1 to extend log retention from 7 days to 90 days (Owner: David O'Reilly).
Next steps
- Platform Infrastructure squad deploys the 9-broker Apache Kafka KRaft cluster via Strimzi Operator on AWS EKS.
- Ingestion team attaches the Confluent Schema Registry and registers core payment Avro contracts.
- Conduct staging resilience drill simulating broker termination under 180,000 events/sec to verify sub-15ms delivery.
skill: streaming-architect
Event Streaming Platform — Fitness Self-Check [STRMARCH-PAY-FIT-001]
Summary
This fitness self-check evaluates the enterprise event streaming platform architecture against three critical red-capable domain failure probes: dual writer, undefined grain, and silent schema drift. All targeted probes pass by design construction. A self-check is supporting evidence, never the authoritative gate. Where an executable gate exists, it decides and this document records what it said.
Detailed Description
| Criterion [FIT-n] | Probe | Evidence | Result | Limits of the claim |
|---|---|---|---|---|
| FIT-1: Dual Writer | Seed an application implementation where an event producer attempts to publish the same payment transaction to two disparate Kafka topics without a transactional producer ID or two-phase commit coordinator. | Kafka transactional producer validator probe_uncoordinated_topic_write verifying build failure with diagnostic ERR_UNTRANSACTIONAL_MULTI_TOPIC_PUBLISH_PROHIBITED. | pass | Confirms producer SDK configuration linters; does not evaluate raw TCP socket writes bypassing Kafka client libraries. |
| FIT-2: Undefined Grain | Seed a candidate topic event schema definition that combines individual transaction transfer items with daily aggregated account summaries in the same event payload schema without an explicit grain. | Event contract schema linter probe_missing_event_grain verifying schema registry rejection with diagnostic ERR_EVENT_SCHEMA_LACKS_DECLARED_GRAIN. | pass | Confirms Confluent Schema Registry admission gates; does not inspect ad-hoc temporary diagnostic topics. |
| FIT-3: Silent Schema Drift | Seed a producer microservice that changes an existing Avro schema field from a long integer to a string without incrementing the schema version or verifying backward compatibility. | Confluent schema compatibility validator probe_breaking_schema_drift verifying publish rejection with diagnostic ERR_SCHEMA_BACKWARD_COMPATIBILITY_BREACH. | pass | Confirms Schema Registry server-side compatibility checks; does not evaluate bypasses using raw byte arrays. |
Residual Risk
- Latency spikes (up to 20 ms) in consumer group processing during AWS cloud network fiber optic degradation between AZ-1 and AZ-2. Accepted by Elena Rostova with local pod-level consumer thread scaling.
Traceability
| Claim | Classification | Source | Freshness |
|---|---|---|---|
| Rejection of uncoordinated multi-topic writes | derived | FIT-1 probe result | 2026-09-15 |
| Rejection of event schemas lacking declared grain | derived | FIT-2 probe result | 2026-09-15 |
| Rejection of breaking Avro schema mutations | derived | FIT-3 probe result | 2026-09-15 |
Verification
No validator was supplied, so no command was run.
Open Decisions
None.
Next steps
- Architecture Guild incorporates event streaming fitness probes into automated CI build checks.
- Platform squad configures Prometheus alerts monitoring consumer group lag and partition fetch latencies.
- Conduct quarterly disaster recovery drill simulating complete Kafka broker loss during peak payment streaming.
event-streaming-platform-and-kraft-kafka.pdf
PDF · document
Example file from a real run - the skill writes it into your workspace.
Connects securely to your tools. The creator never sees your data.
What you get
About this skill
What it does
This skill owns the architecture for continuously processing ordered or partially ordered event streams across producer, transport/log, processor state and sink boundaries. It defines identity, time, ordering, delivery, state, replay and recovery contracts without assuming that a broker, checkpoint, transaction or dashboard creates end-to-end correctness.
Use it when
- Continuous unbounded inputs require stable event, stream, topic, partition, producer and consumer identities
- Partitioning and keys determine ordering, parallelism, state locality and skew
- Event, ingestion and processing time require explicit watermark, window, trigger and late-data semantics
- Producer, broker, processor and sink delivery guarantees differ and must compose end to end
- Keyed/operator state, timers, checkpoints and sink commits need a consistent recovery cut
- Replay, backfill and correction must coexist with live traffic without duplicate or out-of-order effects
For example: “Our fraud rules run on a stream. During a network blip, 40 minutes of events arrived at once and every rule fired, blocking 12,000 legitimate cards.”
What you get
- architecture/streaming-architect/README.md
- architecture/streaming-architect/00-overview/streaming-architect-overview.md
- architecture/streaming-architect/verification/fitness-self-check.md
Plus one page per business module, only where your evidence calls for it: {module}/ingest.md, {module}/storage.md, {module}/serving.md, {module}/lineage.md, {module}/retention.md, {module}/quality.md.
All paths are relative to the output folder you choose.
What it will not do
Do not use merely to add one Kafka topic/consumer, implement a Flink/Beam job, configure Pub/Sub/Kinesis, capture database changes, publish an application event, build a queue worker, schedule/orchestrate batch ETL, or diagnose lag.
How it works
- Check latency genuinely requires streaming.
- Fix event time versus processing time.
- Define windows, watermarks and allowed lateness together.
- State the delivery and state semantics.
- Design backpressure and restart.
- Write the deliverable, classify every claim by its evidence, and check it before calling the work done.
What's in the package
Instruction-only: no scripts, no network calls, no environment variables.
- LICENSE.txt
- SKILL.md
- agents/openai.yaml
- assets/output-template-artifact.md
- assets/output-template-contract.md
- assets/output-template-domain.md
- assets/output-template-fitness.md
- assets/output-template-mechanism.md
- references/domain-rules.md
- references/operating-rules.md
- references/output-contract.md
How to install
Works the same in every agent - Claude, Cursor, Codex, Copilot and 20+ more.
- 1
Download the ZIP
Free skills download straight away. Paid skills unlock right after purchase.
- 2
Unzip into your skills folder
Every agent reads skills from one folder on your machine. Drop the unzipped folder in there.
- 3
Ask your agent to use it
Restart the agent if it was already running. It picks the skill up automatically - no config needed.
Skills folder by agent
Click the path to copy it. Create the folder if it does not exist yet.
Reviews
No reviews yet
Be one of the first to try it. Every listed skill passes our trust checks below.
Security scanned
Passed our 8-point scan before listing
Fresh listing
Recently published to Agensi
30-day refund
Not a fit? Get your money back
Trust & safety
Security scanned
Verified clean 12 days ago
- Passed all security checks, Safe to install