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]

Reply via email to