diff --git a/providers/apache/kafka/src/airflow/providers/apache/kafka/operators/produce.py b/providers/apache/kafka/src/airflow/providers/apache/kafka/operators/produce.py index b8b28b929df84..90ab306bc04db 100644 --- a/providers/apache/kafka/src/airflow/providers/apache/kafka/operators/produce.py +++ b/providers/apache/kafka/src/airflow/providers/apache/kafka/operators/produce.py @@ -100,15 +100,13 @@ def __init__( self.synchronous = synchronous self.poll_timeout = poll_timeout + def execute(self, context) -> None: if not (self.topic and self.producer_function): raise AirflowException( "topic and producer_function must be provided. Got topic=" f"{self.topic} and producer_function={self.producer_function}" ) - return - - def execute(self, context) -> None: # Get producer and callable producer = KafkaProducerHook(kafka_config_id=self.kafka_config_id).get_producer() diff --git a/providers/apache/kafka/tests/unit/apache/kafka/operators/test_produce.py b/providers/apache/kafka/tests/unit/apache/kafka/operators/test_produce.py index 92f408408aa58..b90ace4663d7e 100644 --- a/providers/apache/kafka/tests/unit/apache/kafka/operators/test_produce.py +++ b/providers/apache/kafka/tests/unit/apache/kafka/operators/test_produce.py @@ -20,10 +20,13 @@ import logging from typing import Any +import pendulum import pytest from airflow.models import Connection +from airflow.models.dag import DAG from airflow.providers.apache.kafka.operators.produce import ProduceToTopicOperator +from airflow.providers.common.compat.sdk import AirflowException log = logging.getLogger(__name__) @@ -85,3 +88,18 @@ def test_operator_callable(self): ) operator.execute(context={}) + + def test_execute_rejects_empty_rendered_topic(self): + with DAG("kafka_produce", schedule=None, start_date=pendulum.datetime(2020, 1, 1, tz="UTC")): + operator = ProduceToTopicOperator( + kafka_config_id="kafka_d", + topic="{{ params.topic }}", + producer_function=_simple_producer, + producer_function_args=(b"test", b"test"), + task_id="test", + synchronous=False, + ) + operator.render_template_fields({"params": {"topic": ""}}) + assert operator.topic == "" + with pytest.raises(AirflowException, match="topic and producer_function must be provided"): + operator.execute(context={}) diff --git a/scripts/ci/prek/validate_operators_init_exemptions.txt b/scripts/ci/prek/validate_operators_init_exemptions.txt index ec735e38dde53..6adf724c9c9d9 100644 --- a/scripts/ci/prek/validate_operators_init_exemptions.txt +++ b/scripts/ci/prek/validate_operators_init_exemptions.txt @@ -21,7 +21,6 @@ providers/amazon/src/airflow/providers/amazon/aws/transfers/base.py::AwsToAwsBas providers/amazon/src/airflow/providers/amazon/aws/transfers/gcs_to_s3.py::GCSToS3Operator providers/amazon/src/airflow/providers/amazon/aws/transfers/s3_to_redshift.py::S3ToRedshiftOperator providers/anthropic/src/airflow/providers/anthropic/operators/agent.py::AnthropicAgentSessionOperator -providers/apache/kafka/src/airflow/providers/apache/kafka/operators/produce.py::ProduceToTopicOperator providers/cncf/kubernetes/src/airflow/providers/cncf/kubernetes/operators/pod.py::KubernetesPodOperator providers/docker/src/airflow/providers/docker/operators/docker.py::DockerOperator providers/google/src/airflow/providers/google/cloud/operators/bigquery.py::BigQueryInsertJobOperator