This is an automated email from the ASF dual-hosted git repository.
potiuk 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 353d856281c Stop Kafka hook tests from leaving live clients behind
(#74285)
353d856281c is described below
commit 353d856281c8002621906b6abb60a3baec9cf062
Author: Jarek Potiuk <[email protected]>
AuthorDate: Mon Oct 5 19:53:33 2026 +0200
Stop Kafka hook tests from leaving live clients behind (#74285)
test_get_consumer and test_get_producer patched AdminClient, which the
consumer and producer hooks never use, so each built a real
confluent_kafka client pointed at localhost:9092. The client stays
cached on the hook held by the test instance, so its native threads keep
retrying the connection for the rest of the pytest session. In the
provider test runs this shows up as a stream of "rdkafka#consumer-1" and
"rdkafka#producer-2" connection errors across unrelated tests, and leaves
native threads running in a process that later segfaulted.
Generated-by: Claude Opus 5
---
.../apache/kafka/tests/unit/apache/kafka/hooks/test_consume.py | 9 +++------
.../apache/kafka/tests/unit/apache/kafka/hooks/test_produce.py | 9 +++------
2 files changed, 6 insertions(+), 12 deletions(-)
diff --git
a/providers/apache/kafka/tests/unit/apache/kafka/hooks/test_consume.py
b/providers/apache/kafka/tests/unit/apache/kafka/hooks/test_consume.py
index bc48fa4dfae..120b43a91ef 100644
--- a/providers/apache/kafka/tests/unit/apache/kafka/hooks/test_consume.py
+++ b/providers/apache/kafka/tests/unit/apache/kafka/hooks/test_consume.py
@@ -17,10 +17,9 @@
from __future__ import annotations
import json
-from unittest.mock import MagicMock, patch
+from unittest.mock import patch
import pytest
-from confluent_kafka.admin import AdminClient
from airflow.models import Connection
@@ -56,10 +55,8 @@ class TestConsumerHook:
)
self.hook = KafkaConsumerHook(["test_1"], kafka_config_id="kafka_d")
- @patch("airflow.providers.apache.kafka.hooks.base.AdminClient")
- def test_get_consumer(self, mock_client):
- mock_client_spec = MagicMock(spec=AdminClient)
- mock_client.return_value = mock_client_spec
+ @patch("airflow.providers.apache.kafka.hooks.consume.Consumer")
+ def test_get_consumer(self, mock_consumer):
assert self.hook.get_consumer() == self.hook.get_conn
@patch("airflow.providers.apache.kafka.hooks.consume.Consumer")
diff --git
a/providers/apache/kafka/tests/unit/apache/kafka/hooks/test_produce.py
b/providers/apache/kafka/tests/unit/apache/kafka/hooks/test_produce.py
index 6c35eeb606f..6e571f4cd8a 100644
--- a/providers/apache/kafka/tests/unit/apache/kafka/hooks/test_produce.py
+++ b/providers/apache/kafka/tests/unit/apache/kafka/hooks/test_produce.py
@@ -18,10 +18,9 @@ from __future__ import annotations
import json
import logging
-from unittest.mock import MagicMock, patch
+from unittest.mock import patch
import pytest
-from confluent_kafka.admin import AdminClient
from airflow.models import Connection
from airflow.providers.apache.kafka.hooks.produce import KafkaProducerHook
@@ -57,10 +56,8 @@ class TestProducerHook:
)
self.hook = KafkaProducerHook(kafka_config_id="kafka_d")
- @patch("airflow.providers.apache.kafka.hooks.base.AdminClient")
- def test_get_producer(self, mock_client):
- mock_client_spec = MagicMock(spec=AdminClient)
- mock_client.return_value = mock_client_spec
+ @patch("airflow.providers.apache.kafka.hooks.produce.Producer")
+ def test_get_producer(self, mock_producer):
assert self.hook.get_producer() == self.hook.get_conn
@conf_vars({("apache_kafka", "callback_allowlist"): "json.dumps"})