This is an automated email from the ASF dual-hosted git repository.

shahar1 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new 471df2dc220 Rename Kafka Listener to Kafka Event Producer and update 
the documentation (#70014)
471df2dc220 is described below

commit 471df2dc2201f90b0874c731dd6662e47a7b396d
Author: Christos Bisias <[email protected]>
AuthorDate: Sat Jul 18 15:06:44 2026 +0300

    Rename Kafka Listener to Kafka Event Producer and update the documentation 
(#70014)
---
 .../core_api/routes/public/test_plugins.py         |   2 +-
 providers/apache/kafka/docs/configurations-ref.rst |  42 +++--
 providers/apache/kafka/provider.yaml               |  16 +-
 providers/apache/kafka/pyproject.toml              |   2 +-
 .../providers/apache/kafka/get_provider_info.py    |  12 +-
 .../plugins/{listener.py => event_producer.py}     |  33 ++--
 .../{test_listener.py => test_event_producer.py}   |  14 +-
 .../{test_listener.py => test_event_producer.py}   | 176 ++++++++++-----------
 scripts/ci/docker-compose/integration-kafka.yml    |   8 +-
 9 files changed, 161 insertions(+), 144 deletions(-)

diff --git 
a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_plugins.py 
b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_plugins.py
index 1d9ba2175fa..38b2cc8fa9f 100644
--- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_plugins.py
+++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_plugins.py
@@ -50,7 +50,7 @@ class TestGetPlugins:
                     "edge_executor",
                     "hitl_review",
                     "hive",
-                    "kafka_listener",
+                    "kafka_event_producer",
                     "plugin-a",
                     "plugin-b",
                     "plugin-c",
diff --git a/providers/apache/kafka/docs/configurations-ref.rst 
b/providers/apache/kafka/docs/configurations-ref.rst
index 80d0d373f8f..de52624721a 100644
--- a/providers/apache/kafka/docs/configurations-ref.rst
+++ b/providers/apache/kafka/docs/configurations-ref.rst
@@ -24,27 +24,41 @@
 Highlighted configurations
 ===========================
 
-The ``[kafka_listener]`` section configures the ``KafkaListenerPlugin``,
+The ``[kafka_event_producer]`` section configures the 
``KafkaEventProducerPlugin``,
 which publishes Airflow DagRun and TaskInstance state-change events to a
 Kafka topic. DagRun and TaskInstance events are separated and enabled by
 distinct flags. Both event-type flags default to ``False``.
 
-.. _configuration_kafka_listener_activation:kafka:
+.. _configuration_kafka_event_producer_use_cases:kafka:
 
-Activating the listener
------------------------
+Common use-cases
+----------------
+
+  * Consume the Kafka events by an external observability or analytics tool 
and gather info about the state
+    of multiple Airflow instances without polling their metadata DBs.
+  * Based on the state of a DagRun, trigger a downstream external 
system/pipeline (notifications, alerting,
+    cross-team handoffs) without direct interaction with Airflow.
+  * Coordinate Dags across multiple Airflow instances over a shared Kafka 
service.
+
+    * For example, team_A with Airflow instance_A has a deferred task which is 
triggered
+      when a task from team_B with Airflow instance_B finishes.
+
+.. _configuration_kafka_event_producer_activation:kafka:
+
+Activating the plugin
+---------------------
 
 To enable event publishing you need to
 
   * enable at least one event-type flag
-  * point the listener at an Airflow Kafka connection via ``kafka_config_id``
+  * point the plugin at an Airflow Kafka connection via ``kafka_config_id``
     (defaults to ``kafka_default``) that carries the broker address and any
     other confluent-kafka client options on its extras
   * have a pre-existing kafka topic
 
 .. code-block:: ini
 
-    [kafka_listener]
+    [kafka_event_producer]
     dag_run_events_enabled = True
     task_instance_events_enabled = True
     kafka_config_id = kafka_events
@@ -68,20 +82,20 @@ Environment-variable equivalents:
 
 .. code-block:: ini
 
-    AIRFLOW__KAFKA_LISTENER__DAG_RUN_EVENTS_ENABLED=True
-    AIRFLOW__KAFKA_LISTENER__TASK_INSTANCE_EVENTS_ENABLED=True
-    AIRFLOW__KAFKA_LISTENER__KAFKA_CONFIG_ID=kafka_events
-    AIRFLOW__KAFKA_LISTENER__TOPIC=airflow.events
+    AIRFLOW__KAFKA_EVENT_PRODUCER__DAG_RUN_EVENTS_ENABLED=True
+    AIRFLOW__KAFKA_EVENT_PRODUCER__TASK_INSTANCE_EVENTS_ENABLED=True
+    AIRFLOW__KAFKA_EVENT_PRODUCER__KAFKA_CONFIG_ID=kafka_events
+    AIRFLOW__KAFKA_EVENT_PRODUCER__TOPIC=airflow.events
 
 The two event flags are independent, users can opt-in to get only DagRun
 event messages or only TaskInstance event messages or both.
 
 The topic must already exist on the broker, it's not auto-created. On a missing
-topic, broker connection failure, or any other producer init error, the 
listener
+topic, broker connection failure, or any other producer init error, the plugin
 doesn't fail, instead it logs a warning and retries the init after 
``topic_check_retry_interval``
-seconds (default ``60``). Once the topic is created on the broker the listener 
will pick it up.
+seconds (default ``60``). Once the topic is created on the broker the plugin 
will pick it up.
 
-.. _configuration_kafka_listener_filtering:kafka:
+.. _configuration_kafka_event_producer_filtering:kafka:
 
 Filtering events
 ----------------
@@ -92,7 +106,7 @@ comma-separated list of ``fnmatch`` glob patterns; an empty 
list means
 
 .. code-block:: ini
 
-    [kafka_listener]
+    [kafka_event_producer]
     dag_run_dag_id_allowlist = sales_*,marketing_*
     dag_run_dag_id_denylist = sales_internal_*
 
diff --git a/providers/apache/kafka/provider.yaml 
b/providers/apache/kafka/provider.yaml
index 8e23df57309..a0ae3a8aebd 100644
--- a/providers/apache/kafka/provider.yaml
+++ b/providers/apache/kafka/provider.yaml
@@ -131,14 +131,14 @@ queues:
   - airflow.providers.apache.kafka.queues.kafka.KafkaMessageQueueProvider
 
 plugins:
-  - name: kafka_listener
-    plugin-class: 
airflow.providers.apache.kafka.plugins.listener.KafkaListenerPlugin
+  - name: kafka_event_producer
+    plugin-class: 
airflow.providers.apache.kafka.plugins.event_producer.KafkaEventProducerPlugin
 
 config:
-  kafka_listener:
+  kafka_event_producer:
     description: |
-      Settings for the Kafka listener that publishes Airflow DagRun and
-      TaskInstance state-change events to a Kafka topic.
+      Settings for the Kafka event producer plugin that publishes Airflow
+      DagRun and TaskInstance state-change events to a Kafka topic.
     options:
       dag_run_events_enabled:
         description: |
@@ -161,7 +161,7 @@ config:
         default: "False"
       kafka_config_id:
         description: |
-          Airflow connection used to build the listener's Kafka producer.
+          Airflow connection used to build the plugin's Kafka producer.
           When unset, the producer hook falls back to its default
           connection (``kafka_default``).
         version_added: 1.14.1
@@ -170,8 +170,8 @@ config:
         default: ""
       topic:
         description: |
-          Topic the listener publishes events to. The topic must already
-          exist on the broker; the listener will not auto-create it.
+          Topic the plugin publishes events to. The topic must already
+          exist on the broker; the plugin will not auto-create it.
         version_added: 1.14.1
         type: string
         example: ~
diff --git a/providers/apache/kafka/pyproject.toml 
b/providers/apache/kafka/pyproject.toml
index 3c5ead9a4b8..b2e3a1fa598 100644
--- a/providers/apache/kafka/pyproject.toml
+++ b/providers/apache/kafka/pyproject.toml
@@ -131,7 +131,7 @@ apache-airflow-providers-standard = {workspace = true}
 provider_info = 
"airflow.providers.apache.kafka.get_provider_info:get_provider_info"
 
 [project.entry-points."airflow.plugins"]
-kafka_listener = 
"airflow.providers.apache.kafka.plugins.listener:KafkaListenerPlugin"
+kafka_event_producer = 
"airflow.providers.apache.kafka.plugins.event_producer:KafkaEventProducerPlugin"
 
 [tool.flit.module]
 name = "airflow.providers.apache.kafka"
diff --git 
a/providers/apache/kafka/src/airflow/providers/apache/kafka/get_provider_info.py
 
b/providers/apache/kafka/src/airflow/providers/apache/kafka/get_provider_info.py
index d8c8510cb5d..31162c55425 100644
--- 
a/providers/apache/kafka/src/airflow/providers/apache/kafka/get_provider_info.py
+++ 
b/providers/apache/kafka/src/airflow/providers/apache/kafka/get_provider_info.py
@@ -103,13 +103,13 @@ def get_provider_info():
         "queues": 
["airflow.providers.apache.kafka.queues.kafka.KafkaMessageQueueProvider"],
         "plugins": [
             {
-                "name": "kafka_listener",
-                "plugin-class": 
"airflow.providers.apache.kafka.plugins.listener.KafkaListenerPlugin",
+                "name": "kafka_event_producer",
+                "plugin-class": 
"airflow.providers.apache.kafka.plugins.event_producer.KafkaEventProducerPlugin",
             }
         ],
         "config": {
-            "kafka_listener": {
-                "description": "Settings for the Kafka listener that publishes 
Airflow DagRun and\nTaskInstance state-change events to a Kafka topic.\n",
+            "kafka_event_producer": {
+                "description": "Settings for the Kafka event producer plugin 
that publishes Airflow\nDagRun and TaskInstance state-change events to a Kafka 
topic.\n",
                 "options": {
                     "dag_run_events_enabled": {
                         "description": "Publish DagRun state-change events 
(``dag_run.running``,\n``dag_run.success``, ``dag_run.failed``). When False 
the\nDagRun listener is not registered.\n",
@@ -126,14 +126,14 @@ def get_provider_info():
                         "default": "False",
                     },
                     "kafka_config_id": {
-                        "description": "Airflow connection used to build the 
listener's Kafka producer.\nWhen unset, the producer hook falls back to its 
default\nconnection (``kafka_default``).\n",
+                        "description": "Airflow connection used to build the 
plugin's Kafka producer.\nWhen unset, the producer hook falls back to its 
default\nconnection (``kafka_default``).\n",
                         "version_added": "1.14.1",
                         "type": "string",
                         "example": "kafka_default",
                         "default": "",
                     },
                     "topic": {
-                        "description": "Topic the listener publishes events 
to. The topic must already\nexist on the broker; the listener will not 
auto-create it.\n",
+                        "description": "Topic the plugin publishes events to. 
The topic must already\nexist on the broker; the plugin will not auto-create 
it.\n",
                         "version_added": "1.14.1",
                         "type": "string",
                         "example": None,
diff --git 
a/providers/apache/kafka/src/airflow/providers/apache/kafka/plugins/listener.py 
b/providers/apache/kafka/src/airflow/providers/apache/kafka/plugins/event_producer.py
similarity index 92%
rename from 
providers/apache/kafka/src/airflow/providers/apache/kafka/plugins/listener.py
rename to 
providers/apache/kafka/src/airflow/providers/apache/kafka/plugins/event_producer.py
index 342b37139a0..20f6f1f6e44 100644
--- 
a/providers/apache/kafka/src/airflow/providers/apache/kafka/plugins/listener.py
+++ 
b/providers/apache/kafka/src/airflow/providers/apache/kafka/plugins/event_producer.py
@@ -40,7 +40,7 @@ if TYPE_CHECKING:
 
 log = logging.getLogger(__name__)
 
-CONFIG_SECTION = "kafka_listener"
+CONFIG_SECTION = "kafka_event_producer"
 SCHEMA_VERSION = 1
 
 
@@ -152,7 +152,7 @@ def _task_instance_event_allowed(dag_id: str, task_id: str) 
-> bool:
     )
 
 
-# Listener-owned producer and topic state for the lifetime of the process.
+# Plugin-owned producer and topic state for the lifetime of the process.
 # The producer is built once via KafkaProducerHook and kept for the process 
lifetime;
 # the topic flags track whether the topic exists on the broker and whether 
we're
 # currently in a back-off window after a failed topic check.
@@ -179,7 +179,7 @@ os.register_at_fork(after_in_child=_reset_state_after_fork)
 
 
 def _get_producer() -> Producer | None:
-    """Build (once) and return the listener's producer, or ``None`` on init 
failure."""
+    """Build (once) and return the plugin's Kafka producer, or ``None`` on 
init failure."""
     global _producer
     if _producer is not None:
         return _producer
@@ -194,7 +194,7 @@ def _get_producer() -> Producer | None:
             hook_kwargs["kafka_config_id"] = _get_kafka_config_id()
         producer = KafkaProducerHook(**hook_kwargs).get_producer()
     except Exception as exc:
-        log.warning("Kafka listener: failed to initialize producer (%s).", exc)
+        log.warning("Kafka event producer: failed to initialize producer 
(%s).", exc)
         return None
 
     atexit.register(_flush_producer_at_exit, producer)
@@ -224,7 +224,7 @@ def _check_topic_exists() -> bool:
         topics = 
producer.list_topics(timeout=_get_topic_check_timeout()).topics
     except Exception as exc:
         log.warning(
-            "Kafka listener: topic check failed (%s). Will retry after %ds.",
+            "Kafka event producer: topic check failed (%s). Will retry after 
%ds.",
             exc,
             _get_topic_check_retry_interval(),
         )
@@ -233,7 +233,7 @@ def _check_topic_exists() -> bool:
 
     if _get_topic() not in topics:
         log.warning(
-            "Kafka listener: topic %r not found on the broker. Will retry 
after %ds. "
+            "Kafka event producer: topic %r not found on the broker. Will 
retry after %ds. "
             "Create the topic on the broker to enable publishing.",
             _get_topic(),
             _get_topic_check_retry_interval(),
@@ -242,7 +242,7 @@ def _check_topic_exists() -> bool:
         return False
 
     log.info(
-        "Kafka listener attached: pid=%s source=%r topic=%r",
+        "Kafka event producer attached: pid=%s source=%r topic=%r",
         os.getpid(),
         _get_source(),
         _get_topic(),
@@ -255,16 +255,16 @@ def _flush_producer_at_exit(producer: Producer) -> None:
     try:
         producer.flush(5)
     except Exception:
-        log.debug("Kafka listener: error flushing producer on exit", 
exc_info=True)
+        log.debug("Kafka event producer: error flushing producer on exit", 
exc_info=True)
 
 
 def _on_delivery(err, _msg) -> None:
     if err is None:
         return
-    log.warning("Kafka listener: delivery failed: %s", err)
+    log.warning("Kafka event producer: delivery failed: %s", err)
     # If the broker or local metadata says the topic is gone, the confirmation 
cached at
     # attach-time is stale. Flip it back so the next _check_topic_exists 
re-verifies with
-    # the broker, gated by the same cooldown used elsewhere. Lets the listener 
auto-recover
+    # the broker, gated by the same cooldown used elsewhere. Lets the plugin 
auto-recover
     # if the topic gets recreated instead of requiring a component restart.
     from confluent_kafka import KafkaError
 
@@ -291,7 +291,7 @@ def _produce_dr_message(event: str, dag_run: DagRun, msg: 
str) -> None:
             _get_dr_payload(dag_run, msg),
         )
     except Exception:
-        log.exception("Kafka listener: %s failed", event)
+        log.exception("Kafka event producer: %s failed", event)
 
 
 def _get_dr_payload(dag_run, msg) -> dict[str, Any]:
@@ -325,7 +325,7 @@ def _produce_ti_message(
             _get_ti_payload(task_instance, previous_state, error=error),
         )
     except Exception:
-        log.exception("Kafka listener: %s failed", event)
+        log.exception("Kafka event producer: %s failed", event)
 
 
 def _get_ti_payload(ti, previous_state, error=None) -> dict[str, Any]:
@@ -373,9 +373,12 @@ def _produce_message(event: str, dag_id: str, run_id: str, 
payload: dict[str, An
         )
         producer.poll(0)
     except Exception as ex:
-        log.warning("Kafka listener: failed to enqueue %s for %s/%s: %s", 
event, dag_id, run_id, ex)
+        log.warning("Kafka event producer: failed to enqueue %s for %s/%s: 
%s", event, dag_id, run_id, ex)
 
 
+# DagRunListener / TaskListener are pluggy hookimpl classes that plug into 
Airflow's
+# listener API (``AirflowPlugin.listeners``). These classes listen for Airflow 
events
+# and then produce a message for every event.
 class DagRunListener:
     """Publishes DagRun state-change event messages to Kafka."""
 
@@ -463,8 +466,8 @@ def _get_enabled_listeners() -> list[object]:
     return listeners
 
 
-class KafkaListenerPlugin(AirflowPlugin):
+class KafkaEventProducerPlugin(AirflowPlugin):
     """Publishes Airflow DagRun and TaskInstance event messages to a defined 
Kafka topic."""
 
-    name = "kafka_listener"
+    name = "kafka_event_producer"
     listeners = _get_enabled_listeners()
diff --git 
a/providers/apache/kafka/tests/integration/apache/kafka/plugins/test_listener.py
 
b/providers/apache/kafka/tests/integration/apache/kafka/plugins/test_event_producer.py
similarity index 92%
rename from 
providers/apache/kafka/tests/integration/apache/kafka/plugins/test_listener.py
rename to 
providers/apache/kafka/tests/integration/apache/kafka/plugins/test_event_producer.py
index ce900f79b87..81cbc4ca931 100644
--- 
a/providers/apache/kafka/tests/integration/apache/kafka/plugins/test_listener.py
+++ 
b/providers/apache/kafka/tests/integration/apache/kafka/plugins/test_event_producer.py
@@ -42,7 +42,7 @@ log = logging.getLogger(__name__)
 
 client_config = {
     "bootstrap.servers": "broker:29092",
-    "group.id": "kafka-listener-integration-test",
+    "group.id": "kafka-event-producer-integration-test",
     "enable.auto.commit": False,
     "auto.offset.reset": "earliest",
 }
@@ -59,7 +59,7 @@ def _wait_for_assignment(consumer, timeout: float = 10.0) -> 
None:
 
 @pytest.mark.integration("kafka")
 @pytest.mark.backend("postgres")
-class TestEventListener:
+class TestEventProducer:
     test_dir = os.path.dirname(os.path.abspath(__file__))
     dag_folder = os.path.join(test_dir, "dags")
 
@@ -83,12 +83,12 @@ class TestEventListener:
         os.environ["AIRFLOW__CORE__LOAD_EXAMPLES"] = "False"
         os.environ["AIRFLOW__CORE__UNIT_TEST_MODE"] = "False"
 
-        os.environ["AIRFLOW__KAFKA_LISTENER__DAG_RUN_EVENTS_ENABLED"] = "True"
-        os.environ["AIRFLOW__KAFKA_LISTENER__TASK_INSTANCE_EVENTS_ENABLED"] = 
"True"
-        os.environ["AIRFLOW__KAFKA_LISTENER__TOPIC"] = cls.TOPIC
-        os.environ["AIRFLOW__KAFKA_LISTENER__SOURCE"] = "dev-breeze"
+        os.environ["AIRFLOW__KAFKA_EVENT_PRODUCER__DAG_RUN_EVENTS_ENABLED"] = 
"True"
+        
os.environ["AIRFLOW__KAFKA_EVENT_PRODUCER__TASK_INSTANCE_EVENTS_ENABLED"] = 
"True"
+        os.environ["AIRFLOW__KAFKA_EVENT_PRODUCER__TOPIC"] = cls.TOPIC
+        os.environ["AIRFLOW__KAFKA_EVENT_PRODUCER__SOURCE"] = "dev-breeze"
 
-        # Shared Kafka connection: used by the listener producer, the topic
+        # Shared Kafka connection: used by the event producer plugin, the topic
         # creation/deletion, and the consumer.
         kafka_default_conn = Connection(
             conn_id=cls.KAFKA_CONFIG_ID,
diff --git 
a/providers/apache/kafka/tests/unit/apache/kafka/plugins/test_listener.py 
b/providers/apache/kafka/tests/unit/apache/kafka/plugins/test_event_producer.py
similarity index 72%
rename from 
providers/apache/kafka/tests/unit/apache/kafka/plugins/test_listener.py
rename to 
providers/apache/kafka/tests/unit/apache/kafka/plugins/test_event_producer.py
index 38bdc056637..8e4ed0fcf4e 100644
--- a/providers/apache/kafka/tests/unit/apache/kafka/plugins/test_listener.py
+++ 
b/providers/apache/kafka/tests/unit/apache/kafka/plugins/test_event_producer.py
@@ -24,8 +24,8 @@ import pytest
 from confluent_kafka import KafkaError
 
 from airflow.models import Connection
-from airflow.providers.apache.kafka.plugins import listener
-from airflow.providers.apache.kafka.plugins.listener import DagRunListener, 
TaskListener
+from airflow.providers.apache.kafka.plugins import event_producer
+from airflow.providers.apache.kafka.plugins.event_producer import 
DagRunListener, TaskListener
 
 from tests_common.test_utils.config import conf_vars
 from tests_common.test_utils.version_compat import AIRFLOW_V_3_0_PLUS
@@ -48,7 +48,7 @@ _PRODUCER_HOOK_CLS = 
"airflow.providers.apache.kafka.hooks.produce.KafkaProducer
 
 @pytest.fixture(autouse=True)
 def _default_kafka_conn(create_connection_without_db):
-    # The listener builds its producer through ``KafkaProducerHook``, which 
requires
+    # The plugin builds its producer through ``KafkaProducerHook``, which 
requires
     # a resolvable ``kafka_default`` connection.
     create_connection_without_db(
         Connection(
@@ -63,24 +63,24 @@ def _default_kafka_conn(create_connection_without_db):
 def _clear_cached_values():
     """Invalidates caches before each test run."""
     # Reset config cache.
-    listener._dag_run_events_enabled.cache_clear()
-    listener._task_instance_events_enabled.cache_clear()
-    listener._get_topic.cache_clear()
-    listener._get_kafka_config_id.cache_clear()
-    listener._get_source.cache_clear()
-    listener._get_dag_run_dag_id_allowlist.cache_clear()
-    listener._get_dag_run_dag_id_denylist.cache_clear()
-    listener._get_task_instance_dag_id_allowlist.cache_clear()
-    listener._get_task_instance_dag_id_denylist.cache_clear()
-    listener._get_task_instance_task_id_allowlist.cache_clear()
-    listener._get_task_instance_task_id_denylist.cache_clear()
-    listener._get_topic_check_timeout.cache_clear()
-    listener._get_topic_check_retry_interval.cache_clear()
+    event_producer._dag_run_events_enabled.cache_clear()
+    event_producer._task_instance_events_enabled.cache_clear()
+    event_producer._get_topic.cache_clear()
+    event_producer._get_kafka_config_id.cache_clear()
+    event_producer._get_source.cache_clear()
+    event_producer._get_dag_run_dag_id_allowlist.cache_clear()
+    event_producer._get_dag_run_dag_id_denylist.cache_clear()
+    event_producer._get_task_instance_dag_id_allowlist.cache_clear()
+    event_producer._get_task_instance_dag_id_denylist.cache_clear()
+    event_producer._get_task_instance_task_id_allowlist.cache_clear()
+    event_producer._get_task_instance_task_id_denylist.cache_clear()
+    event_producer._get_topic_check_timeout.cache_clear()
+    event_producer._get_topic_check_retry_interval.cache_clear()
 
     # Reset the module-level producer and topic state.
-    listener._producer = None
-    listener._topic_exists = False
-    listener._topic_check_retry_after = 0.0
+    event_producer._producer = None
+    event_producer._topic_exists = False
+    event_producer._topic_check_retry_after = 0.0
 
 
 @pytest.fixture
@@ -126,7 +126,7 @@ def _assert_common_message_fields(kafka_producer_mock, 
expected_event: str) -> d
     assert kwargs["key"] == f"{_DAG_ID}/{_DAG_RUN_ID}".encode()
 
     body = json.loads(kwargs["value"].decode("utf-8"))
-    assert body["schema_version"] == listener.SCHEMA_VERSION
+    assert body["schema_version"] == event_producer.SCHEMA_VERSION
     assert body["source"] == _SOURCE
     assert body["event"] == expected_event
     assert body["dag_id"] == _DAG_ID
@@ -140,22 +140,22 @@ def _assert_common_message_fields(kafka_producer_mock, 
expected_event: str) -> d
     [
         pytest.param(
             {
-                ("kafka_listener", "dag_run_events_enabled"): "True",
-                ("kafka_listener", "task_instance_events_enabled"): "True",
+                ("kafka_event_producer", "dag_run_events_enabled"): "True",
+                ("kafka_event_producer", "task_instance_events_enabled"): 
"True",
             },
             [DagRunListener, TaskListener],
             id="both_listeners_enabled",
         ),
         pytest.param(
             {
-                ("kafka_listener", "dag_run_events_enabled"): "True",
+                ("kafka_event_producer", "dag_run_events_enabled"): "True",
             },
             [DagRunListener],
             id="only_dr_events",
         ),
         pytest.param(
             {
-                ("kafka_listener", "task_instance_events_enabled"): "True",
+                ("kafka_event_producer", "task_instance_events_enabled"): 
"True",
             },
             [TaskListener],
             id="only_ti_events",
@@ -169,7 +169,7 @@ def _assert_common_message_fields(kafka_producer_mock, 
expected_event: str) -> d
 )
 def test_get_enabled_listeners(configs, expected_listener_classes):
     with conf_vars(configs):
-        registered_listeners = listener._get_enabled_listeners()
+        registered_listeners = event_producer._get_enabled_listeners()
         # Both lists should have the same size.
         assert len(registered_listeners) == len(expected_listener_classes)
 
@@ -187,9 +187,9 @@ def test_get_enabled_listeners(configs, 
expected_listener_classes):
 )
 @conf_vars(
     {
-        ("kafka_listener", "dag_run_events_enabled"): "True",
-        ("kafka_listener", "topic"): _KAFKA_TOPIC,
-        ("kafka_listener", "source"): _SOURCE,
+        ("kafka_event_producer", "dag_run_events_enabled"): "True",
+        ("kafka_event_producer", "topic"): _KAFKA_TOPIC,
+        ("kafka_event_producer", "source"): _SOURCE,
     }
 )
 def test_produce_dr_message(
@@ -224,9 +224,9 @@ def test_produce_dr_message(
 )
 @conf_vars(
     {
-        ("kafka_listener", "task_instance_events_enabled"): "True",
-        ("kafka_listener", "topic"): _KAFKA_TOPIC,
-        ("kafka_listener", "source"): _SOURCE,
+        ("kafka_event_producer", "task_instance_events_enabled"): "True",
+        ("kafka_event_producer", "topic"): _KAFKA_TOPIC,
+        ("kafka_event_producer", "source"): _SOURCE,
     }
 )
 def test_produce_ti_message(
@@ -261,15 +261,15 @@ class TestGetProducer:
 
     @pytest.fixture(autouse=True)
     def _producer_conf(self):
-        with conf_vars({("kafka_listener", "topic"): _KAFKA_TOPIC}):
+        with conf_vars({("kafka_event_producer", "topic"): _KAFKA_TOPIC}):
             yield
 
     def test_producer_cached_on_success(self):
         kafka_producer_mock = MagicMock()
 
         with patch(_PRODUCER_CLS, return_value=kafka_producer_mock) as 
producer_cls_mock:
-            producer1 = listener._get_producer()
-            producer2 = listener._get_producer()
+            producer1 = event_producer._get_producer()
+            producer2 = event_producer._get_producer()
 
         assert producer1 is kafka_producer_mock
         assert producer2 is producer1
@@ -277,10 +277,10 @@ class TestGetProducer:
 
     def test_init_failure_warns_and_returns_none(self, monkeypatch):
         log_warning_mock = MagicMock()
-        monkeypatch.setattr(listener.log, "warning", log_warning_mock)
+        monkeypatch.setattr(event_producer.log, "warning", log_warning_mock)
 
         with patch(_PRODUCER_CLS, side_effect=ValueError("boom")):
-            assert listener._get_producer() is None
+            assert event_producer._get_producer() is None
 
         log_warning_mock.assert_called_once()
         warning_format, exc_arg = log_warning_mock.call_args.args
@@ -293,30 +293,30 @@ class TestGetProducer:
         kafka_producer_mock = MagicMock()
 
         register_mock = MagicMock()
-        monkeypatch.setattr(listener.atexit, "register", register_mock)
+        monkeypatch.setattr(event_producer.atexit, "register", register_mock)
 
         with patch(_PRODUCER_CLS, return_value=kafka_producer_mock):
-            listener._get_producer()
+            event_producer._get_producer()
 
-        
register_mock.assert_called_once_with(listener._flush_producer_at_exit, 
kafka_producer_mock)
+        
register_mock.assert_called_once_with(event_producer._flush_producer_at_exit, 
kafka_producer_mock)
 
-    def test_hook_built_from_listener_config(self):
-        """The listener section's ``kafka_config_id`` is forwarded to the 
hook."""
+    def test_hook_built_from_event_producer_config(self):
+        """The plugin section's ``kafka_config_id`` is forwarded to the 
hook."""
         hook_mock = MagicMock()
 
-        with conf_vars({("kafka_listener", "kafka_config_id"): 
"kafka_events"}):
+        with conf_vars({("kafka_event_producer", "kafka_config_id"): 
"kafka_events"}):
             with patch(_PRODUCER_HOOK_CLS, return_value=hook_mock) as 
hook_cls_mock:
-                producer = listener._get_producer()
+                producer = event_producer._get_producer()
 
         hook_cls_mock.assert_called_once_with(kafka_config_id="kafka_events")
         assert producer is hook_mock.get_producer.return_value
 
     def test_hook_defaults_when_options_unset(self):
-        """Without listener options the hook falls back to its own defaults."""
+        """Without plugin options the hook falls back to its own defaults."""
         hook_mock = MagicMock()
 
         with patch(_PRODUCER_HOOK_CLS, return_value=hook_mock) as 
hook_cls_mock:
-            listener._get_producer()
+            event_producer._get_producer()
 
         # Unset options must not be passed to the hook at all (e.g. as ""), so 
the
         # hook's own defaults (the "kafka_default" connection) apply.
@@ -328,7 +328,7 @@ class TestCheckTopicExists:
 
     @pytest.fixture(autouse=True)
     def _topic_conf(self):
-        with conf_vars({("kafka_listener", "topic"): _KAFKA_TOPIC}):
+        with conf_vars({("kafka_event_producer", "topic"): _KAFKA_TOPIC}):
             yield
 
     def test_topic_exists_kept_for_process_lifetime(self):
@@ -336,9 +336,9 @@ class TestCheckTopicExists:
         
kafka_producer_mock.list_topics.return_value.topics.__contains__.return_value = 
True
 
         with patch(_PRODUCER_CLS, return_value=kafka_producer_mock):
-            assert listener._check_topic_exists() is True
+            assert event_producer._check_topic_exists() is True
             # Once confirmed, the broker is not queried again.
-            assert listener._check_topic_exists() is True
+            assert event_producer._check_topic_exists() is True
 
         kafka_producer_mock.list_topics.assert_called_once()
 
@@ -347,29 +347,29 @@ class TestCheckTopicExists:
         
kafka_producer_mock.list_topics.return_value.topics.__contains__.return_value = 
False
 
         log_warning_mock = MagicMock()
-        monkeypatch.setattr(listener.log, "warning", log_warning_mock)
+        monkeypatch.setattr(event_producer.log, "warning", log_warning_mock)
 
         with patch(_PRODUCER_CLS, return_value=kafka_producer_mock):
-            assert listener._check_topic_exists() is False
+            assert event_producer._check_topic_exists() is False
 
         log_warning_mock.assert_called_once()
         # Format string + the two formatting args are passed positionally to 
log.warning.
         warning_format, topic_arg, retry_interval_arg = 
log_warning_mock.call_args.args
         assert "topic %r not found on the broker" in warning_format
         assert topic_arg == _KAFKA_TOPIC
-        assert retry_interval_arg == listener._get_topic_check_retry_interval()
+        assert retry_interval_arg == 
event_producer._get_topic_check_retry_interval()
 
     def test_list_topics_failure_warns_and_sets_cooldown(self, monkeypatch):
         kafka_producer_mock = MagicMock()
         kafka_producer_mock.list_topics.side_effect = ValueError("broker down")
 
         log_warning_mock = MagicMock()
-        monkeypatch.setattr(listener.log, "warning", log_warning_mock)
+        monkeypatch.setattr(event_producer.log, "warning", log_warning_mock)
 
         with patch(_PRODUCER_CLS, return_value=kafka_producer_mock):
-            assert listener._check_topic_exists() is False
+            assert event_producer._check_topic_exists() is False
             # Within the cooldown window the check is not retried.
-            assert listener._check_topic_exists() is False
+            assert event_producer._check_topic_exists() is False
 
         kafka_producer_mock.list_topics.assert_called_once()
         log_warning_mock.assert_called_once()
@@ -377,7 +377,7 @@ class TestCheckTopicExists:
 
     def test_retried_after_failure_cooldown(self, monkeypatch):
         time_now = [1000.0]
-        monkeypatch.setattr(listener.time, "monotonic", lambda: time_now[0])
+        monkeypatch.setattr(event_producer.time, "monotonic", lambda: 
time_now[0])
 
         kafka_producer_mock = MagicMock()
         # The topic doesn't exist yet.
@@ -385,36 +385,36 @@ class TestCheckTopicExists:
 
         with patch(_PRODUCER_CLS, return_value=kafka_producer_mock):
             # First check fails: the topic is missing.
-            assert listener._check_topic_exists() is False
+            assert event_producer._check_topic_exists() is False
             assert kafka_producer_mock.list_topics.call_count == 1
 
             # The interval hasn't passed and the check isn't retried.
             time_now[0] += 1
-            assert listener._check_topic_exists() is False
+            assert event_producer._check_topic_exists() is False
             assert kafka_producer_mock.list_topics.call_count == 1
 
             # Past the interval. The check is retried and succeeds this time.
             
kafka_producer_mock.list_topics.return_value.topics.__contains__.return_value = 
True
-            time_now[0] += listener._get_topic_check_retry_interval() + 1
-            assert listener._check_topic_exists() is True
+            time_now[0] += event_producer._get_topic_check_retry_interval() + 1
+            assert event_producer._check_topic_exists() is True
             assert kafka_producer_mock.list_topics.call_count == 2
 
     def test_producer_init_failure_sets_cooldown(self, monkeypatch):
         time_now = [1000.0]
-        monkeypatch.setattr(listener.time, "monotonic", lambda: time_now[0])
+        monkeypatch.setattr(event_producer.time, "monotonic", lambda: 
time_now[0])
 
         with patch(_PRODUCER_CLS, side_effect=ValueError("boom")) as 
producer_cls_mock:
-            assert listener._check_topic_exists() is False
+            assert event_producer._check_topic_exists() is False
             assert producer_cls_mock.call_count == 1
 
             # The interval hasn't passed and the producer initialization isn't 
retried.
             time_now[0] += 1
-            assert listener._check_topic_exists() is False
+            assert event_producer._check_topic_exists() is False
             assert producer_cls_mock.call_count == 1
 
             # Past the interval. The producer initialization is retried.
-            time_now[0] += listener._get_topic_check_retry_interval() + 1
-            assert listener._check_topic_exists() is False
+            time_now[0] += event_producer._get_topic_check_retry_interval() + 1
+            assert event_producer._check_topic_exists() is False
             assert producer_cls_mock.call_count == 2
 
     @pytest.mark.parametrize(
@@ -434,27 +434,27 @@ class TestCheckTopicExists:
         """
 
         time_now = [1000.0]
-        monkeypatch.setattr(listener.time, "monotonic", lambda: time_now[0])
+        monkeypatch.setattr(event_producer.time, "monotonic", lambda: 
time_now[0])
 
         kafka_producer_mock = MagicMock()
         
kafka_producer_mock.list_topics.return_value.topics.__contains__.return_value = 
True
         with patch(_PRODUCER_CLS, return_value=kafka_producer_mock):
-            assert listener._check_topic_exists() is True
-            assert listener._topic_exists is True
+            assert event_producer._check_topic_exists() is True
+            assert event_producer._topic_exists is True
 
             err = MagicMock()
             err.code.return_value = kafka_error
-            listener._on_delivery(err, MagicMock())
+            event_producer._on_delivery(err, MagicMock())
 
-            assert listener._topic_exists is topic_exists_expected
+            assert event_producer._topic_exists is topic_exists_expected
 
             if not topic_exists_expected:
                 # Cooldown holds off the immediate re-check.
-                assert listener._check_topic_exists() is False
+                assert event_producer._check_topic_exists() is False
                 # After the cooldown, the check re-verifies; the topic is 
still on the broker,
                 # so confirmation is restored.
-                time_now[0] += listener._get_topic_check_retry_interval() + 1
-                assert listener._check_topic_exists() is True
+                time_now[0] += 
event_producer._get_topic_check_retry_interval() + 1
+                assert event_producer._check_topic_exists() is True
 
 
 class TestFilters:
@@ -464,9 +464,9 @@ class TestFilters:
     def _filters_conf(self):
         with conf_vars(
             {
-                ("kafka_listener", "dag_run_events_enabled"): "True",
-                ("kafka_listener", "task_instance_events_enabled"): "True",
-                ("kafka_listener", "topic"): _KAFKA_TOPIC,
+                ("kafka_event_producer", "dag_run_events_enabled"): "True",
+                ("kafka_event_producer", "task_instance_events_enabled"): 
"True",
+                ("kafka_event_producer", "topic"): _KAFKA_TOPIC,
             }
         ):
             yield
@@ -482,55 +482,55 @@ class TestFilters:
         ],
     )
     def test_id_is_allowed(self, id_to_check, allowlist, denylist, 
expected_result):
-        assert listener._id_is_allowed(id_to_check, allowlist, denylist) is 
expected_result
+        assert event_producer._id_is_allowed(id_to_check, allowlist, denylist) 
is expected_result
 
     @conf_vars(
         {
-            ("kafka_listener", "task_instance_dag_id_denylist"): "dag1_*",
+            ("kafka_event_producer", "task_instance_dag_id_denylist"): 
"dag1_*",
         }
     )
     def test_dag_run_filter_ignores_ti_lists(self):
         # The dag_id lists for DR events, are separate from the dag_id lists 
for TI events.
         # The dag_id denylist for task instances would match "dag1_task1" but 
DR filter must ignore it.
-        assert listener._dag_run_event_allowed("dag1_task1") is True
+        assert event_producer._dag_run_event_allowed("dag1_task1") is True
 
     @conf_vars(
         {
-            ("kafka_listener", "dag_run_dag_id_denylist"): "dag1_*",
+            ("kafka_event_producer", "dag_run_dag_id_denylist"): "dag1_*",
         }
     )
     def test_task_instance_filter_ignores_dr_lists(self):
         # The dag_id lists for DR events, are separate from the dag_id lists 
for TI events.
         # DR denylist would match "dag1_task1" but TI filter must ignore it.
-        assert listener._task_instance_event_allowed("dag1_task1", "load") is 
True
+        assert event_producer._task_instance_event_allowed("dag1_task1", 
"load") is True
 
     @conf_vars(
         {
-            ("kafka_listener", "task_instance_dag_id_allowlist"): "dag1_*",
-            ("kafka_listener", "task_instance_task_id_denylist"): "*_cleanup",
+            ("kafka_event_producer", "task_instance_dag_id_allowlist"): 
"dag1_*",
+            ("kafka_event_producer", "task_instance_task_id_denylist"): 
"*_cleanup",
         }
     )
     def test_task_instance_requires_both_dag_id_and_task_id(self):
         # dag_id passes (matches allowlist), task_id passes (no deny match).
-        assert listener._task_instance_event_allowed("dag1_task1", "task2") is 
True
+        assert event_producer._task_instance_event_allowed("dag1_task1", 
"task2") is True
         # dag_id passes, task_id denied.
-        assert listener._task_instance_event_allowed("dag1_task1", 
"task2_cleanup") is False
+        assert event_producer._task_instance_event_allowed("dag1_task1", 
"task2_cleanup") is False
         # dag_id denied (no allowlist match), task_id passes.
-        assert listener._task_instance_event_allowed("other_dag", "load") is 
False
+        assert event_producer._task_instance_event_allowed("other_dag", 
"load") is False
 
-    @conf_vars({("kafka_listener", "dag_run_dag_id_denylist"): _DAG_ID})
+    @conf_vars({("kafka_event_producer", "dag_run_dag_id_denylist"): _DAG_ID})
     def test_dr_filter_blocks_emission(self, dr_mock, kafka_producer_mock):
         DagRunListener().on_dag_run_running(dag_run=dr_mock, msg=None)
         kafka_producer_mock.produce.assert_not_called()
 
-    @conf_vars({("kafka_listener", "task_instance_dag_id_denylist"): _DAG_ID})
+    @conf_vars({("kafka_event_producer", "task_instance_dag_id_denylist"): 
_DAG_ID})
     def test_ti_filter_blocks_emission_via_dag_id(self, ti_mock, 
kafka_producer_mock):
         TaskListener().on_task_instance_running(
             previous_state="queued", task_instance=ti_mock, **_TI_SESSION_KWARG
         )
         kafka_producer_mock.produce.assert_not_called()
 
-    @conf_vars({("kafka_listener", "task_instance_task_id_denylist"): 
_TASK_ID})
+    @conf_vars({("kafka_event_producer", "task_instance_task_id_denylist"): 
_TASK_ID})
     def test_ti_filter_blocks_emission_via_task_id(self, ti_mock, 
kafka_producer_mock):
         TaskListener().on_task_instance_running(
             previous_state="queued", task_instance=ti_mock, **_TI_SESSION_KWARG
diff --git a/scripts/ci/docker-compose/integration-kafka.yml 
b/scripts/ci/docker-compose/integration-kafka.yml
index f0ba7fab7d6..da67875dfb7 100644
--- a/scripts/ci/docker-compose/integration-kafka.yml
+++ b/scripts/ci/docker-compose/integration-kafka.yml
@@ -60,9 +60,9 @@ services:
   airflow:
     environment:
       - INTEGRATION_KAFKA=true
-      - AIRFLOW__KAFKA_LISTENER__DAG_RUN_EVENTS_ENABLED=True
-      - AIRFLOW__KAFKA_LISTENER__TASK_INSTANCE_EVENTS_ENABLED=True
-      - AIRFLOW__KAFKA_LISTENER__TOPIC=airflow.events
-      - AIRFLOW__KAFKA_LISTENER__SOURCE=dev-breeze
+      - AIRFLOW__KAFKA_EVENT_PRODUCER__DAG_RUN_EVENTS_ENABLED=True
+      - AIRFLOW__KAFKA_EVENT_PRODUCER__TASK_INSTANCE_EVENTS_ENABLED=True
+      - AIRFLOW__KAFKA_EVENT_PRODUCER__TOPIC=airflow.events
+      - AIRFLOW__KAFKA_EVENT_PRODUCER__SOURCE=dev-breeze
     depends_on:
       - broker


Reply via email to