Product
Product11 min readBy The Data Workers Team

You're on Debezium: It Streams Every Committed Change Out of Your Databases. Data Workers Owns Whether What Those Changes Become Downstream Is Right

Debezium streams every committed row change from your databases into Kafka. Data Workers catches what a schema change or incremental snapshot does to dbt, BI and the jobs downstream, and gets the fix approved.

Your platform team thinks in connectors, change events and offsets. A Debezium connector reads the PostgreSQL write-ahead log, the MySQL binlog, Oracle LogMiner or the SQL Server change tables, and turns each committed row change into an event with a before image, an after image, an op code and a source block. Most estates run the connectors on Kafka Connect, with schemas in Schema Registry and a sink, such as the Snowflake Kafka connector or the Debezium JDBC sink, landing events in the warehouse; others send changes from Debezium Server straight to Kinesis, Pub/Sub or HTTP. Connectors that need one keep a schema history topic, and when you need history again you send a signal: an execute-snapshot row in the signaling table, and the connector re-reads a table in chunks while streaming carries on.

Debezium is built to copy changes faithfully: a new column flows into the next event, a re-read table arrives as READ events next to live changes, and the connector stays RUNNING throughout. That is right for a change stream. It also means a dbt model and a fee job can turn a perfectly valid stream into a wrong number with every Debezium metric green. Data Workers watches what those events become in the warehouse, traces the change through dbt, BI and the jobs that read it, and gets the fix approved before anything acts on it.

Key takeaways

  • •Debezium keeps its job. Connectors, signals, snapshots, offsets, schema history and sinks stay with your platform team. Data Workers works on what the change events become downstream.
  • •RUNNING connectors get checked too. Uniqueness, row-count baselines, freshness and schema diffs on the landed tables turn a valid snapshot or schema change into a diagnosed incident before the next job reads it.
  • •Native where Debezium lives. Data Workers reads Kafka Connect connector status and registers approved Schema Registry versions natively, alongside Snowflake, PostgreSQL, Airflow and dbt, and connects to Debezium over the Kafka Connect REST API today.
  • •Every fix goes through a named person. The owner approves the holds and the dbt diff and applies them; signals and connector settings stay with the platform owner; each step leaves a receipt.
  • •Autonomy is set per domain. Start at L1 observe, move to L2 propose, and open L3 act reversibly for narrow classes once the record earns it.

Debezium is the change stream out of your databases. Data Workers is the owner of what it becomes downstream.

Debezium's job is to read every committed change once, in order, and never lose one. The job after the event lands is different: notice that a correct stream has produced a wrong table, tie it to a change in the connector, size what it touches, hold anything that would act on it, get the fix approved, rebuild and prove the number. One morning at a payments company, as an illustration, not a customer case:

TimeSystemWhat happens
Tue 10:05Kafka Connect + Schema RegistryThe fraud team wants risk_score in the warehouse. The payments platform owner removes it from the pay-postgres connector's column.exclude.list through the Connect REST API. New change events for public.payments carry the column, and version 7 of pay.public.payments-value adds it as an optional field; the BACKWARD compatibility check passes
10:12PostgreSQL + DebeziumSo that the 41 million existing payments carry the column too, he inserts an execute-snapshot signal (type incremental, data collection public.payments) into the debezium_signal table. Debezium re-reads the table in chunks while streaming continues, emitting each existing row as a READ event (op: r). The connector and its task stay RUNNING
10:13SnowflakeThe Snowflake Kafka connector appends every event to RAW_PAY.PAYMENTS_CDC, as it always has
11:00Airflow + dbtThe hourly DAG dbt_payments_hourly builds stg_pay__payments, an incremental model written when the connector was first set up. It appends events with op in c or r loaded since the last run and has no unique key: back then, READ events only came from the one-time initial snapshot. 9.4 million existing payments are appended again, and fct_merchant_volume, lifetime processed volume per merchant, jumps for the merchants with the longest history. Its not-null tests pass
11:08Data Workers + Snowflake + Datadogrun_quality_check finds payment_id in stg_pay__payments no longer unique, for the first time since the model shipped. The hourly row count the team records with monitor_metrics reads 9.4 million against a baseline near 2,500. On the team's Datadog dashboard of connector JMX metrics, TotalNumberOfReadEventsSeen, the snapshot read-event metric new in Debezium 3.7, has climbed since 10:12. Data Workers opens an incident
11:15Data Workers + Kafka Connect + dbtget_connector_status shows pay-postgres and its task RUNNING, so the stream is healthy. The run_quality_check profile shows RISK_SCORE added to the raw table after the 10:05 filter change, a safe additive change. Every duplicate is a READ event whose original arrived as a create months earlier. diagnose_incident names the cause: a snapshot re-read appended as new payments by an append-only model, with the snapshot about a quarter done. Blast radius: fct_merchant_volume, the Power BI "Merchant volume" report (a context-graph note the team recorded), and the Airflow DAG merchant_tier_review, which at 16:00 writes fee tiers to the billing service's PostgreSQL. On these numbers it would cut fees for 2,140 merchants; on a normal day it moves about a dozen
11:20PagerDuty + SlackData Workers raises a PagerDuty incident with the diagnosis and sends the approval request to the analytics engineering owner, with the payments platform owner copied. He confirms the 10:12 signal in the thread
11:45SpellbookThe owner reviews four proposals: pause dbt_payments_hourly and merchant_tier_review; a dbt diff that makes stg_pay__payments a merge on payment_id, keeping the latest event per payment, so a READ event updates the row and fills risk_score instead of adding one; a run plan to full-refresh stg_pay__payments and its children once Debezium reports the snapshot complete; and a runbook note for the platform owner on what to check before the next execute-snapshot signal. She approves all four, and the platform owner takes the note
11:55Airflow + GitHubShe pauses both DAGs and merges the dbt diff
13:52DebeziumThe connector's notification channel reports the incremental snapshot COMPLETED for public.payments; the platform owner posts it to the incident
13:55Data Workers + AirflowData Workers queues the approved full-refresh DAG with trigger_airflow_dag; Airflow finishes it at 14:30
14:40Data Workers + SnowflakeData Workers verifies: payment_id is unique, the hourly row count and the merchant-volume metric are back on their baselines, load lag is back on its baseline and risk_score is filled on new and re-read payments. The receipt records cause, approvals, runs, checks and the undo (revert the dbt diff and refresh again; the raw change log is untouched)
14:45AirflowThe owner resumes both DAGs. At 16:00 merchant_tier_review moves 11 merchants
Incident timeline across the stack: what Debezium, your team and Data Workers each do, step by step

Every part of Debezium did what it was asked to do: the column filter changed, the signal started an incremental snapshot, and the snapshot re-emitted each existing row once beside live streaming. Catching the problem takes knowledge Debezium was never meant to hold: that a staging model in another repo assumes READ events happen once, and that a fee job reads that model at 16:00.

JobWhat Debezium doesWhat Data Workers does
The captureReads the transaction log and emits every committed change with before and after images, from PostgreSQL, MySQL, Oracle, SQL Server, MongoDB and moreReads what the events land as, natively in the warehouse, and connects to Debezium over the Kafka Connect REST API
The schemaRecords DDL in the schema history topic and evolves Avro or Protobuf schemas through Schema RegistryDiffs the landed schema against a stored baseline, classifies the change and traces every model, dashboard and job that reads the column
The re-readRuns incremental snapshots on a signal, in chunks beside live streaming, and reports STARTED, IN_PROGRESS and COMPLETED notificationsChecks what the re-read did to every table built from the stream: uniqueness, row counts against a baseline, freshness and nulls
The signalJMX metrics, notifications and, in the Debezium Platform (preview releases), dashboards and threshold alerts on pipelinesOpens an incident on the data itself, so a RUNNING connector with a broken downstream model still gets a cause and an owner
The fixApplies whatever connector config and signals the platform owner sendsProposes the holds and the dbt diff to the owner, and a runbook note to the platform owner, who approve and apply them
The proofKeeps offsets and schema historyRe-checks the tables and writes a receipt: what changed, who approved it, how it was checked, how to undo it

Why doesn't Debezium just do this itself?

Because Debezium is a change data capture project, and it made the right choices for that job. A connector's promise is that every committed change reaches the stream in order, without locking your database. Incremental snapshots exist so you can re-read a table without stopping streaming, and when a snapshot read collides with a live change, Debezium keeps the live one. What happens to an event after the sink writes it, in dbt, BI, a fee job or a reverse ETL sync, is outside the connector by design. Debezium can't know that one consumer treats a READ event as a new payment.

Debezium's newer features point the same way. Release 3.7, generally available on September 29, 2026, adds incubating TiDB, Milvus and SQLite connectors, the first official Debezium CLI, alerting and deployment outside Kubernetes in the Debezium Platform, and a dedicated metric for snapshot read events, because dashboards used to show zero activity while a snapshot ran. Its AI work is pipeline work: embeddings transforms that add vectors to change events, and PyDebeziumAI, which keeps vector stores in step with your database for LangChain and LangGraph apps. The project publishes no MCP server, and its Platform AI chatbot is an unreleased experiment. All of that makes the change stream better.

Owning whether what the stream becomes is right, across the database, Kafka, the warehouse, dbt, BI and the jobs that act on them, is a different product: a context graph of every table and consumer, blast-radius scoping, named approvers, a recorded undo and receipts an auditor can read. That is Data Workers. More in is it safe to let AI agents change production data.

Every tool owns a slice. Data Workers covers the whole lifecycle

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 Debezium connectors already there.

Spider chart of ten jobs a data team does: Data Workers covers the whole list, Debezium goes deep on its own area
StageData WorkersDebeziumWhy we scored it this way
Catalog & Context92The schema history topic records each captured table's DDL for the connector. Data Workers keeps one governed context graph of what each table means, who owns it and what reads it.
Analytics & Insights81Not Debezium's job: it emits change events, while the numbers live in BI. Data Workers answers data questions from governed definitions with lineage behind every number.
Data Quality82Change events carry each row faithfully, before and after. Data Workers checks the landed tables for nulls, uniqueness and volume, tracks lateness against a baseline your team records, and turns a bad load into an incident.
Observability & Incidents8.54JMX metrics, notifications and, since 3.7, Platform alerting show connector and snapshot health. Data Workers diagnoses the data incident while every task is RUNNING and verifies the fix.
Pipelines & Ingestion8.59Debezium's home stage: log-based CDC from PostgreSQL, MySQL, Oracle, SQL Server, MongoDB and more, with signals and incremental snapshots. Data Workers plans the downstream reruns for the owner.
Schema & Migration85Captures DDL and evolves Avro or Protobuf schemas through Schema Registry converters. Data Workers traces the landed change through dbt, BI and the jobs that read it.
Governance & Access8.52Kafka Connect and Platform access decide who changes a connector. Data Workers routes every data change to a named approver.
Security & Privacy83Column filters and hash or character masking keep fields out of the stream. Data Workers leaves a receipt on every data change.
Cost / FinOps82Open source; the cost sits in Kafka and the warehouse. Data Workers traces Snowflake credits to the dbt model behind them through query tags.
MLOps & Models7.52Embeddings SMTs and PyDebeziumAI keep vector stores in step with the database. Data Workers keeps the data under models and agents healthy.

How Debezium and Data Workers work together

How Data Workers fits with Debezium: your coding agent on top, Data Workers in the middle, your estate underneath

Data Workers reads the open parts of a Debezium estate natively. Point it at the Kafka Connect and Schema Registry URLs 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, it registers it with register_stream_schema only after a named approval. The Schema Evolution agent traces a landed schema change to its readers, the Incident Debugging agent builds the diagnosis, and the Streaming agent reviews stream configurations with recommendations a person approves.

By design, Data Workers does not send signals, run snapshots, restart connectors, reset offsets, edit connector config or touch the schema history topic. Those steps belong to the platform owner, through the signaling table, the Connect REST API, the Debezium CLI or the Debezium Platform; Data Workers proposes them, scopes what they touch and verifies afterwards. Connector and snapshot metrics come from the monitoring you already run, with Datadog and OpenTelemetry as native reads. The rest of this incident connects natively too: PostgreSQL, Snowflake, Airflow, dbt, PagerDuty, GitHub (pull-request review) and Slack, among 50+ connectors. Power BI connects over its API or MCP server today, and the report enters the blast radius through the team's context-graph note.

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-incidents

List the tools with your client's own command (/mcp in Claude Code). In this incident: run_quality_check (dw-quality) runs the uniqueness and null checks on Snowflake; monitor_metrics and diagnose_incident (dw-incidents) flag the row count against its baseline and name the cause; get_connector_status, send_pagerduty_alert and trigger_airflow_dag (dw-connectors) read the connector state, raise the incident and queue the approved full refresh on the team's Airflow 2 deployment; assess_impact (dw-schema) sizes the landed change; blast_radius_analysis and trace_cross_platform_lineage (dw-context-catalog) map what the duplicates reach; remediate re-checks the quality assertions after the rebuild and escalates any failure to a person.

The agents run in your infrastructure, next to Kafka Connect, and hold the warehouse credentials and model key. Your data stays in your systems; the hosted Conductor sees workflow metadata only. More in where does our data go.

One incident, L0 to L4, set per domain:

The autonomy ladder: L0 manual, L1 observe, L2 propose, L3 act reversibly, L4 autonomous
  • •L0 manual. The fee job runs at 16:00 and 2,140 merchants pay less; finance spots the revenue dip two days later.
  • •L1 observe. Data Workers flags the duplicates at 11:08 with the cause and blast radius. Nothing changes.
  • •L2 propose. Data Workers proposes the holds, the dbt diff, the run plan and the runbook note; nothing moves until the named owner approves.
  • •L3 act reversibly. For a class with a clean record, Data Workers queues the rebuild itself and sends any failed check to a person. Signals stay with the platform owner.
  • •L4 autonomous. For a scoped domain, Data Workers checks the CDC tables on every load and handles the classes of fix you have opened.

What changes for your team

Six jobs that run on autopilot with Data Workers next to Debezium, with a concrete example of each
  • •Snapshots stop being a gamble. A signal on a busy table becomes a watched event: models built from the stream are checked while the re-read runs.
  • •On-call starts from a cause. The incident carries the connector state, the schema change, the duplicates, the damage and the proposed fix.
  • •Platform and analytics engineers share one record. Both see the same incident and receipt in Spellbook Data Catalog (in preview), linked from the PagerDuty incident.
  • •Consumer assumptions get written down. Each incident leaves a runbook note on what the next snapshot, filter change or schema version will do downstream.

Keep Debezium, or consolidate?

Keep Debezium if you love it; Data Workers works with it from day one. Many teams consolidate once Data Workers runs that slice too.

For most Debezium teams the answer is keep it: Apache 2.0, log-based CDC across a dozen databases, governed in the Commonhaus Foundation and available as a supported Red Hat build, with no per-row price and full control of where the events go. What teams consolidate is the tooling around the landed data: a separate observability tool, hand-written duplicate checks and "watch the dashboards after a snapshot" runbooks. Many estates run Debezium on Apache Kafka or Confluent, and pair it with a managed mover such as Estuary or Fivetran; one incident record spans them.

Weighing a build on the Connect REST API and a coding agent? Reading connector status from a chat is easy; the context graph, approvals, undo and receipts are the work. See build it ourselves with Claude Code and MCP servers.

The case for your CFO

The outcome. When a database change or a re-read of history reaches the tables behind fees, revenue or risk, it is caught and corrected before a job charges, pays or reports on it.

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. Signals, snapshots and connector settings stay with the platform owner. No agent can promote its own work, and an org-wide stop halts all autonomous dispatch. Every change carries a receipt with the blast radius and the undo.

Why now. CDC streams now feed pricing, risk and AI applications as well as dashboards, so a misread stream reaches a customer in hours.

The first win. L1 on the CDC tables that feed money numbers: the next snapshot or column change becomes a diagnosed incident before the next job reads it.

What stays the same. Nothing migrates: Debezium, Kafka Connect, your sinks, dbt, the warehouse and the on-call rota. For the numbers, see the ROI of agentic data operations.

The sentence for upstairs: "Debezium streams every change out of our databases; Data Workers makes sure what those changes become is right, and gets it fixed with our approval before we charge, pay or report on it."

Getting started

Start with a pilot. Pick the Debezium connectors whose tables feed the numbers leaders and customers see, give Data Workers read access to Kafka Connect, Schema Registry, those warehouse schemas and the dbt project, 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 earn it. 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 Debezium? It reads Kafka Connect connector and task status and Schema Registry natively, and connects to Debezium itself over the Kafka Connect REST API today. The tables Debezium's sinks write, in Snowflake, BigQuery, Databricks or PostgreSQL, are native reads, as are dbt and Airflow, Dagster or Prefect.

Will Data Workers send signals or run snapshots for us? No. Signals, incremental and blocking snapshots, connector config, offsets and the schema history topic stay with the platform owner. Data Workers proposes the step, says what it will touch downstream and verifies afterwards.

Our connectors are RUNNING and the Platform alerts are quiet. How can the data be wrong? RUNNING means Debezium is capturing changes. An incremental snapshot re-emits existing rows as READ events, a column filter change adds a field, a new schema version passes compatibility; each is correct, and each can break a consumer that assumed otherwise. Data Workers checks what the events became.

Does a schema change in the source break our pipeline? Usually not at the connector: Debezium records the DDL and evolves the schema in Schema Registry. The break shows up later, in a staging model's cast, a join or a BI field. Data Workers diffs the landed schema and traces every model, dashboard and job that reads the column.

We run Debezium Server, not Kafka Connect. Does this still work? Yes. Data Workers reads the tables where Debezium Server's sink writes and the metrics your monitoring collects. Since 3.4, Debezium Server can also emit OpenLineage events to the lineage backend your team browses.

Sources

  • •Debezium, Releases overview (3.7 2026-09-29, 3.6 2026-09-18, 3.5 2026-06-02, 3.4 2026-03-30 "OpenLineage output set for Debezium Server"; 2.7 stable), https://debezium.io/releases/ (checked Oct 3, 2026)
  • •Debezium blog, "Debezium 3.7 Release Summary" (Sep 29, 2026: TiDB, Milvus, SQLite incubating; Debezium CLI; Platform alerting and host-based deployment; TotalNumberOfReadEventsSeen; HTTP notification channel; Embeddings SMT for the OpenAI API; PyDebeziumAI), https://debezium.io/blog/2026/09/29/debezium-3-7-final-released/ (checked Oct 3, 2026)
  • •Debezium blog, "Debezium 3.7.0.CR1 Released" (Sep 22, 2026: PyDebeziumAI for LangChain and LangGraph), https://debezium.io/blog/2026/09/22/debezium-3-7-cr1-released/ (checked Oct 3, 2026)
  • •Debezium blog, "Debezium Platform: Mid-Year Update (2026)" (Jul 10, 2026: monitoring, AI chatbot POC not released), https://debezium.io/blog/2026/07/10/debezium-platform-status-update/ and "Recent improvements in the Debezium Platform UI" (Sep 28, 2026), https://debezium.io/blog/2026/09/28/platform-stage-ui-improvements/ (checked Oct 3, 2026)
  • •Debezium docs, Signaling (channels; execute-snapshot, stop-snapshot, pause-snapshot, resume-snapshot; additional-conditions), https://debezium.io/documentation/reference/stable/configuration/signalling.html (checked Oct 3, 2026)
  • •Debezium docs, Notifications (sink, log, JMX channels; incremental snapshot STARTED, IN_PROGRESS, TABLE_SCAN_COMPLETED, COMPLETED), https://debezium.io/documentation/reference/stable/configuration/notification.html (checked Oct 3, 2026)
  • •Debezium docs, PostgreSQL connector (READ events op: r; incremental snapshots in parallel with streaming; chunk size 1024; collision handling), https://debezium.io/documentation/reference/stable/connectors/postgresql.html (checked Oct 3, 2026)
  • •Debezium, Governance (steering committee, committers, Commonhaus Foundation code of conduct), https://debezium.io/community/governance/ and Commonhaus Foundation projects, https://www.commonhaus.org/ (checked Oct 3, 2026)
  • •Red Hat, Red Hat build of Debezium documentation, https://docs.redhat.com/en/documentation/red_hat_build_of_debezium (checked Oct 3, 2026)
  • •Debezium on GitHub (Apache 2.0; tag v3.7.0.Final), https://github.com/debezium/debezium (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)