xBis7 commented on code in PR #73276:
URL: https://github.com/apache/airflow/pull/73276#discussion_r4199625642


##########
providers/apache/kafka/docs/configurations-ref.rst:
##########


Review Comment:
   There is this and a few other `topic` references in this doc that are 
ambiguous. They make sense as they are but we could also add some extra 
explanation. I leave it up to you to decide.



##########
providers/apache/kafka/src/airflow/providers/apache/kafka/plugins/event_producer.py:
##########
@@ -202,15 +243,15 @@ def _get_producer() -> Producer | None:
     return _producer
 
 
-def _check_topic_exists() -> bool:
+def _check_topic_exists(topic_type: EventProducerKafkaTopic) -> bool:
     """
     Verify the configured topic exists on the broker, with a 
retry-after-failure cooldown.
 
     Once the topic has been confirmed, the result is kept for the process 
lifetime; until
     then, failed checks are retried on a configured interval.
     """
-    global _topic_exists, _topic_check_retry_after
-    if _topic_exists:
+    global _topic_existence_map, _topic_check_retry_after  # noqa: PLW0602

Review Comment:
   `_topic_check_retry_after` is the same for both topics but we need one per 
type. To explain why with an example:
   
   * Dag run topic does NOT exist
   * Task instance topic exists
   *  Call `_check_topic_exists` for the Dag run topic
       * It doesn't exist,
       * Setting `_topic_check_retry_after = time.monotonic() + 
_get_topic_check_retry_interval()`
       * Function returns `False`
   * Call `_check_topic_exists` for the Task instance topic
       * Now the global `_topic_check_retry_after` is set and until the 
interval passes, the Task instance topic check will remain `False`
       * It will get to the part that converts it to `True` only when the 
interval passes and the Dag run check hasn't run again to reset the interval
           ```
           _topic_existence_map[topic_type] = True
           return True
           ```



##########
providers/apache/kafka/docs/configurations-ref.rst:
##########


Review Comment:
   ```suggestion
       AIRFLOW__KAFKA_EVENT_PRODUCER__DAGRUN_TOPIC=airflow.dagrun
       AIRFLOW__KAFKA_EVENT_PRODUCER__TASK_INSTANCE_TOPIC=airflow.task_instance
   ```



##########
providers/apache/kafka/src/airflow/providers/apache/kafka/plugins/event_producer.py:
##########
@@ -366,10 +417,10 @@ def _produce_message(event: str, dag_id: str, run_id: 
str, payload: dict[str, An
     key = f"{dag_id}/{run_id}".encode()
     try:
         producer.produce(
-            _get_topic(),
+            _get_topic(topic_type),
             key=key,
             value=json.dumps(body, default=str).encode("utf-8"),
-            on_delivery=_on_delivery,
+            on_delivery=_on_delivery_map[topic_type],

Review Comment:
   Wouldn't this be simpler if you used a lambda? This is the only place where 
the `_on_delivery_map` is used.
   
   ```suggestion
               on_delivery=lambda err, msg: _on_delivery(topic_type, err, msg),
   ```
    
   or
   
   ```suggestion
               on_delivery=functools.partial(_on_delivery, topic_type),
   ```
   
   I think the 2nd approach is easier to test.



-- 
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