You're on Apache Kafka: It Is the Event Backbone Your Platform Team Runs. Data Workers Owns Whether What Consumers Land Is Right
Apache Kafka 4.x carries every event. Data Workers checks what your consumers land downstream, traces a wrong number back to the consumer change and gets the fix approved.
Your platform team runs Apache Kafka itself, on virtual machines or on Kubernetes with Strimzi. Since Kafka 4.0 KRaft runs the metadata and ZooKeeper is gone. App teams produce to topics, partitions spread the load, and consumer groups read each partition in order. Kafka Connect moves topics into Postgres, Iceberg or the warehouse, and a Schema Registry sits beside the cluster. Since Kafka 4.2 in February 2026, share groups (Queues for Kafka) are production-ready, so a team can put more consumers on a topic than it has partitions.
Kafka is excellent at its job: every event, durable, replicated, in partition order. What a consumer does with that event is a different job: it can read every record and still write the wrong row. Data Workers watches what your consumers land, traces a wrong number to the change behind it, and gets the fix approved before a customer or an invoice builds on it.
Key takeaways
- •Kafka keeps its job. Brokers, topics, groups, Kafka Connect and Kafka Streams stay with your platform team. Data Workers works on what consumers write and what reads it.
- •Native where Kafka is open. Kafka, Kafka Connect and Schema Registry connect natively. Data Workers reads connector and task status, checks proposed schemas against the registry's rule and reads the tables consumers write in Postgres, Snowflake and BigQuery.
- •Correct delivery, checked downstream. When a consumer change, a connector or a schema version leaves wrong rows behind, Data Workers turns it into one incident with the likely cause and the jobs that would spread the damage.
- •Every fix goes through a named person. The owner approves the hold, the consumer diff and the repair, and runs the Kafka side. Data Workers registers a schema version only after a named approval.
- •Autonomy is set per domain. Start at L1 observe, move to L2 propose, and open L3 act reversibly for narrow classes once the record supports it.
Kafka is the event backbone your platform team runs. Data Workers owns whether what your consumers land is right.
The job downstream of Kafka is to notice that a table built from its events is wrong, find the change behind it, hold anything that would act on it, get the fix made and prove the result. Here is one Tuesday at a regional parcel carrier during peak season. This is an illustration, not a customer case.
The carrier runs Kafka 4.2 on Strimzi. Every depot scan lands on the topic shipment-events, 12 partitions keyed by parcel ID. A service called status-writer reads it and upserts each parcel's latest status into the PostgreSQL table shipments_current, last write wins; the tracking page and the late-delivery credit job read that table. The Apache Iceberg sink connector on Kafka Connect also lands every event in the Iceberg table logistics.shipment_events, registered in Polaris and read by Snowflake.
| Time | System | What happens |
|---|---|---|
| Tue 10:40 | Kafka + Kubernetes | Peak scan volume is outrunning 12 consumers. The platform team ships status-writer 3.0: it moves from the consumer group to a new share group and scales from 12 pods to 36. Throughput triples and lag drops to near zero |
| 10:41 | PostgreSQL | In a share group, several pods now take records from the same partition. A parcel's out_for_delivery and delivered scans, seconds apart, go to different pods; when the older one commits last, it overwrites delivered in shipments_current |
| 10:41 | Kafka Connect + Polaris | The Iceberg sink connector, still on its own consumer group, writes every event in partition order. Its connector and all four tasks stay RUNNING |
| 13:05 | Data Workers | Parcels marked delivered per hour in shipments_current, a metric the team records hourly with monitor_metrics, reads 68% of its baseline for the 12:00 hour while scan volume on the Iceberg table in Snowflake is normal. Data Workers opens an incident |
| 13:20 | Data Workers + Kafka Connect + PostgreSQL | get_connector_status shows the Iceberg sink and its tasks RUNNING, so the event log is complete. The drop starts in the first full hour after the 10:40 rollout, recorded as the only recent change to the consumer, and diagnose_incident ranks a consumer regression first. Blast radius: the Prefect deployment issue_late_credits, which credits customers at 18:00 for parcels not delivered on time, the tracking page, and the Tableau "Carrier on-time" workbook |
| 13:25 | Slack | The approval request goes to the owner of status-writer, with the logistics data owner copied |
| 13:35 | PostgreSQL + Snowflake | The owner runs the team's reconciliation query against shipments_current and the event log: 8,412 parcels updated since 10:40 show an older status than their latest scan; none before 10:40 |
| 13:45 | Spellbook | The owner reviews three proposals: pause issue_late_credits; a diff to status-writer that keeps the share group but only overwrites a row when the incoming scan is newer (WHERE shipments_current.event_ts < EXCLUDED.event_ts); and a repair statement that resets the 8,412 rows from the event log, saving their current values to a before-image table first. She approves all three |
| 13:50 | Prefect | She pauses the issue_late_credits schedule |
| 14:15 | GitHub + Kubernetes | She merges the diff and rolls out status-writer 3.0.1, still at 36 pods |
| 14:30 | PostgreSQL | She runs the repair in one transaction; the before-image table holds the undo. Her reconciliation query now returns zero parcels behind their latest scan |
| 14:35 | Data Workers + Prefect | Data Workers queues the approved rerun of the hourly_delivery_metrics flow with trigger_prefect_flow to recompute the metric for 11:00 to 14:00 |
| 15:00 | Data Workers | It verifies: delivered per hour is back on baseline for every hour since 10:00, and run_quality_check finds normal volume in the Iceberg table. The receipt records cause, approvals, runs, checks and the undo (restore from the before-image table, then roll back to 3.0) |
| 15:05 | Prefect | The owner resumes the schedule. At 18:00 credits go only to parcels that were late |
| Wed 08:30 | Tableau | The carrier review opens on the correct on-time rate |

Every part of Kafka did its job. Kafka's documentation is explicit: consumer groups pair ordering with scale, while in a share group "partitions may be assigned to multiple consumers", and KIP-932 puts the focus "on sharing rather than ordering". The choice was reasonable for throughput; it broke an assumption one table relied on. Catching it takes knowledge outside the cluster: that shipments_current is last write wins, that a credit job reads it at 18:00, and that a full log exists to repair it from.
| Job | What Kafka does | What Data Workers does |
|---|---|---|
| The log | Stores every event durably, replicated, in partition order, with KRaft managing the cluster | Reads Kafka, Kafka Connect and Schema Registry natively, plus the databases and warehouses consumers write |
| The consumers | Assigns partitions to consumer groups, shares records across share groups, tracks offsets and lag | Checks what consumers land: volume, uniqueness, freshness and the team's key ratios against a baseline |
| The connectors | Runs Kafka Connect sinks and sources, with connector and task state over its REST API | Reads connector and task status with get_connector_status and ties each connector to the tables it feeds |
| The schema | Carries bytes; the registry beside it enforces each subject's compatibility rule | Checks a proposed schema against the registry's rule, traces the change to every table that reads it, registers a version after approval |
| The fix | Applies whatever consumer, topic, offset or connector change the owner makes | Proposes the hold, the code diff and the repair to a named owner, then queues the approved downstream runs |
| The proof | Keeps offsets, logs and metrics | Re-checks the data and writes a receipt: what changed, who approved it, how it was checked, how to undo it |
Why doesn't Kafka just do this itself?
Because Kafka is a log, by design, and that is why it scales. The broker does not read your payloads or know what a consumer writes to Postgres. Its guarantees are precise and strong: durability, replication, ordering within a partition for a consumer group, delivery counting and acknowledgements for a share group. None of them can say whether the row a consumer wrote is the right row, and they shouldn't try: that would make the broker depend on every application built on it.
Apache Kafka also ships no AI assistant or MCP server of its own. Vendors and the community offer MCP servers for Kafka administration, each scoped to operating the cluster: topics, configs, connectors. That is the right scope for cluster tooling.
Owning whether the data is right from a consumer service to Postgres, Snowflake, Prefect and Tableau is a different product with a different liability: a context graph from topic to dashboard, blast-radius scoping, approvals that name a person, a recorded undo and receipts an auditor can read. That is Data Workers. For how that holds up in production, read is it safe to let AI agents change production data.
Every tool owns a slice. Data Workers covers the whole lifecycle
Kafka owns carrying every event between producers and consumers. Each point tool adds another console, contract and handoff. Data Workers covers the whole lifecycle with one context, one approval flow and one audit trail, and builds on the Kafka cluster already there.

| Stage | Data Workers | Apache Kafka | Why we scored it this way |
|---|---|---|---|
| Catalog & Context | 9 | 3 | Topics and consumer groups are named in the cluster; meaning and ownership live in wikis and registries beside it. Data Workers keeps one governed context graph from the topic to the table, the job and the dashboard. |
| Analytics & Insights | 8 | 1 | Not Kafka's job: it carries events, while business numbers live in the warehouse and BI. Data Workers answers data questions from governed definitions with lineage behind every number. |
| Data Quality | 8 | 3 | Kafka stores what producers write; schema rules sit in a registry next to it. Data Workers checks what consumers land, so an event processed in the wrong order still raises an incident. |
| Observability & Incidents | 8.5 | 4 | JMX metrics, consumer and share partition lag show the cluster's health. Data Workers diagnoses the data incident across the consumer, the tables and the jobs, proposes the fix and verifies it. |
| Pipelines & Ingestion | 8.5 | 9 | Kafka's home stage: durable partitioned logs, consumer and share groups, Kafka Connect and Kafka Streams. Data Workers plans the repair and queues the downstream reruns for the owner. |
| Schema & Migration | 8 | 4 | Kafka carries bytes; schema evolution is managed in a registry and in client code. Data Workers checks a proposed schema against the registry's rule and traces the change to every table that reads it. |
| Governance & Access | 8.5 | 4 | ACLs and the KRaft-managed cluster control who reads and writes each topic. Data Workers routes every data change to a named approver and records the decision. |
| Security & Privacy | 8 | 5 | TLS, SASL and ACLs protect the stream in transit and at the broker. Data Workers leaves a receipt on every data change. |
| Cost / FinOps | 8 | 2 | Brokers, partitions and retention are sized and paid for by the team that runs them. Data Workers traces Snowflake spend to the dbt model behind it. |
| MLOps & Models | 7.5 | 2 | Kafka feeds features and models with live events. Data Workers keeps the data under models and agents healthy. |
How Kafka and Data Workers work together

Data Workers reads the open parts of your Kafka estate natively. Point it at the Kafka Connect and Schema Registry endpoints your platform team runs. It reads each connector's state, its tasks and their worker assignment with get_connector_status, and checks a proposed schema against the subject's compatibility rule before anyone registers it. When the fix is a new schema version, Data Workers proposes it and registers it with register_stream_schema only after a named approval; that action sits in the stricter class reserved for changes to the live data stack. The Schema Evolution agent traces a schema change to everything that reads it, and the Streaming agent reviews stream configurations with recommendations a person approves.
By design, Data Workers does not restart connectors, reset offsets, create or change topics, alter partitions or retention, change ACLs or redeploy consumers. Those steps belong to the owner, through kafka-consumer-groups, the Connect REST API or Strimzi custom resources. Data Workers proposes them, scopes what they touch and verifies afterwards. Consumer and share-group lag comes from the monitoring you already run, such as Prometheus and Grafana or Datadog; Datadog and OpenTelemetry are native reads. The rest of the estate in this incident connects natively too: PostgreSQL, Snowflake, Prefect, Tableau (read), GitHub (pull-request review) and Slack, among 50+ connectors.
The Data Workers side follows the documented client setup: clone the open-source repo and add each agent's start-agent.sh entry to your client.
# Example: Data Workers agents in Claude Code, from a clone of the open-source repo
claude mcp add --scope user dw-connectors -- "$(pwd)/start-agent.sh" dw-connectors
claude mcp add --scope user dw-schema -- "$(pwd)/start-agent.sh" dw-schema
claude mcp add --scope user dw-catalog -- "$(pwd)/start-agent.sh" dw-context-catalog
claude mcp add --scope user dw-quality -- "$(pwd)/start-agent.sh" dw-quality
claude mcp add --scope user dw-incidents -- "$(pwd)/start-agent.sh" dw-incidentsList the tools with your client's own command (/mcp in Claude Code). In this incident: monitor_metrics and diagnose_incident (dw-incidents) flag the drop and rank the cause; get_connector_status (dw-connectors) confirms the sink; trace_cross_platform_lineage and blast_radius_analysis (dw-context-catalog) map what reads shipments_current; run_quality_check (dw-quality) checks the Iceberg table in Snowflake; trigger_prefect_flow queues the approved rerun. For a schema change, check_compatibility and assess_impact (dw-schema) size it downstream. A Kafka MCP server for cluster administration can sit in the same client.
The agents run in your infrastructure with your credentials and model key; your data stays in your systems and the hosted Conductor sees workflow metadata only. More in where does our data go.
One incident, L0 to L4, set per domain:

- •L0 manual. Customer service hears on Wednesday about late credits for parcels on doorsteps, after the money has gone out.
- •L1 observe. Data Workers flags the drop at 13:05 with the likely cause and the blast radius. Nothing changes; the owner starts from the answer.
- •L2 propose. Data Workers proposes the hold, the consumer diff and the repair. Nothing reaches production until the named owner approves; an unanswered request expires and escalates, never auto-grants.
- •L3 act reversibly. For a class with a clean record, Data Workers carries out reversible steps inside the domain you open, such as queueing the approved reruns, verifies them and sends any failed check to a person. Consumers, topics, offsets and repairs to Postgres stay with the owner.
- •L4 autonomous. For a scoped domain, once the owner's fix lands, Data Workers queues the metric reruns, re-checks every hour and posts the receipt without waiting for a person. Consumer code, topics and repairs still go through the owner.
Who sets those levels: who owns the agents, how approvals work for AI data agents and autonomy levels L0 to L4 explained.
What changes for your team

- •Producers, consumers and data owners share one record. The platform team, the service team and the logistics data team see the same incident and receipt in Spellbook Data Catalog (in preview), linked from Slack.
- •Consumer changes get a downstream check. A new group type or client version is weighed against the tables it writes and the jobs that read them.
- •Jobs that act on data get a guard. When a credit run, a reverse ETL sync or a partner feed reads a table that just went wrong, the hold is proposed before it runs.
- •Repairs have a plan. The code fix, the repair from the log and the reruns arrive as one approved sequence with a recorded undo.
For the operating and governing angles, see Kafka operations automation with AI agents, top Kafka topic governance tools, how to monitor Kafka consumer lag and Kafka vs Kinesis.
Keep Kafka, or consolidate?
Keep Apache Kafka if you love it; Data Workers works with it from day one. Many teams consolidate once Data Workers runs that slice too.
For almost every Kafka team the answer is keep it: an open, proven log with share groups, KRaft and a large connector ecosystem is hard to beat. What teams consolidate is the tooling around the landed data: a separate data observability tool and runbooks that say "restart the connector and eyeball the dashboard". If you buy Kafka as a managed service, see you're on Confluent or you're on Amazon Kinesis and MSK; if Flink processes your topics, see you're on Apache Flink. Data Workers integrations lists what connects natively.
If you are weighing building this yourself on a Kafka MCP server and a coding agent, read build it ourselves with Claude Code and MCP servers. Listing topics from a chat is easy; the context graph, approvals, undo and receipts are the work.
The case for your CFO
The outcome. When a change around the stream quietly writes wrong rows into the tables behind credits, billing or customer pages, it is caught within hours and corrected before money moves, with a record of how it was checked.
The risk story. At L0 and L1, agents only read. At L2 they propose and a named person approves; an unanswered request expires and escalates, never auto-grants. At L3 they act on reversible changes inside the domains you open; L4 is a later choice per domain. Brokers, topics, consumers, connectors and offsets stay with the owner. No agent can promote its own work, and an org-wide stop halts all autonomous dispatch. Every change carries a receipt: what changed, who approved it, the blast radius and how to undo it.
Why now. Share groups and Iceberg sinks put more consumers on the same events, writing more tables. A wrong row reaches an automated job before a person sees it.
The first win. L1 on the consumers that write tables money depends on: each one's output is checked against its baseline, and a wrong number becomes a diagnosed incident before the next job acts on it.
What stays the same. Kafka and its operators, your consumers, connectors, databases, warehouse and on-call rota. For the numbers, see the ROI of agentic data operations.
The sentence for upstairs: "Kafka carries our events; Data Workers makes sure what our systems write from them is right, and gets it fixed with our approval when it isn't, before customers or finance act on it."
Getting started
Start with a pilot. Pick the consumers and connectors that write the tables behind credits, billing or customer-facing numbers, give Data Workers read access to Kafka Connect, Schema Registry and those tables, record the ratios that define "right" for them, and run at L1 for a few weeks. Then turn on L2 for one domain, and open L3 for a narrow class once the receipts show the agents were right. The pilot path is on the pricing page, and the pilot is credited in full against the first year.
FAQ
How does Data Workers connect to Apache Kafka? Kafka, Kafka Connect and Schema Registry connect natively. Data Workers reads connector and task status over the Connect REST API, checks proposed schemas against the registry's compatibility rule, and reads the tables consumers write in Postgres, Snowflake and BigQuery.
Our Kafka Connect connector shows RUNNING. Why would the data still be wrong? RUNNING means the connector and its tasks are alive, not that a field landed, a record skipped the dead-letter topic or no other consumer wrote over the rows. Data Workers checks the tables themselves and uses connector status as one input to the diagnosis.
Should we move consumers to share groups? Often yes, for work-queue patterns and throughput: Kafka 4.2 made share groups production-ready. Check what each consumer writes first. Several consumers can take records from the same partition, so a consumer that upserts the latest state per key loses the per-key order a consumer group gives it. Guard the write (only newer events overwrite) or keep that consumer in a consumer group.
Will Data Workers restart connectors, reset offsets or change topics? No. It reads status, subject versions and the tables. Restarts, offset resets, topics, partitions, retention, ACLs and consumer deploys stay with the owner. The one Kafka-side write Data Workers makes is registering a schema version, and only after a named approval.
We are planning a Kafka 4 upgrade or a KRaft migration. Does Data Workers help? The cluster work stays with your platform team. Data Workers watches what consumers land through the change and turns any shift into an incident with a likely cause, so the data side of the upgrade has a record.
Sources
- •Apache Kafka, downloads (4.3.1 released Jun 25, 2026; 4.3.0 May 22, 2026; 4.2.2 Sep 29, 2026; 4.2.0 Feb 17, 2026; 4.0.0 Mar 18, 2025), https://kafka.apache.org/community/downloads/ (checked Oct 3, 2026)
- •Apache Kafka blog, "Apache Kafka 4.0.0 Release Announcement" (Mar 18, 2025; "the first major release to operate entirely without Apache ZooKeeper", KRaft by default), https://kafka.apache.org/blog/2025/03/18/apache-kafka-4.0.0-release-announcement/ (checked Oct 3, 2026)
- •Apache Kafka blog, "Apache Kafka 4.2.0 Release Announcement" (Feb 17, 2026; share groups "production-ready"; multiple consumers "concurrently process records from the same partitions"), https://kafka.apache.org/blog/2026/02/17/apache-kafka-4.2.0-release-announcement/ (checked Oct 3, 2026)
- •Apache Kafka blog, "Apache Kafka 4.3.0 Release Announcement" (May 22, 2026; 25 KIPs; additional share group configurations), https://kafka.apache.org/blog/2026/05/22/apache-kafka-4.3.0-release-announcement/ (checked Oct 3, 2026)
- •Apache Kafka 4.3 documentation, Design: share groups ("partitions may be assigned to multiple consumers"; consumer groups combine ordering with scale), https://kafka.apache.org/43/design/design/ (checked Oct 3, 2026)
- •KIP-932, Queues for Kafka (Ordering: records in a share-partition "can be delivered out of order"; "The focus in this KIP is on sharing rather than ordering"; key-based ordering listed as future work), https://cwiki.apache.org/confluence/display/KAFKA/KIP-932%3A+Queues+for+Kafka (checked Oct 3, 2026)
- •Data Workers open-source repository (tool registrations in dw-connectors, dw-schema, dw-context-catalog, dw-quality, dw-incidents), https://github.com/DataWorkersProject/dataworkers-claw-community (checked Oct 3, 2026)
- •Data Workers client setup guide, https://dataworkers.io/opensource-docs/client-setup/ (checked Oct 3, 2026)