seanghaeli commented on code in PR #74293:
URL: https://github.com/apache/airflow/pull/74293#discussion_r4213182018
##########
providers/amazon/src/airflow/providers/amazon/aws/triggers/s3.py:
##########
@@ -22,11 +22,15 @@
from typing import TYPE_CHECKING, Any
from airflow.providers.amazon.aws.hooks.s3 import S3Hook
+from airflow.providers.common.compat.triggers import BaseEventTrigger
from airflow.triggers.base import BaseTrigger, TriggerEvent
if TYPE_CHECKING:
from datetime import datetime
+# Key under which the last-reported object fingerprint is persisted in the
asset state store.
+WATERMARK_KEY = "etag"
Review Comment:
I think we should make a unique key per S3 bucket + object key, in the same
way KinesisTrigger._asset_store_checkpoint_key does:
https://github.com/apache/airflow/blob/3e302fd257a0cfe8deba2de529bdcae6606b83c8/providers/amazon/src/airflow/providers/amazon/aws/triggers/kinesis.py#L133-L144
Currently, multiple triggers attached to the same asset would both be
writing to the same key "etag", which means they can overwrite each other.
It's not a big consequence, I think all it can lead to is unnecessary extra
Dag runs if `Trigger1` updates the "etag" key, and then the next time
`Trigger2` reads it after a restart, it will think it needs to fire.
##########
providers/amazon/src/airflow/providers/amazon/aws/triggers/s3.py:
##########
@@ -22,11 +22,15 @@
from typing import TYPE_CHECKING, Any
from airflow.providers.amazon.aws.hooks.s3 import S3Hook
+from airflow.providers.common.compat.triggers import BaseEventTrigger
Review Comment:
`common.compat.triggers` was added in `AMPP 1.19.0`, so I believe this would
force us to require in `pyproject.toml` that the user has at least this version.
--
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]