From 9ea18c00843af941c085fb4291ea624545b1fef8 Mon Sep 17 00:00:00 2001 From: LIU ZHE YOU Date: Sun, 3 May 2026 16:41:13 +0800 Subject: [PATCH 1/2] Fix AwaitMessageTrigger missing _task_instance attribute --- .../providers/apache/kafka/triggers/await_message.py | 1 + .../unit/apache/kafka/triggers/test_await_message.py | 10 ++++++++++ 2 files changed, 11 insertions(+) diff --git a/providers/apache/kafka/src/airflow/providers/apache/kafka/triggers/await_message.py b/providers/apache/kafka/src/airflow/providers/apache/kafka/triggers/await_message.py index 196a1d8aa4899..3dbc5e467d419 100644 --- a/providers/apache/kafka/src/airflow/providers/apache/kafka/triggers/await_message.py +++ b/providers/apache/kafka/src/airflow/providers/apache/kafka/triggers/await_message.py @@ -82,6 +82,7 @@ def __init__( poll_interval: float = 5, commit_offset: bool = True, ) -> None: + super().__init__() self.topics = topics self.apply_function = apply_function self.apply_function_args = apply_function_args or () diff --git a/providers/apache/kafka/tests/unit/apache/kafka/triggers/test_await_message.py b/providers/apache/kafka/tests/unit/apache/kafka/triggers/test_await_message.py index 397992ec9ea40..65d843884ffd6 100644 --- a/providers/apache/kafka/tests/unit/apache/kafka/triggers/test_await_message.py +++ b/providers/apache/kafka/tests/unit/apache/kafka/triggers/test_await_message.py @@ -77,6 +77,16 @@ def setup_connections(self, create_connection_without_db): ) ) + def test_trigger_initializes_base_state(self): + trigger = AwaitMessageTrigger( + kafka_config_id="kafka_d", + apply_function="test.noop", + topics=["noop"], + ) + + assert trigger._task_instance is None + assert trigger.task_instance is None + def test_trigger_serialization(self): trigger = AwaitMessageTrigger( kafka_config_id="kafka_d", From a6237e21a9587b15ebc3cb786856950a122186f1 Mon Sep 17 00:00:00 2001 From: LIU ZHE YOU Date: Thu, 7 May 2026 16:25:25 +0800 Subject: [PATCH 2/2] CI: Fix compat test --- .../unit/apache/kafka/triggers/test_await_message.py | 10 ++++++++-- 1 file changed, 8 insertions(+), 2 deletions(-) diff --git a/providers/apache/kafka/tests/unit/apache/kafka/triggers/test_await_message.py b/providers/apache/kafka/tests/unit/apache/kafka/triggers/test_await_message.py index 65d843884ffd6..374fdb0f2f96d 100644 --- a/providers/apache/kafka/tests/unit/apache/kafka/triggers/test_await_message.py +++ b/providers/apache/kafka/tests/unit/apache/kafka/triggers/test_await_message.py @@ -30,6 +30,7 @@ collect_queue_param_deprecation_warning, mark_common_msg_queue_test, ) +from tests_common.test_utils.version_compat import AIRFLOW_V_3_3_PLUS USED_FIXTURES = [collect_queue_param_deprecation_warning] @@ -84,8 +85,13 @@ def test_trigger_initializes_base_state(self): topics=["noop"], ) - assert trigger._task_instance is None - assert trigger.task_instance is None + if AIRFLOW_V_3_3_PLUS: + # The trigger._task_instance attribute was introduced in https://github.com/apache/airflow/pull/55068 + assert trigger._task_instance is None + assert trigger.task_instance is None + else: + assert not hasattr(trigger, "_task_instance") + assert trigger.task_instance is None def test_trigger_serialization(self): trigger = AwaitMessageTrigger(