diff --git a/docs/superpowers/specs/2026-10-03-turbostats-kafka-connect-design.md b/docs/superpowers/specs/2026-10-03-turbostats-kafka-connect-design.md index 4119417d..2e37cc27 100644 --- a/docs/superpowers/specs/2026-10-03-turbostats-kafka-connect-design.md +++ b/docs/superpowers/specs/2026-10-03-turbostats-kafka-connect-design.md @@ -376,9 +376,39 @@ incremental snapshot's rows appear in both the snapshot context's `TotalNumberOfEventsSeen` and the streaming context's, so their sum counts each of those rows twice. -During a task restart, Debezium retries registering its metrics for up to a -minute while the old ones are still registered, and their values stand still -meanwhile. A receiver's stall rule must outlast that. +**Which Debezium MBeans belong to a task.** Debezium names its metrics +`debezium.:type=connector-metrics,context=,server=`. +Single-task connectors, such as Postgres and MySQL, add no task key, and their +MBeans belong to task 0. SQL Server and MongoDB add `task=`, and SQL Server +adds `database=` per database. The reporter matches `server` to the +connector's `topic.prefix` and `task` to the task. Several MBeans for one task, +one per database, are shards of one measurement: lag is their maximum, rows +their sum, and the backfill runs while any of them does. + +**A Debezium MBean registered before the task's last start belongs to the run +before it.** Restarted mid-snapshot, a task's new producer configured at once, +while the old run's MBeans stayed registered for 80 seconds. Debezium retried +registering the new ones every 5 seconds, twelve times, then logged "Failed to +register metrics MBean, metrics will not be available", and the new task ran +its whole life without a snapshot MBean. The reporter records when each +Debezium MBean registered, from the MBean server's notifications, and treats +one registered before the task's last start as stale. + +**A Debezium connector whose snapshot metrics are missing or stale sends +`backfill.state: unknown`.** Leaving `backfill` out would break the rule that +it is present from the first report; `none` would claim no snapshot ran, and +one did. The state vocabulary is open, so a receiver reads `unknown` as +unknown. `blocks_stream` means nothing under that state. Event lag and +`source_connected` are absent while the streaming MBean is missing or stale. + +A snapshot MBean registered fresh after a completed snapshot reads all zeros, +with neither `SnapshotRunning` nor `SnapshotCompleted` set: no backfill in this +run, so `state: none`. + +`MilliSecondsBehindSource` is -1 until the task processes an event, and event +lag is absent until then. `event_lag_observed_at` is the time of the last +event, now less `MilliSecondsSinceLastEvent`, because the reading is that +event's lag and does not move while the source is quiet. `backfill.state` comes from `SnapshotRunning`, `SnapshotPaused`, `SnapshotCompleted` and `SnapshotAborted`. A skipped snapshot is `none`. @@ -429,6 +459,18 @@ broker's message lag keeps growing. Read together, they catch a stalled sink. Consumer metrics replace the dots in a topic name with underscores, so a reporter reading them per topic maps the names back. +The broker knows which partitions each task holds. A sink connector's group +is `connect-`, or the connector's `consumer.override.group.id`; +each member's `client.id` is `connector-consumer--`, and its +assignment lists its partitions. A task's lag is the sum, over its +partitions, of the end offset less the committed offset. A failed sink task +leaves the group, so it has no assignment and sends no lag; its `state: +failed` already says what is wrong. + +The admin client uses the worker's own connection settings: the extension's +`configure` receives the worker's properties, `bootstrap.servers`, +`security.protocol` and the `sasl.*` and `ssl.*` settings among them. + `sink_rows_written` is required in v1, so a sink task must send a true number. Records read less records still in flight is the records Connect has finished: flushed, committed, or dropped by a transform. That last part makes