potiuk commented on code in PR #69646:
URL: https://github.com/apache/airflow/pull/69646#discussion_r4175292629
##########
providers/apache/kafka/src/airflow/providers/apache/kafka/operators/produce.py:
##########
@@ -58,7 +62,16 @@ class ProduceToTopicOperator(BaseOperator):
:param synchronous: If writes to kafka should be fully synchronous,
defaults to True
:param poll_timeout: How long of a delay should be applied when calling
poll after production to kafka,
defaults to 0
- :raises AirflowException: _description_
+ :param raise_on_delivery_failure: Kafka reports delivery outcomes
asynchronously through the
+ delivery callback, so a message that is rejected by the broker (e.g.
too large, unknown
+ topic, insufficient permissions) does not make the task fail by
default - the error is only
+ passed to the delivery callback. Set to True to fail the task with an
``AirflowException``
Review Comment:
This raises `KafkaMessageDeliveryError`, not a bare `AirflowException` —
naming the concrete class tells users what to catch:
```suggestion
passed to the delivery callback. Set to True to fail the task with a
``KafkaMessageDeliveryError``
```
##########
providers/apache/kafka/tests/unit/apache/kafka/operators/test_produce.py:
##########
@@ -103,3 +118,81 @@ def test_execute_rejects_empty_rendered_topic(self):
assert operator.topic == ""
with pytest.raises(AirflowException, match="topic and
producer_function must be provided"):
operator.execute(context={})
+
+ @mock.patch(GET_PRODUCER_PATH)
Review Comment:
Two things on the new tests:
- Please use `autospec=True` on these patches
(`@mock.patch(GET_PRODUCER_PATH, autospec=True)`), so a wrong call signature on
the producer fails the test instead of passing silently.
- `test_operator_raise_on_delivery_failure`, `..._successful_delivery` and
`test_operator_delivery_failure_ignored_by_default` differ only in the delivery
error, the flag value and the expected outcome — please fold them into one
`@pytest.mark.parametrize` test. The custom-callback test can stay separate
since it asserts something different.
##########
providers/apache/kafka/src/airflow/providers/apache/kafka/operators/produce.py:
##########
@@ -58,7 +62,16 @@ class ProduceToTopicOperator(BaseOperator):
:param synchronous: If writes to kafka should be fully synchronous,
defaults to True
:param poll_timeout: How long of a delay should be applied when calling
poll after production to kafka,
defaults to 0
- :raises AirflowException: _description_
+ :param raise_on_delivery_failure: Kafka reports delivery outcomes
asynchronously through the
+ delivery callback, so a message that is rejected by the broker (e.g.
too large, unknown
+ topic, insufficient permissions) does not make the task fail by
default - the error is only
+ passed to the delivery callback. Set to True to fail the task with an
``AirflowException``
+ when at least one message could not be delivered. Because the task
fails after the batch has
+ already been partially produced, a retry re-produces the whole batch,
so messages that were
+ delivered successfully the first time will be produced again. Defaults
to False to keep
+ backwards compatibility.
+ :raises AirflowException: If ``raise_on_delivery_failure`` is True and at
least one message
Review Comment:
Same here:
```suggestion
:raises KafkaMessageDeliveryError: If ``raise_on_delivery_failure`` is
True and at least one message
```
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]