A Custom Deferrable Sensor

Writing a Custom Deferrable Sensor

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

include/booknest_triggers.py (excerpt): a trigger and the sensor that defers to itPython
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:

Output of 20
 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.