- Why:
- keep Airflow artifacts under a single airflow/ directory
- align repository layout with intended project structure
- What:
- move dags/ to airflow/dags/ and update compose mounts
- make SQL root resolution work in container and local runs
- update DAG path references in README, AGENTS, ARCHITECTURE, and plans
- remove tracked Python cache artifacts from old DAG location
- Check:
- airflow dags list
- airflow dags list-import-errors
- e2e success: ddl_init, kafka_load(limit=50), etl_pipeline
- Why:
- For DE task we only need full ingest or limit-based sample.
- load_* and full_load params were redundant and unclear in current flow.
- What:
- Remove full_load and load_* params from kafka_load DAG contract.
- Simplify kafka helpers (validate/check files) to fixed 4-stream ingest.
- Sync AGENTS, README, runbook, architecture and airflow plan docs.
- Check:
- python3 -m py_compile dags/kafka_load_dag.py dags/utils/kafka_helpers.py
- Airflow smoke/full runs: ddl_init -> kafka_load -> etl_pipeline (all success).
- Legacy path: make data && make transform (success).
Move DDL files from flat ddl/ directory to sql/ddl/ with layer-based
subdirectories (stg, ods, dds, dm). Move batch transformation SQL from
jobs/ to sql/ layer directories. Update scripts and documentation to
reflect new paths for improved organization and Airflow integration.
Split Kafka ingestion into a dedicated `kafka_load` DAG to enable independent
experimentation with data loading without triggering the full ETL pipeline.
Restructure implementation phases: Stage 1 uses `make data` for MVP, Stage 2
adds the standalone Kafka DAG, Stage 3 adds monitoring. Update DAG numbering,
parameters, task groups, and acceptance criteria to reflect the new
architecture.
Consolidate the orchestration strategy by merging `kafka_load` and
`etl_batch_transform` into a unified `etl_pipeline`. Replace BashOperator
dependencies on Kafka CLI with PythonOperators utilizing `kafka-python`.
Add detailed technical specifications for helper functions, MVP stages,
and validation checks to align with current infrastructure constraints.
Detail the architecture for migrating ETL orchestration from make to
Airflow. Define DAG structures for database initialization, Kafka data
ingestion, batch transformation, and quality monitoring. Include
technical specifications, operator details, and implementation phases.