diff --git a/providers/apache/kafka/src/airflow/providers/apache/kafka/hooks/base.py b/providers/apache/kafka/src/airflow/providers/apache/kafka/hooks/base.py index 5d02903a4d692..0b389535007d7 100644 --- a/providers/apache/kafka/src/airflow/providers/apache/kafka/hooks/base.py +++ b/providers/apache/kafka/src/airflow/providers/apache/kafka/hooks/base.py @@ -22,7 +22,6 @@ from confluent_kafka.admin import AdminClient from airflow.hooks.base import BaseHook -from airflow.providers.google.cloud.hooks.managed_kafka import ManagedKafkaHook class KafkaBaseHook(BaseHook): @@ -70,6 +69,15 @@ def get_conn(self) -> Any: and bootstrap_servers.find("cloud.goog") != -1 and bootstrap_servers.find("managedkafka") != -1 ): + try: + from airflow.providers.google.cloud.hooks.managed_kafka import ManagedKafkaHook + except ImportError: + from airflow.exceptions import AirflowOptionalProviderFeatureException + + raise AirflowOptionalProviderFeatureException( + "Failed to import ManagedKafkaHook. For using this functionality google provider version " + ">= 14.1.0 should be pre-installed." + ) self.log.info("Adding token generation for Google Auth to the confluent configuration.") hook = ManagedKafkaHook() token = hook.get_confluent_token