Airflow Summit 2026 is coming August 31 - September 2 in Austin, TX. Register now to secure your spot!

airflow.providers.apache.kafka.plugins.event_producer

Attributes

log

CONFIG_SECTION

SCHEMA_VERSION

Classes

DagRunListener

Publishes DagRun state-change event messages to Kafka.

TaskListener

Publishes TaskInstance state-change event messages to Kafka.

KafkaEventProducerPlugin

Publishes Airflow DagRun and TaskInstance event messages to a defined Kafka topic.

Module Contents

airflow.providers.apache.kafka.plugins.event_producer.log[source]
airflow.providers.apache.kafka.plugins.event_producer.CONFIG_SECTION = 'kafka_event_producer'[source]
airflow.providers.apache.kafka.plugins.event_producer.SCHEMA_VERSION = 1[source]
class airflow.providers.apache.kafka.plugins.event_producer.DagRunListener[source]

Publishes DagRun state-change event messages to Kafka.

on_dag_run_running(dag_run, msg)[source]
on_dag_run_success(dag_run, msg)[source]
on_dag_run_failed(dag_run, msg)[source]
class airflow.providers.apache.kafka.plugins.event_producer.TaskListener[source]

Publishes TaskInstance state-change event messages to Kafka.

on_task_instance_running(previous_state, task_instance)[source]
class airflow.providers.apache.kafka.plugins.event_producer.KafkaEventProducerPlugin[source]

Bases: airflow.providers.common.compat.sdk.AirflowPlugin

Publishes Airflow DagRun and TaskInstance event messages to a defined Kafka topic.

name = 'kafka_event_producer'[source]
listeners = [][source]

Was this entry helpful?