From ba4743ed520ae344b4debf861123a88e2cddc621 Mon Sep 17 00:00:00 2001 From: Dmitry Dementiev Date: Thu, 20 Aug 2026 16:32:36 +0300 Subject: [PATCH] =?UTF-8?q?feat(stand):=20=D0=B4=D0=BE=D0=B1=D0=B0=D0=B2?= =?UTF-8?q?=D0=BB=D0=B5=D0=BD=D1=8B=20=D1=81=D1=82=D0=B0=D1=80=D1=82=D0=BE?= =?UTF-8?q?=D0=B2=D1=8B=D0=B5=20=D1=81=D0=BB=D0=B5=D0=BF=D0=BA=D0=B8=20?= =?UTF-8?q?=D0=B7=D0=B0=D0=BA=D0=B0=D0=B7=D0=BE=D0=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - свежий стенд должен показывать отставание ночной выгрузки заказов от кликстрима. - Что: - после событий генератор отправляет семь слепков дней 0…6 в общий топик orders. - make up запускает orders_ingest, ограниченно ждёт его завершения и оставляет день 7 первому ходу мира. - README и карта проверок описывают второй источник и измеренную цену подъёма. - Проверка: - make clean && make up; make smoke; make check-clickhouse; make check-services. - cd generator && make lint && make typecheck && make test. --- Makefile | 1 + README.md | 19 +++++-- compose.yaml | 30 ++++++++-- docs/architecture/testing.md | 27 ++++++--- scripts/ingest-starting-orders.sh | 94 +++++++++++++++++++++++++++++++ 5 files changed, 150 insertions(+), 21 deletions(-) create mode 100755 scripts/ingest-starting-orders.sh diff --git a/Makefile b/Makefile index 2f5b8a4..f340831 100644 --- a/Makefile +++ b/Makefile @@ -13,6 +13,7 @@ GENERATOR_SPEED ?= up: $(COMPOSE) up --detach --build --wait --wait-timeout 600 COMPOSE_BIN="$(COMPOSE)" ./scripts/wait-for-world.sh + COMPOSE_BIN="$(COMPOSE)" ./scripts/ingest-starting-orders.sh down: $(COMPOSE) down --remove-orphans diff --git a/README.md b/README.md index 68f5f1d..b0b3b08 100644 --- a/README.md +++ b/README.md @@ -126,13 +126,20 @@ Kafka по той же причине спрашивают снаружи. Её Стенд поднимается не пустым: разовая служба `world-init` играет в топик `hits` первые восемь дней модельного времени — понедельник по понедельник, 401 185 -событий. Дальше их обычным путём разбирает хранилище, и к концу `make up` они -лежат в `ods.event`. Так у всякой лабы есть данные, и всегда одни и те же. +событий. Следом второй запуск того же генератора отправляет в топик `orders` +семь слепков — за дни 0…6. Дальше оба источника разбирает хранилище, и к концу +`make up` события лежат в `ods.event`, а заказы — в `ods.order_v`. Так у всякой +лабы есть данные, и всегда одни и те же. -Ждать приходится дольше, чем работает заливка: приём асинхронный, поэтому -вторым шагом `make up` зовёт `scripts/wait-for-world.sh` — тот опрашивает -ClickHouse, пока мир не доедет. Повторный `make up` заливает мир заново; это -не ошибка, а свойство: номера событий те же, и повтор схлопнет +Последний день намеренно неполон: события дня 7 уже приехали, а его слепок +уедет только с первым ходом мира. Поэтому на свежем стенде правый край графика +выручки отстаёт от кликстрима и дозаполняется по мере движения мира — так +виден разный темп потокового трекера и ночной выгрузки магазина. + +Ждать приходится дольше, чем работает заливка. `make up` сначала опрашивает +ClickHouse, пока события не доедут, затем разово запускает даг `orders_ingest` +и ждёт его конца. Повторный `make up` заливает мир заново; это не ошибка, а +свойство: номера событий и версии заказов те же, и повторы схлопнут `ReplacingMergeTree`. Сам мир в git не хранится — он чистая функция зерна, и держать его в diff --git a/compose.yaml b/compose.yaml index a8b2218..d8a208f 100644 --- a/compose.yaml +++ b/compose.yaml @@ -36,11 +36,19 @@ x-clickhouse-common: &clickhouse-common # пульт мира. Названы по одному разу — иначе однажды разойдутся, и заметит это # не проверка, а менти с пустым топиком. x-generator-image: &generator-image clickstream-generator:local +x-orders-topic: &orders-topic orders + +x-kafka-common: &kafka-common + KAFKA_BOOTSTRAP_SERVERS: kafka:9092 x-kafka-target: &kafka-target - KAFKA_BOOTSTRAP_SERVERS: kafka:9092 + <<: *kafka-common KAFKA_TOPIC: hits +x-orders-kafka-target: &orders-kafka-target + <<: *kafka-common + KAFKA_TOPIC: *orders-topic + x-world-starting-days: &world-starting-days "8" x-airflow-common: &airflow-common @@ -71,7 +79,7 @@ x-airflow-common: &airflow-common GENERATOR_IMAGE: *generator-image STAND_NETWORK: ${COMPOSE_PROJECT_NAME}_default WORLD_STARTING_DAYS: *world-starting-days - KAFKA_ORDERS_TOPIC: orders + KAFKA_ORDERS_TOPIC: *orders-topic volumes: - ./dags:/opt/airflow/dags:ro # Запросы дагов: даг называет файл, а текст читает Airflow при исполнении @@ -307,6 +315,17 @@ services: <<: *generator-common command: ["batch", "--day", "0", "--days", *world-starting-days] + # Второй источник стартового мира идёт после событий тем же генератором. + # Сдвиг команды snapshot оставляет день 7 первому ходу мира: за восемь + # дней уезжают семь слепков, дней 0…6. + world-snapshots: + <<: *generator-common + depends_on: + world-init: + condition: service_completed_successfully + environment: *orders-kafka-target + command: ["snapshot", "--day", "0", "--days", *world-starting-days] + # Тот же образ для ручных прогонов: `make generate-batch`, `make generate-live`. # Под профилем — чтобы обычный подъём стенда её не трогал. generator: @@ -355,10 +374,9 @@ services: # упавшим для --wait. clickhouse-init: condition: service_completed_successfully - # По той же причине. Заливка стартового мира идёт последней в цепи - # разовых служб, и зависимого ей взять больше негде: долгоживущие службы - # данных не ждут, а `airflow-init` эту цепь и так замыкает. - world-init: + # По той же причине. Слепки идут за событиями и замыкают цепь разовых + # служб; долгоживущие службы данных её не ждут, а `airflow-init` ждёт. + world-snapshots: condition: service_completed_successfully entrypoint: ["/bin/bash"] command: ["/opt/airflow/init.sh"] diff --git a/docs/architecture/testing.md b/docs/architecture/testing.md index 65ad23f..1145354 100644 --- a/docs/architecture/testing.md +++ b/docs/architecture/testing.md @@ -71,9 +71,9 @@ тогда окружение окупится, и условие возврата — именно это, а не «стало неудобно». -Цена самого `make up` — 2 м 50 с с нуля (`make clean` перед ним — ещё 23 с) и -1 м 5 с на живом стенде. В цену целей она не входит, но записана здесь по той -же причине: с #42 подъём стенда заливает стартовый мир и стал заметно дороже. +Цена самого `make up` — 8 м 18 с с нуля и 3 м 35 с на живом стенде. В цену +целей она не входит, но записана здесь по той же причине: подъём пересчитывает +оба источника стартового мира и ждёт их приёма. ## Правило: смоук обязан оставаться быстрым @@ -194,12 +194,12 @@ Kafka → STG → ODS, и живёт он в `make check-clickhouse`. | `make check-services` | `scripts/stand-services.sh` | | `make check-clickhouse` | `scripts/check-clickhouse.sh` | -Особняком — `scripts/wait-for-world.sh`: он не проверка, а вторая половина -`make up`. `docker compose up --wait` дожидается служб, в том числе успешно -отработавшей заливки, но «заливка кончилась» значит «события в топике», а не -«события в хранилище»: приём асинхронный. Скрипт ждёт доезда ограниченным -циклом опроса. Без него `make check-clickhouse` следом краснел бы по -устройству, а не по поломке. +Особняком — два сценария продолжения `make up`, оба не проверки. +`scripts/wait-for-world.sh` ждёт, пока события доедут из Kafka в ODS: +`docker compose up --wait` видит конец отправки, но приём идёт асинхронно. +Затем `scripts/ingest-starting-orders.sh` разово запускает единственный даг +приёма заказов `orders_ingest` и ждёт его конца. Оба ждут ограниченным циклом, +а не паузой наугад. Общее у смоука и `check-services` — счёт проверок, обращение к Compose и две проверки — вынесено в `scripts/stand-common.sh`; сам он не запускается. @@ -221,6 +221,15 @@ Kafka → STG → ODS, и живёт он в `make check-clickhouse`. ## Что проверено +**Перезамер 20 августа 2026 года при исполнении #96.** `make up` после +`make clean` прошёл за 8 м 18 с, повтор на живом стенде — за 3 м 35 с. Служба +семи слепков заняла около 28 секунд между концом `world-init` и своим концом. +Прежний замер ниже — 2 м 50 с и 1 м 5 с до появления приёма заказов и пульта +мира. Весь рост приписывать слепкам нельзя: на холодном прогоне основное время +ушло ещё до сценариев ожидания: `airflow-init` завершился примерно на 320-й +секунде, обработчик DAG стал здоров примерно на 458-й. Числа здесь называют +цену нынешнего стенда, а не выносят вердикт о регрессии. + **Перезамер 19 августа 2026 года при исполнении #95.** После добавления десятой проверки словаря `make check-clickhouse` на чисто поднятом стенде прошёл за 24 секунды. Это цена всей цели, а не одной новой строки. diff --git a/scripts/ingest-starting-orders.sh b/scripts/ingest-starting-orders.sh new file mode 100755 index 0000000..b816914 --- /dev/null +++ b/scripts/ingest-starting-orders.sh @@ -0,0 +1,94 @@ +#!/usr/bin/env bash +# Последний шаг `make up`: запускает единственный путь приёма заказов и ждёт +# его конца. Слепки к этому моменту уже лежат в Kafka, а события доехали до ODS. +set -Eeuo pipefail + +readonly ROOT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" +readonly DAG_ID='orders_ingest' +# Потолок каждого ожидания — 300 с: больше чем всемеро от замеренных 40 с +# между запуском сценария и концом дага (#96). Это предел, а не бюджет времени. +readonly ATTEMPTS=100 +readonly PAUSE_SECONDS=3 +readonly REPORT_EVERY=5 + +read -r -a COMPOSE_CMD <<<"${COMPOSE_BIN:-docker compose}" + +fail() { + printf 'ОШИБКА: %s\n' "$1" >&2 + exit 1 +} + +compose() { + "${COMPOSE_CMD[@]}" --project-directory "$ROOT_DIR" "$@" +} + +config="$(compose config --format json)" +user="$(jq -r '.services["airflow-apiserver"].environment.AIRFLOW_ADMIN_USER // empty' <<<"$config")" +password="$(jq -r '.services["airflow-apiserver"].environment.AIRFLOW_ADMIN_PASSWORD // empty' <<<"$config")" +binding="$(compose port airflow-apiserver 8080)" +port="${binding##*:}" +[[ -n "$port" ]] || fail 'не найден отображённый порт API Airflow' + +response="$(curl -sf --max-time 10 -X POST \ + -H 'Content-Type: application/json' \ + -d "$(jq -cn --arg username "$user" --arg password "$password" \ + '{username: $username, password: $password}')" \ + "http://127.0.0.1:${port}/auth/token" 2>/dev/null || true)" +token="$(jq -r '.access_token // empty' <<<"$response")" +[[ -n "$token" ]] || fail 'Airflow не принял учётные данные администратора' + +# После `up --wait` обработчик DAG здоров, но файл мог ещё не попасть в список. +# Ограниченный опрос закрывает эту гонку без паузы наугад. +dag='' +for ((attempt = 1; attempt <= ATTEMPTS; attempt++)); do + dag="$(curl -sf --max-time 10 \ + -H "Authorization: Bearer ${token}" \ + "http://127.0.0.1:${port}/api/v2/dags/${DAG_ID}" \ + 2>/dev/null || true)" + jq -e --arg dag_id "$DAG_ID" '.dag_id == $dag_id' \ + >/dev/null 2>&1 <<<"$dag" && break + sleep "$PAUSE_SECONDS" +done +jq -e --arg dag_id "$DAG_ID" '.dag_id == $dag_id' \ + >/dev/null 2>&1 <<<"$dag" || \ + fail "даг ${DAG_ID} не появился в Airflow за $((ATTEMPTS * PAUSE_SECONDS)) с" + +# В API v2 Airflow 3.3 logical_date=null означает событийный ручной запуск. +# Это уместно для дага без расписания: календарь Airflow не подменяет дату +# слепка, которая уже записана в сообщениях. +response="$(curl -sf --max-time 10 -X POST \ + -H "Authorization: Bearer ${token}" \ + -H 'Content-Type: application/json' \ + -d '{"logical_date":null}' \ + "http://127.0.0.1:${port}/api/v2/dags/${DAG_ID}/dagRuns" \ + 2>/dev/null || true)" +run_id="$(jq -r '.dag_run_id // empty' <<<"$response")" +[[ -n "$run_id" ]] || fail "Airflow не запустил даг ${DAG_ID}" + +encoded_run_id="$(jq -rn --arg value "$run_id" '$value | @uri')" +printf 'Ждём приём стартовых слепков дагом %s...\n' "$DAG_ID" +state='' +for ((attempt = 1; attempt <= ATTEMPTS; attempt++)); do + response="$(curl -sf --max-time 10 \ + -H "Authorization: Bearer ${token}" \ + "http://127.0.0.1:${port}/api/v2/dags/${DAG_ID}/dagRuns/${encoded_run_id}" \ + 2>/dev/null || true)" + state="$(jq -r '.state // empty' <<<"$response")" + case "$state" in + success) + printf 'ЗЕЛЁНО: стартовые слепки приняты в ODS.\n' + exit 0 + ;; + failed) + fail "даг ${DAG_ID} завершился с ошибкой; запуск ${run_id}" + ;; + esac + if ((attempt % REPORT_EVERY == 0)); then + printf 'Приём ещё не завершён: состояние %s, ждём дальше (%d с)...\n' \ + "${state:-неизвестно}" "$((attempt * PAUSE_SECONDS))" + fi + sleep "$PAUSE_SECONDS" +done + +fail "даг ${DAG_ID} не завершился за $((ATTEMPTS * PAUSE_SECONDS)) с: + состояние ${state:-неизвестно}, запуск ${run_id}"