The extract task is three lines because the work lives in booknest_pipeline.py, shared with BookNest's Daily Pipeline:
def extract_orders(ds, source, landing):
"""Copy one day's orders from the order service's export into the landing zone."""
out_dir = os.path.join(landing, "orders", f"date={ds}")
os.makedirs(out_dir, exist_ok=True)
out, tmp = os.path.join(out_dir, "orders.jsonl"), os.path.join(out_dir, ".orders.tmp")
n = 0
with open(source, encoding="utf-8") as src, open(tmp, "w", encoding="utf-8") as dst:
for line in src:
if json.loads(line)["order_ts"][:10] == ds: # order_ts is UTC ISO 8601
dst.write(line)
n += 1
os.replace(tmp, out) # atomic: readers never see a half-written file
return out, nIts log in the UI (task ingest.extract, try 1) holds the print() output and the returned value:
21:37:45 info extracted 192 orders for 2026-06-29 21:37:45 info Done. Returned value was: /opt/airflow/data/landing/orders/date=2026-06-29/orders .jsonl
Each choice serves the orchestrator: the day comes from ds, so retries and backfills extract the same day; streaming keeps memory flat; the rename means a killed task never leaves a partial file; the date= partition makes a re-run replace exactly that day; and the XCom is a small path.