Before TaskFlow, every task was an operator instance: a class with an execute() method, configured by arguments and wired with >>. Operators remain best when someone already wrote the logic, as providers have for SQL, Bash, Kubernetes 5,150 and cloud services. The extract plus a row count, in classic style:
from datetime import datetime, timezone
from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator
from airflow.providers.standard.operators.python import PythonOperator
from airflow.sdk import DAG
import booknest_pipeline as p
def extract_orders(ds, **_):
path, n = p.extract_orders(ds, "/opt/airflow/data/source/orders.jsonl",
"/opt/airflow/data/landing")
print(f"extracted {n} orders for {ds}")
return path
with DAG("booknest_classic", schedule=None,
start_date=datetime(2026, 6, 1, tzinfo=timezone.utc), tags=["ch07"]):
extract = PythonOperator(task_id="extract", python_callable=extract_orders)
count = SQLExecuteQueryOperator(
task_id="count_orders", conn_id="booknest_pg", show_return_value_in_logs=True,
sql="SELECT count(*) FROM orders WHERE order_ts::date = '{{ ds }}'")
extract >> countairflow dags test runs it once in-process, without the scheduler (dag.test() and Local DAG Runs); the XComs show the result:
$ airflow dags test booknest_classic 2026-06-29
...
extracted 192 orders for 2026-06-29
... DagRun Finished: dag_id=booknest_classic, logical_date=2026-06-29 00:00:00+00:00...
extract | "/opt/airflow/data/landing/orders/date=2026-06-29/orders.jsonl"
count_orders | [{"__data__": [192], "__version__": 1, "__classname__": "builtins.tuple"}]The SQL is a Jinja 15,439 template: {{ ds }} renders per run in an operator's template_fields, while TaskFlow functions simply receive ds. TaskFlow infers dependencies from data flow and suits your own Python; operators need explicit >> and xcom_pull() but bring tested integrations.