uranusjr commented on code in PR #68517:
URL: https://github.com/apache/airflow/pull/68517#discussion_r4204080677


##########
airflow-core/src/airflow/timetables/simple.py:
##########
@@ -219,26 +221,58 @@ class AssetTriggeredTimetable(_TrivialTimetable):
     description: str = "Triggered by assets"
     asset_triggered = True
 
-    def __init__(self, assets: Collection[SerializedAsset] | 
SerializedAssetBase) -> None:
+    def __init__(
+        self,
+        assets: Collection[SerializedAsset] | SerializedAssetBase,
+        *,
+        batch_asset_events: bool | None = None,
+    ) -> None:
         super().__init__()
         # Compatibility: Handle SDK assets if needed so this class works in 
dag files.
         if isinstance(assets, SerializedAssetBase | BaseAsset):
             self.asset_condition = ensure_serialized_asset(assets)
         else:
             self.asset_condition = 
SerializedAssetAll([ensure_serialized_asset(a) for a in assets])
+        if batch_asset_events is None:
+            batch_asset_events = (
+                conf.getboolean("scheduler", "batch_asset_events", 
fallback=False)
+                or self.get_batching_requirement() is not None
+            )
+        self.batch_asset_events = batch_asset_events
 
     @classmethod
     def deserialize(cls, data: dict[str, Any]) -> Timetable:
         from airflow.serialization.decoders import decode_asset_like
 
-        return cls(decode_asset_like(data["asset_condition"]))
+        return cls(
+            decode_asset_like(data["asset_condition"]),
+            batch_asset_events=data.get("batch_asset_events", True),
+        )
 
     @property
     def summary(self) -> str:
         return "Asset"
 
     def serialize(self) -> dict[str, Any]:
-        return {"asset_condition": encode_asset_like(self.asset_condition)}
+        return {
+            "asset_condition": encode_asset_like(self.asset_condition),
+            "batch_asset_events": self.batch_asset_events,
+        }
+
+    def get_batching_requirement(self) -> str | None:
+        """Return why this timetable cannot run without event batching, or 
``None``."""
+        pending = [self.asset_condition]
+        while pending:
+            condition = pending.pop()
+            if isinstance(condition, SerializedAssetAll) and 
len(condition.objects) > 1:

Review Comment:
   Can this be a property on the asset condition classes instead (AssetAll 
returns True when it has more than one object, boolean conditions combine their 
children), like `PartitionMapper.is_rollup`? Then new condition types can 
declare this themselves.



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