A file's existence is a weak signal: the order service appends to its export all day. BookNest needs a watermark: the export already holds an order from after the end of day ds, so no later line can belong to that day. No provider ships that check, so include/booknest_triggers.py writes a trigger and a sensor (newest_order_ts() reads the export's last line; day_end() returns the next midnight as an ISO string):
class ExportWatermarkTrigger(BaseTrigger):
def serialize(self): # how the triggerer re-creates this object
return ("booknest_triggers.ExportWatermarkTrigger",
{"path": self.path, "ds": self.ds, "poll_interval": self.poll_interval})
async def run(self):
while True:
newest = await asyncio.to_thread(newest_order_ts, self.path) # off the event loop
if newest >= day_end(self.ds):
yield TriggerEvent({"ds": self.ds, "newest": newest})
return
self.log.info("export ends at %s, waiting for %s", newest, day_end(self.ds))
await asyncio.sleep(self.poll_interval)
class ExportCompleteSensor(BaseSensorOperator):
def execute(self, context):
newest = newest_order_ts(self.path)
if newest >= day_end(self.ds): # already complete: no need to defer
return newest
self.defer(trigger=ExportWatermarkTrigger(self.path, self.ds, self.poll_interval),
method_name="execute_complete", timeout=timedelta(seconds=self.timeout))serialize() is the contract: the triggerer stores a class path and JSON arguments, so the class must be importable there, and the docs forbid loading it from a DAG bundle (include/ is on every container's PYTHONPATH). run7/w2.sh started runs for 29 and 30 June, then appended a sample order dated 1 July:
ready-29 | success | 1 | ready-30 | deferred | 1 | execute_complete export ends at 2026-06-30T23:55:02Z, waiting for 2026-07-01T00:00:00Z "event":"export complete for 2026-06-30, newest order 2026-07-01T00:03:12Z"
29 June was already complete, so its sensor returned without deferring. 30 June waited in the triggerer, whose log lands beside the task's (attempt=1.log.trigger.13.log), and resumed within one poll in execute_complete(), which logs the event and returns the watermark as the task's XCom. Keep triggers stateless: after a triggerer restart the same trigger can run again.