ConnectorStatusTrigger icon
ConnectorGetStatus icon
OutputValues icon
Switch icon
Set icon
ConnectorRestart icon
LoopUntil icon
If icon
SlackIncomingWebhook icon
Log icon

Self-heal Kafka Connect connectors with a restart budget and escalation

Auto-restart failed Kafka Connect connector tasks with Kestra, verify the recovery, then page on-call once with the stack trace when it keeps failing.

Categories
DataInfrastructure

Kafka Connect does not restart a failed task on its own. A sink that hits a bad record, a revoked credential or a full disk sits in FAILED until someone notices. This blueprint watches one connector, restarts only the failed tasks, waits for a real recovery, then counts restarts in a time window. When the connector keeps failing it stops restarting, pages on-call once with the failed task stack trace, then mutes itself until the window ends.

A connector usually stays RUNNING while one of its tasks is FAILED. The flow checks the connector state and every task state, so this common case is healed and not skipped.

How it works

  1. connector_failed (io.kestra.plugin.kafka.ConnectorStatusTrigger) polls the Connect REST API every minute. It starts an execution only when the connector or one of its tasks is FAILED, so a healthy connector creates no executions. The trigger does not start a new run while one is still in progress.
  2. get_status (io.kestra.plugin.kafka.ConnectorGetStatus) reads the connector state plus each task state and trace.
  3. read_restart_count (io.kestra.plugin.core.output.OutputValues) reads the restart count for the current window with the kv() function. A missing or expired key reads as 0.
  4. route (io.kestra.plugin.core.flow.Switch) picks one of three paths.
    • RESTART while under max_restarts: bump_restart_count (io.kestra.plugin.core.kv.Set) stores the new count with a ttl of flap_window. restart_connector (io.kestra.plugin.kafka.ConnectorRestart, includeTasks: true, onlyFailed: true) restarts the failed tasks. wait_for_recovery (io.kestra.plugin.core.flow.LoopUntil) re-reads the status every 10 seconds for up to 2 minutes. report_restart (io.kestra.plugin.core.flow.If) posts a recovered or still-unhealthy note to Slack.
    • ESCALATE once the budget is spent: page_on_call posts a Slack page with the first failed task stack trace. mute_until_window_ends (io.kestra.plugin.core.kv.Set) stores a connect_paged_<connector_name> key with a ttl of flap_window. The trigger when condition skips polling while that key exists, so on-call gets one page per window and not one per minute.
    • HEALTHY when nothing is FAILED any more: already_healthy logs it and the run ends.
  5. The flow-level errors block pages on-call if a Connect API call fails.

Prerequisites

  • A Kafka Connect worker or cluster reachable from Kestra over its REST API, for example http://connect:8083.
  • The connector to watch, already deployed on that worker.
  • A Slack incoming webhook.

Secrets

  • SLACK_WEBHOOK_URL: the Slack incoming webhook used for recovery notes and pages.

If your Connect REST API uses basic auth, add username and password to every io.kestra.plugin.kafka task and the trigger, for example {{ secret('KAFKA_CONNECT_USER') }} and {{ secret('KAFKA_CONNECT_PASSWORD') }}.

Variables

  • connect_url (string, default http://connect:8083): Connect REST API base URL.
  • connector_name (string, default orders_jdbc_sink): the connector to watch.
  • max_restarts (int, default 3): restarts allowed in one flap window before escalating.
  • flap_window (ISO 8601 duration, default PT30M): TTL of the restart counter. The budget resets after a full window with no restart.

These are flow variables and not inputs because trigger properties are rendered before an execution exists, so {{ inputs.* }} is not available in a trigger.

Quick start

  1. Add the SLACK_WEBHOOK_URL secret.
  2. Set connect_url and connector_name in variables.
  3. Save the flow. The trigger starts polling.
  4. To test it, break the connector. For a FileStreamSinkConnector, point file at a directory that does not exist. Then create the directory and watch the next run restart the task and report it recovered.
  5. After an escalation, fix the cause. To re-arm the watchdog before the window ends, delete the connect_paged_<connector_name> and connect_restarts_<connector_name> KV keys in the namespace. Otherwise both expire after flap_window and the next failed poll starts a new restart budget.

Expected outputs

  • outputs.get_status.tasks: the task states and traces that started the run.
  • outputs.final_status.connectorState: the state after a restart.
  • KV key connect_restarts_<connector_name>: the restart count for the current window.
  • KV key connect_paged_<connector_name>: present while on-call is paged and the trigger is muted.

How to extend

  • Watch every connector on a worker with io.kestra.plugin.kafka.ConnectorList and a Loop over a subflow that holds this logic.
  • Send the page to PagerDuty or Opsgenie in place of Slack.
  • Add UNASSIGNED handling with a second trigger on that targetState.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.