Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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()

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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__)

Expand Down Expand Up @@ -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={})
1 change: 0 additions & 1 deletion scripts/ci/prek/validate_operators_init_exemptions.txt
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down