Events land in raw_events (event_id, user_id, event_type, event_ts, ingested_at): event_ts is when the event happened on the device, ingested_at when it reached the warehouse. Delivery is at-least-once, so one event_id can arrive more than once, and phones that were offline deliver old events hours late.
A job appends new events to fct_events (event_id, user_id, event_type, event_ts). etl_watermark (table_name, high_water_mark) holds the job's mark under 'raw_events': the latest ingested_at it has fully processed. Rows ingested exactly at the mark were processed. The last run failed partway, after inserting some events but before moving the mark.
Return the rows the next run should insert, so that re-running it never creates a duplicate in fct_events:
- •everything that arrived after the mark, however old its
event_ts - •one row per
event_id, from its earliest delivery after the mark - •nothing already in
fct_events
Columns: event_id, user_id, event_type, event_ts, ingested_at. Sort by event_id.