You're on Apache Flink: It Computes on Every Event as It Arrives. Data Workers Owns Whether What Your Jobs Write Is Right
Apache Flink 2.3 computes on events as they arrive. Data Workers checks what your Flink jobs write, traces a wrong number to the release or restore behind it and gets the fix approved.
Your streaming team runs Apache Flink. Jobs read Kafka topics, keep state per key and write results to Kafka, Iceberg, Paimon or a database with exactly-once checkpoints, mostly in Flink SQL. Before every upgrade the job takes a savepoint and the new version restores from it, so running totals, sessions and windows carry across releases. You run it yourself on the Flink Kubernetes Operator (1.16, September 15, 2026), or buy it managed: Confluent Cloud for Apache Flink, Amazon Managed Service for Apache Flink (supports Flink 2.3 since July 20, 2026) or Ververica. Flink 2.3, released June 25, 2026, is the stable line, with 1.20 as the LTS.
Flink is excellent at its job: correct state, exactly once, at very high event rates. Whether the numbers a job writes are right is a different job: a job can restore cleanly, stay RUNNING with every checkpoint green, and still write wrong results into the tables your warehouse, app and marketing tools read. Data Workers watches what your Flink jobs land, traces a wrong number to the release or restore behind it, and gets the fix approved before a customer or a campaign acts on it.
Key takeaways
- •Flink keeps its job. Jobs, state, checkpoints, savepoints and deploys stay with your streaming team. Data Workers works on what jobs write and on everything that reads it.
- •Connected from day one. Data Workers connects to Flink over its REST API or an MCP server today, and reads the Kafka Connect sinks, Schema Registry and warehouse tables around your jobs natively.
- •Wrong results, caught downstream. When a release, a restore or a plan change leaves wrong numbers behind, Data Workers raises one incident with the likely cause and everything that would act on the damage.
- •Every fix goes through a named person. The owner approves the holds and the rollback, and runs the Flink side: deploys, stops, savepoints and restores.
- •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.
Flink is the engine that computes on every event. Data Workers owns whether what its jobs write is right.
Here is one Tuesday at a grocery chain's loyalty program. This is an illustration, not a customer case.
The chain runs Flink 2.2 on the Kubernetes Operator. Every purchase lands on the Kafka topic points-events. The Flink SQL job points-balance keeps each member's running balance in state (SUM(points) GROUP BY member_id) and writes every change to the compacted topic member-balances through the upsert-kafka connector. The mobile app's points screen reads that topic, and the BigQuery sink connector on Kafka Connect merges each update into loyalty.member_balances. Every night the dbt Cloud job nightly_tiers rebuilds fct_member_tier, and a Hightouch sync sends tier changes to Braze at 22:00, which emails each member whose tier moved. The FlinkDeployment uses upgradeMode: savepoint, and allowNonRestoredState: true was set months ago to get past one restore, then left on.
| Time | System | What happens |
|---|---|---|
| Tue 14:00 | Kubernetes Operator + Flink | The team ships points-balance v4.2: the SQL now filters out test-store transactions and adds a last_store_id column. The operator takes a savepoint and stops v4.1 |
| 14:02 | Flink | v4.2 starts from the 14:00 savepoint. The edit gave the query a new physical plan, so the aggregate no longer matches its saved state, and with allowNonRestoredState on the restore skips that state instead of failing. The Kafka source resumes from the consumer group's committed offsets. The job is RUNNING and checkpoints succeed |
| 14:02 | Kafka | Every balance restarts from zero. A member with 6,400 points who buys milk now shows 12, upsert-kafka writes it to member-balances, and the app shows 12 |
| 14:05 | Kafka Connect + BigQuery | The sink and its tasks stay RUNNING and merge the small balances into loyalty.member_balances |
| 15:05 | Data Workers + BigQuery | The average balance of members active in the last hour, a metric the team records hourly with monitor_metrics, reads 9% of its baseline. run_quality_check finds member_id unique, no nulls and volume normal, so this is a value problem. Data Workers opens an incident |
| 15:15 | Data Workers + Flink REST API | Over an MCP server in front of the JobManager REST API, Data Workers reads that points-balance restored from the 14:00 savepoint at 14:02 and that the JobManager skipped savepoint state for one operator. With the v4.2 release recorded as the only recent change, diagnose_incident ranks the restore first. From the context graph (the dbt project and the team's notes on the Flink job and the Hightouch sync), blast_radius_analysis returns fct_member_tier and the 22:00 Braze sync, which would downgrade about 61,000 members who shopped after 14:02 |
| 15:20 | Slack | The approval request goes to the owner of points-balance, with the loyalty data owner copied |
| 15:40 | Spellbook | Three proposals: hold the Hightouch sync and the nightly_tiers schedule; a FlinkDeployment diff that redeploys v4.1 from the 14:00 savepoint (initialSavepointPath plus a new savepointRedeployNonce) with allowNonRestoredState: false; and v4.2 held until its state migration is planned. The owner approves all three |
| 15:50 | Hightouch + dbt Cloud | The owner pauses the sync to Braze and the nightly_tiers schedule |
| 16:05 | GitHub + Kubernetes Operator | She merges the diff. The operator redeploys v4.1 from the 14:00 savepoint, every balance intact as of 14:00 |
| 16:25 | Kafka + BigQuery | v4.1 replays points-events from the 14:00 offsets and catches up. upsert-kafka rewrites every balance touched since 14:00, the app is right again, and the sink merges the corrections into BigQuery |
| 16:45 | Data Workers + BigQuery | Data Workers verifies: the average-balance metric is back on baseline and run_quality_check passes. The owner's spot check of 200 members against the points ledger finds no differences |
| 17:10 | dbt Cloud | The owner starts the approved nightly_tiers run; it finishes at 17:24. Tier downgrades, recorded with monitor_metrics, sit at 1,120, in line with baseline. The receipt records cause, approvals, runs, checks and the undo (revert the diff; both savepoints are recorded) |
| 17:30 | Hightouch | The owner resumes the sync and the schedule. At 22:00 only real tier changes reach Braze |

Every part of Flink did what it was told. The savepoint documentation says --allowNonRestoredState skips state that cannot be mapped to the new program, and "improper usage of this feature could result in significant issues with the correctness of the application". The upgrade guide adds that for Table API and SQL, "any change to both the query and the Flink version could lead to state incompatibility". Catching the cost takes knowledge outside the job: that the balance lives only in state, that a nightly email reads it, and that the 14:00 savepoint is the clean way back.
| Job | What Flink does | What Data Workers does |
|---|---|---|
| The computation | Runs stateful jobs in Flink SQL, Table API or DataStream, with exactly-once checkpoints | Reads Flink over its REST API or an MCP server, and the sinks, lineage and tables around it natively |
| The state | Snapshots state in checkpoints and savepoints and restores it by the rules the job author sets | Ties each job to the tables its results feed, so a restore that changes a number becomes a data incident |
| The results | Writes to Kafka, Iceberg, Paimon or a database through its connectors | Checks what lands: nulls, uniqueness, freshness and volume on Snowflake and BigQuery, plus the team's key metrics against a baseline |
| The fix | Deploys, stops, rescales and restores whatever version the owner chooses | Proposes the holds, the deployment diff and the release plan to a named owner, then queues the approved downstream runs |
| The proof | Keeps job status, checkpoint history, exceptions 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 Flink just do this itself?
Because Flink is a processing engine, by design, and that focus is why it is trusted with state at scale. Its guarantees are precise: exactly-once state, consistent snapshots, restores that follow the job author's rules. Whether the right state was restored for the business is a question about the job's meaning, which Flink rightly leaves to the people who write the job.
Flink's AI work follows the same scope. Flink Agents (0.3, a preview release from June 2026; 0.3.1 in July added Flink 2.3 support) runs agents as operators inside a Flink job for event-driven work, with Agent Skills support, Mem0 memory and a YAML API. Judging whether another job's output is right in BigQuery or Braze sits outside that scope, as it should. The project publishes no official MCP server; its REST API, SQL Gateway and Kubernetes Operator are the integration surface.
Owning whether the data is right from a Flink job through Kafka, BigQuery, dbt Cloud and Hightouch is a different product: a context graph from topic to campaign, blast-radius scoping, approvals that name a person, a recorded undo and receipts an auditor can read. That is Data Workers. See 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 Flink jobs already there.

| Stage | Data Workers | Apache Flink | Why we scored it this way |
|---|---|---|---|
| Catalog & Context | 9 | 2 | Flink SQL catalogs name tables for a job; meaning, owners and consumers live elsewhere. Data Workers keeps one governed context graph from the topic through the job to the table and the dashboard. |
| Analytics & Insights | 8 | 4 | Flink SQL computes live aggregates and materialized tables; business definitions live in the warehouse and BI. Data Workers answers data questions from governed definitions with lineage behind every number. |
| Data Quality | 8 | 3 | Flink guarantees exactly-once state, not correct results; checks are code the team writes. Data Workers checks what jobs land, so a job that is RUNNING with wrong numbers still raises an incident. |
| Observability & Incidents | 8.5 | 4 | Job status, checkpoints, exceptions and metrics show the job's health. Data Workers diagnoses the data incident across the job, the sinks, the tables and the jobs downstream, proposes the fix and verifies it. |
| Pipelines & Ingestion | 8.5 | 9 | Flink's home stage: stateful stream processing with exactly-once checkpoints, savepoints, Flink SQL and Flink CDC. Data Workers plans the repair and queues the downstream reruns for the owner. |
| Schema & Migration | 8 | 4 | Savepoints carry state across upgrades; state compatibility is the job author's to manage. Data Workers traces a schema or plan change to every table and job that reads the result. |
| Governance & Access | 8.5 | 2 | Who may deploy or stop a job is set in Kubernetes or the managed service. Data Workers routes every data change to a named approver and records the decision. |
| Security & Privacy | 8 | 3 | TLS, Kerberos and the platform's IAM protect the cluster. Data Workers leaves a receipt on every data change. |
| Cost / FinOps | 8 | 3 | Parallelism, slots and the autoscaler size the job; spend shows up as compute. Data Workers attributes Snowflake credits to dbt models, totals BigQuery spend from the Jobs API and sets a scan budget per session. |
| MLOps & Models | 7.5 | 5 | Flink computes live features and runs Flink Agents (0.3) inside jobs. Data Workers keeps the data under models and agents healthy. |
How Flink and Data Workers work together

Data Workers connects to Flink over its REST API or an MCP server today. The JobManager REST API exposes jobs, checkpoint and savepoint history, exceptions and logs: what changed in this job, and when.
Everything around the job connects natively. Kafka Connect and Schema Registry: Data Workers reads each sink's connector and task status with get_connector_status, and registers a new schema version with register_stream_schema only after a named approval. Lineage: Flink 2 jobs can report to OpenLineage through Flink's native lineage API, including Flink SQL jobs, with the Kafka connector supported today; in Data Workers, the hop from topic to job to table goes into the context graph as a note your team records once. Then the landing side: BigQuery, Snowflake, Databricks, Iceberg catalogs, dbt Cloud, GitHub (pull-request review) and Slack, among 50+ connectors. Hightouch connects over its REST API or MCP server today.
By design, Data Workers does not deploy, stop, rescale or restore Flink jobs, trigger savepoints, change topics or offsets, or pause syncs. The owner runs those steps; Data Workers proposes them, scopes what they touch and verifies afterwards. Backpressure and consumer lag stay with your monitoring; Datadog and OpenTelemetry are native reads.
Setup 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-incidents -- "$(pwd)/start-agent.sh" dw-incidents
claude mcp add --scope user dw-quality -- "$(pwd)/start-agent.sh" dw-quality
claude mcp add --scope user dw-catalog -- "$(pwd)/start-agent.sh" dw-context-catalog
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-schemaList 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; run_quality_check (dw-quality) checks the BigQuery table; trace_cross_platform_lineage and blast_radius_analysis (dw-context-catalog) map what reads the balances; get_connector_status (dw-connectors) confirms the sink. For a schema change, check_compatibility and assess_impact (dw-schema) size it downstream.
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. See where does our data go.
One incident, L0 to L4, set per domain:

- •L0 manual. The 22:00 emails downgrade 61,000 members before anyone connects support calls to the release.
- •L1 observe. Data Workers flags the drop at 15:05 with the restore as the likely cause and the Braze sync in the blast radius. Nothing changes.
- •L2 propose. Data Workers proposes the holds and the rollback diff; nothing changes until the named owner approves.
- •L3 act reversibly. For a class with a clean record, Data Workers takes reversible steps in the domain you open, such as re-checking the balances table after the approved run.
- •L4 autonomous. Once the owner's rollback lands, Data Workers queues the downstream runs, re-checks the metrics and posts the receipt. Job changes and sync holds stay with the owner at every level.
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

- •One record for every team. The streaming team, the data owner and marketing see the same incident and receipt in Spellbook Data Catalog (in preview), linked from Slack.
- •Releases get a downstream check. A Flink SQL edit, restore or upgrade is weighed against the tables it writes, and the hours after it are watched against baseline.
- •Jobs that act on data get a guard. When an email sync or a payout reads a table that just went wrong, the hold is proposed before it runs.
- •Rollbacks have a plan. The savepoint, the replay and the reruns arrive as one approved sequence with a recorded undo.
For the wider streaming picture, see Flink vs Spark Structured Streaming for data teams and streaming agents for Kafka and Flink.
Keep Flink, or consolidate?
Keep Apache Flink 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 Flink team the answer is keep it: exactly-once state at this scale, Flink SQL, materialized tables that now evolve in place with Flink 2.3, and a choice of managed services are hard to replace. What teams consolidate is the tooling around the results: a separate data observability tool and runbooks that say "restore the last savepoint and eyeball the dashboard". If Kafka carries your events, see you're on Apache Kafka; on Confluent Cloud, you're on Confluent; on AWS, you're on Amazon Kinesis and MSK. If Hightouch sends your results to marketing tools, see you're on Hightouch.
If you are weighing building this yourself on a Flink MCP server and a coding agent, read build it ourselves with Claude Code and MCP servers.
The case for your CFO
The outcome. When a stream release or restore quietly writes wrong numbers into the tables behind balances, offers or customer emails, it is caught from the data in the hours after the release and corrected before customers act on it, with a record of how it was checked.
The risk story. At L1 agents only read. At L2 a named person approves every change; an unanswered request expires and escalates, never auto-grants. At L3 they act on reversible changes inside the domains you open. Flink jobs, deploys and restores 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. App screens, offers and emails now read the stream within seconds, so a wrong number reaches a member before a person sees it.
The first win. L1 on the Flink jobs whose results customers see: each output is checked against its baseline after every release.
What stays the same. Flink, your jobs, savepoints, Kafka, warehouse and reverse ETL. Zero migration. For the numbers, see the ROI of agentic data operations.
The sentence for upstairs: "Flink computes our real-time numbers; Data Workers makes sure what it writes is right, and gets it fixed with our approval when it isn't, before customers or campaigns act on it."
Getting started
Start with a pilot. Pick the Flink jobs whose results feed balances, offers or customer-facing numbers, connect the Flink REST API, Kafka Connect, lineage and the tables those jobs write, record the metrics that define "right", and run at L1 through a few releases. Then turn on L2 for one domain. 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 Flink? Over Flink's REST API or an MCP server today. The Kafka Connect sinks, Schema Registry and warehouse tables around your jobs connect natively.
Our job is RUNNING and every checkpoint succeeds. How can the results be wrong? They mean Flink is snapshotting state consistently, not that it is the right state: a restore can skip state when allowNonRestoredState is on, and a SQL edit can change the plan. Data Workers checks the results where they land and uses job status as one input to the diagnosis.
Should we turn off allowNonRestoredState? For stateful production jobs, keep it off by default and turn it on for a single deploy when you have dropped an operator on purpose. With it off, a restore that cannot map state fails loudly instead of starting clean.
Will Data Workers deploy, restore or stop our Flink jobs? No. It reads job status, restore history and the tables. Deploys, savepoints, restores, rescaling and offsets stay with the owner, through the Kubernetes Operator or your managed service. Data Workers proposes the step as a diff and verifies the data afterwards.
We run Flink as a managed service. Does this still work? Yes. Confluent Cloud for Apache Flink, Amazon Managed Service for Apache Flink and Ververica are read over their APIs or MCP servers; the sinks and tables around the jobs connect natively.
How do Flink Agents and Data Workers fit together? Flink Agents runs agents inside a Flink job. Data Workers owns whether what those jobs write downstream is right; an agent operator's output table gets the same checks, lineage and receipts as any other job's.
Sources
- •Apache Flink blog index (Flink 2.3.0 Jun 25, 2026; 1.20.5 Jun 8, 2026; Kubernetes Operator 1.16.0 Sep 15, 2026; Flink Agents 0.3.1 Jul 25, 2026; nav: Flink 2.3 stable, 1.20 LTS), https://flink.apache.org/posts/ (checked Oct 3, 2026)
- •Apache Flink, "Apache Flink 2.3.0 Release Announcement" (Jun 25, 2026; materialized table evolution and refresh control), https://flink.apache.org/2026/06/25/apache-flink-2.3.0-release-announcement/ (checked Oct 3, 2026)
- •Apache Flink 2.3 documentation, Savepoints (allowNonRestoredState; "Improper usage of this feature could result in significant issues with the correctness of the application"), https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/ops/state/savepoints/ (checked Oct 3, 2026)
- •Apache Flink 2.3 documentation, Upgrading Applications and Flink Versions (Table API & SQL: "any change to both the query and the Flink version could lead to state incompatibility"), https://nightlies.apache.org/flink/flink-docs-release-2.3/docs/ops/upgrading/ (checked Oct 3, 2026)
- •Apache Flink, "Apache Flink Kubernetes Operator 1.16.0 Release Announcement" (Sep 15, 2026), https://flink.apache.org/2026/09/15/apache-flink-kubernetes-operator-1.16.0-release-announcement/ (checked Oct 3, 2026)
- •Flink Kubernetes Operator 1.16 reference (JobSpec upgradeMode, allowNonRestoredState, initialSavepointPath, savepointRedeployNonce), https://nightlies.apache.org/flink/flink-kubernetes-operator-docs-release-1.16/docs/custom-resource/reference/ (checked Oct 3, 2026)
- •Apache Flink, "Apache Flink Agents 0.3.0 Release Announcement" (Jun 19, 2026; "a preview version"; Agent Skills, Mem0 memory, YAML API), https://flink.apache.org/2026/06/19/apache-flink-agents-0.3.0-release-announcement/ (checked Oct 3, 2026)
- •Apache Flink, "Apache Flink Agents 0.3.1 Release Announcement" (Jul 25, 2026; Flink 2.3 distribution support), https://flink.apache.org/2026/07/25/apache-flink-agents-0.3.1-release-announcement/ (checked Oct 3, 2026)
- •OpenLineage docs, Flink 2.x integration (native lineage API from FLIP-314, Flink SQL supported, "only the Kafka connector supports this functionality"), https://openlineage.io/docs/integrations/flink/flink2 (checked Oct 3, 2026)
- •AWS What's New, "Amazon Managed Service for Apache Flink now supports Apache Flink 2.3" (Jul 20, 2026), https://aws.amazon.com/about-aws/whats-new/2026/07/amazon-managed-service-flink-2-3/ (checked Oct 3, 2026)
- •Confluent documentation, Stream Processing with Confluent Cloud for Apache Flink, https://docs.confluent.io/cloud/current/flink/overview.html (checked Oct 3, 2026)
- •Ververica, enterprise stream processing platform (VERA engine, BYOC), https://www.ververica.com/ (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)