Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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>:type=connector-metrics,context=<snapshot|streaming>,server=<topic.prefix>`.
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=<n>`, and SQL Server
adds `database=<name>` 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`.
Expand Down Expand Up @@ -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-<connector>`, or the connector's `consumer.override.group.id`;
each member's `client.id` is `connector-consumer-<connector>-<task>`, 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
Expand Down
Loading