Test Yourself!

These ten questions cover the chapter, and each does something that surprises people who have built pipelines for years: an upsert that rejects its own batch, one cron string that labels a day two ways, a schedule that skips four months, a manual run that cannot render {{ ds }}, a number that loses digits between tasks, a script mistaken for a template, a mapped task with nothing to map, a cache that needs a setting, an update that forgets the old row, and a merge that merges nothing. All ran against BookNest's data: PostgreSQL 18.6 1,289 (l2-pg, wal_level=logical), Airflow 3.3.2 129 in Securing Orchestration's stack, Prefect 3.8.7 74,615 and dlt 1.30.0 876,044 . Write down what each prints and why, then check Appendix I for the real output, the reason, the fix and the section.

check.sh in demos/ch07/quiz/ runs questions.sql through psql, copies quiz_dags.py into the DAG folder, unpauses its DAGs, triggers quiz_ds over the REST API, prints each run's task states and log lines, and then runs the Prefect and dlt scripts.

questions.sql: questions 1 and 9 in PostgreSQLSQL
-- questions.sql: questions 1 and 9 (PostgreSQL 18.6, wal_level = logical; check.sh resets).
-- 1. One batch holds the same payment twice, the next two versions of one. (Section 7.1.2)
INSERT INTO payments VALUES ('Q1A', 100002, 50.00), ('Q1A', 100002, 50.00)
  ON CONFLICT (event_id) DO NOTHING;
INSERT INTO payments VALUES ('Q1B', 100003, 20.00), ('Q1B', 100003, 25.00)
  ON CONFLICT (event_id) DO UPDATE SET amount = excluded.amount;
SELECT 1, event_id, amount FROM payments WHERE event_id IN ('Q1A', 'Q1B');
-- 9. Two updates through a logical slot, the second after REPLICA IDENTITY FULL. (7.13.3)
CREATE TABLE quiz_cdc AS SELECT order_id, status FROM orders ORDER BY order_id LIMIT 1;
ALTER TABLE quiz_cdc ADD PRIMARY KEY (order_id);
SELECT FROM pg_create_logical_replication_slot('quiz', 'test_decoding');
UPDATE quiz_cdc SET status = 'returned';
ALTER TABLE quiz_cdc REPLICA IDENTITY FULL;
UPDATE quiz_cdc SET status = 'delivered';
SELECT 9, data FROM pg_logical_slot_get_changes('quiz', NULL, NULL) WHERE data LIKE 'table%';
quiz_dags.py: questions 2-7 in AirflowPython
# quiz_dags.py: questions 2-7 (Airflow 3.3.2). check.sh unpauses the DAGs and reads their runs.
from datetime import datetime, timezone
from decimal import Decimal
from airflow.providers.standard.operators.bash import BashOperator
from airflow.sdk import DAG, task
from airflow.timetables.interval import CronDataIntervalTimetable
JUNE1, OCT1 = (datetime(2026, m, 1, tzinfo=timezone.utc) for m in (6, 10))
# 2. One cron string, two timetables. Which ds does each give the 02:00 run on 2 October?
#    (Sections 7.3.1 and 7.5.4)
for name, every in [("trigger", "0 2 * * *"),
                    ("interval", CronDataIntervalTimetable("0 2 * * *", timezone="UTC"))]:
    with DAG(f"quiz_cron_{name}", schedule=every, start_date=OCT1):
        BashOperator(task_id="show", bash_command="echo ds={{ ds }}")
# 3. Daily since 1 June, unpaused on 2 October. How many runs appear? (Section 7.3.1)
with DAG("quiz_catchup", schedule="@daily", start_date=JUNE1):
    BashOperator(task_id="noop", bash_command="true")
# 4. Triggered over the REST API with {"logical_date": null}. What happens? (Section 7.3.1)
with DAG("quiz_ds", schedule=None, start_date=JUNE1):
    BashOperator(task_id="show", bash_command="echo ds={{ ds }}")
# 5. Orders, gross and an 18-digit RM-to-USD rate cross XCom. What arrives? (Section 7.5.1)
with DAG("quiz_xcom", schedule="@once", start_date=JUNE1):
    @task
    def totals(): return (186, Decimal("6855.25"), Decimal("0.212345678901234567"))
    @task
    def show(t): print("5:", type(t).__name__, t)
    show(totals())
# 6. A BashOperator that runs the pipeline's shell script. (Section 7.5.2)
with DAG("quiz_bash", schedule="@once", start_date=JUNE1):
    BashOperator(task_id="run_day", bash_command="run_day.sh")
# 7. Map the load over today's files when there are none. (Section 7.6.4)
with DAG("quiz_map", schedule="@once", start_date=JUNE1):
    @task
    def todays_files(): return []
    @task
    def load(path): return path
    @task
    def publish(): return "published"
    load.expand(path=todays_files()) >> publish()
quiz_prefect.py: question 8 in PrefectPython
# quiz_prefect.py: question 8 (Prefect 3.8.7; with no server set, a temporary one starts).
# 8. Two flow runs each call both tasks twice. How often does each body run? (Section 7.10.3)
from prefect import flow, task
calls = {"plain": 0, "persisted": 0}
@task
def plain(ds): calls["plain"] += 1
@task(persist_result=True)
def persisted(ds): calls["persisted"] += 1
@flow
def daily(ds):
    for _ in range(2):
        plain(ds), persisted(ds)
daily("2026-06-28"), daily("2026-06-28")
print("8:", calls)
quiz_dlt.py: question 10 in dltPython
# quiz_dlt.py: question 10 (dlt 1.30.0 into l2-pg; dev_mode makes a fresh dataset each run).
# 10. Two orders, merged twice as their status changes. What is in the table? (7.14.4)
import dlt
pg = dlt.destinations.postgres("postgresql://postgres:booknest@localhost:32543/booknest")
pipe = dlt.pipeline("quiz_merge", destination=pg, dataset_name="dlt_quiz", dev_mode=True)
for status in ("paid", "shipped"):
    pipe.run([{"order_id": 1, "status": status}, {"order_id": 2, "status": status}],
             table_name="orders", write_disposition="merge")
with pipe.sql_client() as sql:
    print("10:", sql.execute_sql("SELECT order_id, status FROM orders ORDER BY 1, 2"))