aaron-y-chen commented on code in PR #73276:
URL: https://github.com/apache/airflow/pull/73276#discussion_r4232045630
##########
providers/apache/kafka/tests/integration/apache/kafka/plugins/test_event_producer.py:
##########
@@ -66,7 +66,8 @@ class TestEventProducer:
KAFKA_CONFIG_ID = "kafka_default"
# Use a unique topic per run to avoid errors on a re-run, in case
# the previous teardown hasn't finished with the topic deletion.
- TOPIC = f"airflow.events.itest.{uuid.uuid4().hex[:8]}"
+ DAGRUN_TOPIC = f"airflow.dagrun.itest.{uuid.uuid4().hex[:8]}"
Review Comment:
I think this breaks the integration test.
The class no longer defines `TOPIC`, but `setup_class` still references
`cls.TOPIC` when creating the topic (line 112). Teardown and the consumer also
reference the removed attribute. This causes setup to fail with
`AttributeError` before any events are tested.
Please update topic creation, cleanup, and consumption to use both new
topics, and verify that each event type is published to its corresponding topic.
##########
providers/apache/kafka/tests/unit/apache/kafka/plugins/test_event_producer.py:
##########
@@ -439,34 +504,90 @@ def
test_delivery_report_invalidates_topic_confirmation_only_on_missing_topic(
kafka_producer_mock = MagicMock()
kafka_producer_mock.list_topics.return_value.topics.__contains__.return_value =
True
with patch(_PRODUCER_CLS, return_value=kafka_producer_mock):
- assert event_producer._check_topic_exists() is True
- assert event_producer._topic_exists is True
+ assert event_producer._check_topic_exists(topic_type) is True
+ assert event_producer._topic_existence_map[topic_type] is True
err = MagicMock()
err.code.return_value = kafka_error
- event_producer._on_delivery(err, MagicMock())
+ event_producer._on_delivery(topic_type, err, MagicMock())
- assert event_producer._topic_exists is topic_exists_expected
+ assert event_producer._topic_existence_map[topic_type] is
topic_exists_expected
if not topic_exists_expected:
# Cooldown holds off the immediate re-check.
- assert event_producer._check_topic_exists() is False
+ assert event_producer._check_topic_exists(topic_type) is False
# After the cooldown, the check re-verifies; the topic is
still on the broker,
# so confirmation is restored.
time_now[0] +=
event_producer._get_topic_check_retry_interval() + 1
- assert event_producer._check_topic_exists() is True
+ assert event_producer._check_topic_exists(topic_type) is True
+
+
[email protected](
+ (
+ "topic",
+ "dagrun_topic_setting",
+ "task_instance_topic_setting",
+ "dagrun_expected",
+ "task_instance_expected"
+ ),
+ [
+ pytest.param(
+ "airflow.events",
+ "airflow.dagrun",
+ "airflow.task",
+ "airflow.events",
+ "airflow.events",
+ id="default_to_topic",
+ ),
+ pytest.param(None, None, None, "airflow.events", "airflow.events",
id="fallback_to_default"),
+ pytest.param(
+ None, "airflow.dagrun", "airflow.task", "airflow.dagrun",
"airflow.task", id="new_settings"
+ ),
+ pytest.param("airflow.foo", None, None, "airflow.foo", "airflow.foo",
id="old_setting"),
+ ],
+)
+def test_get_topic(
+ topic, dagrun_topic_setting, task_instance_topic_setting, dagrun_expected,
task_instance_expected
+):
+ # Check that ``topic`` is properly deprecated
+ ctxt = (
+ pytest.raises(
+ AirflowProviderDeprecationWarning,
+ match=r"""The ``topic`` option in \[kafka_event_producer\] has
been split into ``dagrun_topic`` and ``task_instance_topic`` -
+ please update your config\.""",
+ )
Review Comment:
`AirflowProviderDeprecationWarning` is listed in `forbidden_warnings` in our
pytest config, so this warning is raised as an exception. `pytest.raises`
catches it on the first `_get_topic(...)` call and exits the `with` block, so
the two asserts below never run for the `default_to_topic` and `old_setting`
cases. I checked locally: the test still passes after changing the
`old_setting` expected values to `"WRONG"`.
Could you use `pytest.warns(AirflowProviderDeprecationWarning, match=...)`
instead, so that the deprecated `topic` path is actually checked?
--
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]