- Добавлен kafka-python==2.0.6 в airflow/requirements.txt - Создан dags/utils/kafka_helpers.py с функциями: - check_kafka_ready() — проверка доступности брокера - prepare_topics() — создание/сброс топиков через KafkaAdminClient - load_jsonl() — загрузка данных через KafkaProducer (limit=0 = все) - validate_load_params(), check_input_files() — валидация - Создан dags/kafka_load_dag.py с TaskGroup: - precheck: check_kafka, check_input_files, validate_load_params - ingest: prepare_topics, параллельная загрузка 4 потоков, verify_publish_counts - Параметры DAG: limit (0 = все), reset_topics, load_* (выбор потоков) - Обновлена документация: AGENTS.md, README.md, plans/runbook.md, plans/airflow_dags_plan.md, docs/ARCHITECTURE.md Тестирование: - Подключение к Kafka: ✅ (kafka:29092 доступен, брокер 2.6.0) - Загрузка данных: ✅ (1000 сообщений — полный файл browser_events) - Python синтаксис: ✅ (py_compile проходит) - Структура DAG: ✅ (все 9 задач корректно определены)
11 lines
286 B
Plaintext
11 lines
286 B
Plaintext
# Airflow requirements для учебного ETL-проекта
|
|
|
|
# Metadata DB для Airflow
|
|
psycopg2-binary==2.9.9
|
|
|
|
# ClickHouse operator/hook для DAG'ов
|
|
airflow-clickhouse-plugin==1.6.0
|
|
|
|
# Kafka client для загрузки данных (DAG kafka_load)
|
|
kafka-python==2.0.6
|