Event Streaming Platform and KRaft Kafka Architect

    1

    Architects event streaming: Apache Kafka KRaft clusters, deterministic partition keys, and cooperative sticky rebalancing.

    $9

    Secure checkout via Stripe

    30-day refund guarantee

    Converts to your local currency at checkout

    Security scanned

    Works with the AI tools you already use

    Claude CodeClaude CodeCursorCursorCodex CLICodex CLIMuseMuseOpenClawOpenClaw+21 more

    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; producer acks=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

    CriterionWhy it matters hereWeightSource of the weight
    Zero Event Loss on Broker Failure (Durability)Dropping transaction events caused incident STR-4919 ($3.2M penalty, 850k events lost).0.40David O'Reilly (Chief Messaging Architect)
    Elimination of Consumer Rebalance StormsConsumer group stalls paralyze payment authorization for hours under load.0.30Elena 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.15Fraud Risk Management SLA
    Schema Evolution & Strict CompatibilityBreaking schema changes crash consumer pods simultaneously across 45 services.0.15Enterprise Data Governance Guild

    Comparison

    Streaming Platform ArchitectureDurability GuaranteeTail Latency (180k EPS)Rebalance StabilityEvaluation
    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 QueueHigh (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.v1 is 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 CooperativeStickyAssignor protocol.
      • 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.

    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

    ClaimClassificationSourceFreshness
    180,000 events/sec across 45 microservicesprovidedStreaming platform capacity briefCurrent
    4.8 TB daily event log volumeprovidedVolumetric traffic profileCurrent
    Incident STR-4919 2.5-hour stall ($3.2M penalty)providedOperations forensic post-mortemHistorical
    Delivery latency target p99 <= 15 msprovidedReal-Time Payment Platform SLACurrent
    KRaft Apache Kafka 3.6 cluster selecteddecidedDavid O'Reilly & Elena Rostova2026-09-15
    Mandatory acks=all producer invariant INV-STRM-01decidedArchitectural invariant INV-STRM-012026-09-15

    Verification

    No validator was supplied, so no command was run.

    Reviewer self-check against streaming platform architecture standards:

    • Durability Rigor: PASS. acks=all and 6-way storage ensure zero lost events (STR-4919 resolved).
    • Partition Balance: PASS. MurmurHash2 on account_id across 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

    1. Platform Infrastructure squad deploys the 9-broker Apache Kafka KRaft cluster via Strimzi Operator on AWS EKS.
    2. Ingestion team attaches the Confluent Schema Registry and registers core payment Avro contracts.
    3. 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]ProbeEvidenceResultLimits of the claim
    FIT-1: Dual WriterSeed 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.passConfirms producer SDK configuration linters; does not evaluate raw TCP socket writes bypassing Kafka client libraries.
    FIT-2: Undefined GrainSeed 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.passConfirms Confluent Schema Registry admission gates; does not inspect ad-hoc temporary diagnostic topics.
    FIT-3: Silent Schema DriftSeed 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.passConfirms 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

    ClaimClassificationSourceFreshness
    Rejection of uncoordinated multi-topic writesderivedFIT-1 probe result2026-09-15
    Rejection of event schemas lacking declared grainderivedFIT-2 probe result2026-09-15
    Rejection of breaking Avro schema mutationsderivedFIT-3 probe result2026-09-15

    Verification

    No validator was supplied, so no command was run.

    Open Decisions

    None.

    Next steps

    1. Architecture Guild incorporates event streaming fitness probes into automated CI build checks.
    2. Platform squad configures Prometheus alerts monitoring consumer group lag and partition fetch latencies.
    3. Conduct quarterly disaster recovery drill simulating complete Kafka broker loss during peak payment streaming.

    event-streaming-platform-and-kraft-kafka.pdf

    PDF · document

    Generated

    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

    Design Kafka KRaft clusters with deterministic partition keysDefine event-time watermarks and late-data handling policiesConfigure cooperative sticky rebalancing for consumer groupsArchitect idempotent sinks and exactly-once processing boundariesMap keyed state recovery and checkpointing for Flink/Spark jobs

    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

    1. Check latency genuinely requires streaming.
    2. Fix event time versus processing time.
    3. Define windows, watermarks and allowed lateness together.
    4. State the delivery and state semantics.
    5. Design backpressure and restart.
    6. 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.

    ~30 seconds
    1. 1

      Download the ZIP

      Free skills download straight away. Paid skills unlock right after purchase.

    2. 2

      Unzip into your skills folder

      Every agent reads skills from one folder on your machine. Drop the unzipped folder in there.

    3. 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

    Listed12 days ago

    What's inside

    Frequently Asked Questions