Extracting with TaskFlow

Extracting BookNest's Orders with TaskFlow

The extract task is three lines because the work lives in booknest_pipeline.py, shared with BookNest's Daily Pipeline:

booknest_pipeline.py (excerpt): the extract stepPython
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, n

Its log in the UI (task ingest.extract, try 1) holds the print() output and the returned value:

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