New to Kestra?
Use blueprints to kickstart your first workflows.
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.
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.
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.get_status (io.kestra.plugin.kafka.ConnectorGetStatus) reads the connector state plus each task state and trace.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.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.errors block pages on-call if a Connect API call fails.http://connect:8083.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') }}.
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.
SLACK_WEBHOOK_URL secret.connect_url and connector_name in variables.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.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.outputs.get_status.tasks: the task states and traces that started the run.outputs.final_status.connectorState: the state after a restart.connect_restarts_<connector_name>: the restart count for the current window.connect_paged_<connector_name>: present while on-call is paged and the trigger is muted.io.kestra.plugin.kafka.ConnectorList and a Loop over a subflow that holds this logic.UNASSIGNED handling with a second trigger on that targetState.