Пробные DAG: test_clickhouse и test_kafka вместо демонстрационного #18

Closed
opened 2026-07-31 09:32:15 +03:00 by ddmitry · 0 comments
Owner

Part of #3

Цель

Заменить демонстрационный DAG двумя интеграционными пробниками — test_clickhouse и test_kafka, — чтобы Airflow проверял не сам себя, а настоящие подключения стенда.

example_clickstream_hello доказывает только, что планировщик жив и задачи передают друг другу данные. Подключение clickhouse_default и путь от Airflow до Kafka сейчас не проверяет никто: scripts/stand-smoke.sh лишь дёргает демонстрационный DAG, а доставку ON CLUSTER до обеих нод проверяет отдельный scripts/clickhouse-smoke.sh — снаружи Airflow. Пробник проверяет подключения тем же способом, каким ими будет пользоваться будущий пайплайн (#7), и заодно даёт образец обращения к кластеру из кода.

Главная проверка — сквозная по кластеру: записать данные на одну ноду, а прочитать их с другой через Distributed. Так кластер и задуман работать, и только такая проверка отличает настоящий кластер от двух независимых баз, случайно оказавшихся в одном compose.yaml.

Ловушку разных нод в её полном виде (забытый ON CLUSTER проявляется в дашборде Superset) пробник не воспроизводит: это остаётся ручным сценарием из README. Но забытую доставку таблицы он поймает.

Что уже выяснено на живом стенде 31 июля 2026 года

Перепроверять не нужно, но и игнорировать нельзя — постановка на этом построена.

  • В образе apache/airflow:3.3.0 нет ни clickhouse-connect, ни клиента Kafka. Официального провайдера ClickHouse у Airflow не существует (сверено через Context7); сторонний плагин брать не нужно.
  • Топология кластера — два шарда по одной реплике. Строка, записанная в локальную таблицу на ноде 1, в локальной таблице ноды 2 не появится: это разные шарды, а не копии. Прочитать её с ноды 2 можно только через таблицу Distributed — она опрашивает все шарды. На этом и построена главная проверка.
  • Изнутри контейнера Airflow брокер доступен только как kafka:9092. Порт машины 29092, который используется в scripts/stand-smoke.sh, из контейнера не работает — не скопируйте его по ошибке.
  • У брокера auto.create.topics.enable=true и log.retention.hours=168, оба по умолчанию и в compose.yaml не переопределены. Срок хранения удаляет записи, а не топик: сам топик живёт вечно.
  • tests/stand-smoke-guards.sh запускает make smoke три раза подряд.
  • Запас памяти до порога 3,4 ГБ в scripts/stand-smoke.sh — около 150 MiB, и он колеблется.

Критерии приёмки

Зависимости и образ:

  • Клиенты ClickHouse и Kafka добавлены в Airflow своим образом infra/airflow/Dockerfile по образцу infra/superset/Dockerfile; версии закреплены точно. _PIP_ADDITIONAL_REQUIREMENTS не использовать: он ставит пакеты при каждом старте и официально не предназначен для боя.
  • Пакеты лежат в образе, а не доустанавливаются при старте. Проверяемый признак: docker run --rm <образ> python -c "import clickhouse_connect" (и то же для клиента Kafka) отрабатывает без сети и без установки, а в журнале запуска контейнеров нет ни одной строки установки пакетов.
  • compose.yaml и .env.example согласованы с новым образом; make up по-прежнему поднимает всё одной командой.

test_clickhouse:

  • Работает через существующее подключение clickhouse_default (нода 1). Второго подключения Airflow не заводить.
  • Создаёт ON CLUSTER пару служебных таблиц — локальную и Distributed поверх неё, по образцу scripts/clickhouse-smoke.sh.
  • Проверяет, что обе таблицы появились на обеих нодах, а не только там, где выполнялся DDL.
  • Пишет строку с уникальным для прогона маркером в локальную таблицу на ноде 1.
  • Читает эту строку с ноды 2 через Distributed — данные проходят путь «нода 2 → все шарды → шард ноды 1». Это главная проверка тикета: она доказывает, что кластер работает как кластер. Одного подключения для неё достаточно, например через remote('clickhouse-02:9000', ...) с ноды 1; способ на усмотрение исполнителя, важно, чтобы чтение выполняла именно нода 2.
  • Ищет именно свой маркер: проверка «есть хоть какая-то строка» не годится.
  • Убирает обе таблицы ON CLUSTER за собой и убеждается, что их не осталось ни на одной ноде.
  • Падает, если таблица доехала не до обеих нод или маркер не нашёлся при чтении с ноды 2: молча зелёным не проходит.

test_kafka:

  • Пишет сообщение в топик с постоянным именем и читает обратно именно своё сообщение — по уникальному для прогона маркеру и с ограничением по времени, чтобы не висеть вечно.
  • Топик не удаляется. В коде одной строкой отмечено: стенд опирается на автосоздание топиков, срок хранения чистит записи, а не топик, и в бою автосоздание обычно выключают.
  • Адрес брокера — kafka:9092.

Общее:

  • Оба пробника переживают запуск подряд: остатки прошлого прогона не мешают следующему.
  • Клиенты создаются внутри задач, а не на верхнем уровне файла DAG: обработчик перечитывает файлы каждый цикл, и объекты верхнего уровня осели бы в памяти трёх контейнеров Airflow навсегда.
  • scripts/stand-smoke.sh запускает оба пробника вместо example_clickstream_hello; dags/example_clickstream_hello.py удалён, упоминания в README.md и docs/adr/0001-stand-services.md обновлены.
  • В tests/stand-smoke-guards.sh добавлен случай, доказывающий красный путь: после намеренной поломки (например, служебная таблица удалена только на одной ноде) make smoke падает и называет пробник. Ломать для этого можно — так уже сделано с Prometheus; стенд в конце восстанавливается.
  • make up и make smoke зелены на чистых томах, расход памяти ниже порога 3,4 ГБ. Если порог не выдерживается — не поднимать его молча, а вернуть вопрос заказчику.
  • README.md и docs/adr/0001-stand-services.md обновлены тем же PR.

Границы

  • Это не etl_pipeline (#7): пробники каркаса, без слоёв, партиций, сенсоров и расписания.
  • Не дашборды: Superset — #10, панели Grafana — #11.
  • Модель данных кликстрима не трогаем (#4): таблицы пробника служебные и временные.
  • Второго подключения Airflow к ноде 2 не заводить: чтение с ноды 2 выполнять в рамках существующего clickhouse_default. Ловушку разных нод в её полном виде (Superset) не автоматизировать.
  • Порог памяти в scripts/stand-smoke.sh не менять.

Сначала прочитать

  • docs/adr/0001-stand-services.md — роли нод и решения по стенду;
  • scripts/clickhouse-smoke.sh — уже написанная проверка ON CLUSTER на обеих нодах, ближайший образец;
  • scripts/stand-smoke.sh — раздел проверок Airflow (запуск DAG через API) и измерение памяти;
  • infra/superset/Dockerfile — образец своего образа поверх официального;
  • infra/airflow/init.sh — как заводится подключение clickhouse_default;
  • dags/example_clickstream_hello.py — то, что заменяется.

Спорные API сверить через Context7 до кода: как в Airflow 3.3 правильно выполнять запросы к ClickHouse и какой брать клиент Kafka (официальный провайдер apache-kafka против голого confluent-kafka/kafka-python).

Проверка

make config-test
make clean
cp .env.example .env
make up
make smoke
make smoke-guards
make smoke-cluster
Part of #3 ## Цель Заменить демонстрационный DAG двумя интеграционными пробниками — `test_clickhouse` и `test_kafka`, — чтобы Airflow проверял не сам себя, а настоящие подключения стенда. `example_clickstream_hello` доказывает только, что планировщик жив и задачи передают друг другу данные. Подключение `clickhouse_default` и путь от Airflow до Kafka сейчас не проверяет никто: `scripts/stand-smoke.sh` лишь дёргает демонстрационный DAG, а доставку `ON CLUSTER` до обеих нод проверяет отдельный `scripts/clickhouse-smoke.sh` — снаружи Airflow. Пробник проверяет подключения тем же способом, каким ими будет пользоваться будущий пайплайн (#7), и заодно даёт образец обращения к кластеру из кода. Главная проверка — сквозная по кластеру: записать данные на одну ноду, а прочитать их с другой через `Distributed`. Так кластер и задуман работать, и только такая проверка отличает настоящий кластер от двух независимых баз, случайно оказавшихся в одном `compose.yaml`. Ловушку разных нод в её полном виде (забытый `ON CLUSTER` проявляется в дашборде Superset) пробник не воспроизводит: это остаётся ручным сценарием из README. Но забытую доставку таблицы он поймает. ## Что уже выяснено на живом стенде 31 июля 2026 года Перепроверять не нужно, но и игнорировать нельзя — постановка на этом построена. - В образе `apache/airflow:3.3.0` нет ни `clickhouse-connect`, ни клиента Kafka. Официального провайдера ClickHouse у Airflow не существует (сверено через Context7); сторонний плагин брать не нужно. - Топология кластера — **два шарда по одной реплике**. Строка, записанная в локальную таблицу на ноде 1, в локальной таблице ноды 2 не появится: это разные шарды, а не копии. Прочитать её с ноды 2 можно только через таблицу `Distributed` — она опрашивает все шарды. На этом и построена главная проверка. - Изнутри контейнера Airflow брокер доступен только как `kafka:9092`. Порт машины `29092`, который используется в `scripts/stand-smoke.sh`, из контейнера не работает — не скопируйте его по ошибке. - У брокера `auto.create.topics.enable=true` и `log.retention.hours=168`, оба по умолчанию и в `compose.yaml` не переопределены. Срок хранения удаляет **записи**, а не топик: сам топик живёт вечно. - `tests/stand-smoke-guards.sh` запускает `make smoke` три раза подряд. - Запас памяти до порога 3,4 ГБ в `scripts/stand-smoke.sh` — около 150 MiB, и он колеблется. ## Критерии приёмки Зависимости и образ: - [x] Клиенты ClickHouse и Kafka добавлены в Airflow своим образом `infra/airflow/Dockerfile` по образцу `infra/superset/Dockerfile`; версии закреплены точно. `_PIP_ADDITIONAL_REQUIREMENTS` не использовать: он ставит пакеты при каждом старте и официально не предназначен для боя. - [x] Пакеты лежат в образе, а не доустанавливаются при старте. Проверяемый признак: `docker run --rm <образ> python -c "import clickhouse_connect"` (и то же для клиента Kafka) отрабатывает без сети и без установки, а в журнале запуска контейнеров нет ни одной строки установки пакетов. - [x] `compose.yaml` и `.env.example` согласованы с новым образом; `make up` по-прежнему поднимает всё одной командой. `test_clickhouse`: - [x] Работает через существующее подключение `clickhouse_default` (нода 1). Второго подключения Airflow не заводить. - [x] Создаёт `ON CLUSTER` пару служебных таблиц — локальную и `Distributed` поверх неё, по образцу `scripts/clickhouse-smoke.sh`. - [x] Проверяет, что обе таблицы появились **на обеих нодах**, а не только там, где выполнялся DDL. - [x] Пишет строку с уникальным для прогона маркером **в локальную таблицу на ноде 1**. - [x] **Читает эту строку с ноды 2 через `Distributed`** — данные проходят путь «нода 2 → все шарды → шард ноды 1». Это главная проверка тикета: она доказывает, что кластер работает как кластер. Одного подключения для неё достаточно, например через `remote('clickhouse-02:9000', ...)` с ноды 1; способ на усмотрение исполнителя, важно, чтобы чтение выполняла именно нода 2. - [x] Ищет **именно свой** маркер: проверка «есть хоть какая-то строка» не годится. - [x] Убирает обе таблицы `ON CLUSTER` за собой и убеждается, что их не осталось ни на одной ноде. - [x] Падает, если таблица доехала не до обеих нод или маркер не нашёлся при чтении с ноды 2: молча зелёным не проходит. `test_kafka`: - [x] Пишет сообщение в топик с постоянным именем и читает обратно **именно своё** сообщение — по уникальному для прогона маркеру и с ограничением по времени, чтобы не висеть вечно. - [x] Топик не удаляется. В коде одной строкой отмечено: стенд опирается на автосоздание топиков, срок хранения чистит записи, а не топик, и в бою автосоздание обычно выключают. - [x] Адрес брокера — `kafka:9092`. Общее: - [x] Оба пробника переживают запуск подряд: остатки прошлого прогона не мешают следующему. - [x] Клиенты создаются внутри задач, а не на верхнем уровне файла DAG: обработчик перечитывает файлы каждый цикл, и объекты верхнего уровня осели бы в памяти трёх контейнеров Airflow навсегда. - [x] `scripts/stand-smoke.sh` запускает оба пробника вместо `example_clickstream_hello`; `dags/example_clickstream_hello.py` удалён, упоминания в `README.md` и `docs/adr/0001-stand-services.md` обновлены. - [x] В `tests/stand-smoke-guards.sh` добавлен случай, доказывающий красный путь: после намеренной поломки (например, служебная таблица удалена только на одной ноде) `make smoke` падает и называет пробник. Ломать для этого можно — так уже сделано с Prometheus; стенд в конце восстанавливается. - [x] `make up` и `make smoke` зелены на чистых томах, расход памяти ниже порога 3,4 ГБ. Если порог не выдерживается — не поднимать его молча, а вернуть вопрос заказчику. - [x] `README.md` и `docs/adr/0001-stand-services.md` обновлены тем же PR. ## Границы - Это не `etl_pipeline` (#7): пробники каркаса, без слоёв, партиций, сенсоров и расписания. - Не дашборды: Superset — #10, панели Grafana — #11. - Модель данных кликстрима не трогаем (#4): таблицы пробника служебные и временные. - Второго подключения Airflow к ноде 2 не заводить: чтение с ноды 2 выполнять в рамках существующего `clickhouse_default`. Ловушку разных нод в её полном виде (Superset) не автоматизировать. - Порог памяти в `scripts/stand-smoke.sh` не менять. ## Сначала прочитать - `docs/adr/0001-stand-services.md` — роли нод и решения по стенду; - `scripts/clickhouse-smoke.sh` — уже написанная проверка `ON CLUSTER` на обеих нодах, ближайший образец; - `scripts/stand-smoke.sh` — раздел проверок Airflow (запуск DAG через API) и измерение памяти; - `infra/superset/Dockerfile` — образец своего образа поверх официального; - `infra/airflow/init.sh` — как заводится подключение `clickhouse_default`; - `dags/example_clickstream_hello.py` — то, что заменяется. Спорные API сверить через Context7 до кода: как в Airflow 3.3 правильно выполнять запросы к ClickHouse и какой брать клиент Kafka (официальный провайдер `apache-kafka` против голого `confluent-kafka`/`kafka-python`). ## Проверка ``` make config-test make clean cp .env.example .env make up make smoke make smoke-guards make smoke-cluster ```
ddmitry added the ready-for-agent label 2026-07-31 09:32:15 +03:00
Sign in to join this conversation.
1 Participants
Notifications
Due Date
No due date set.
Dependencies

No dependencies set.

Reference: ddmitry/clickstream-data-platform#18