feat(stand): добавлены стартовые слепки заказов

- Зачем:
  - свежий стенд должен показывать отставание ночной выгрузки заказов от кликстрима.
- Что:
  - после событий генератор отправляет семь слепков дней 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.
This commit is contained in:
2026-08-20 19:12:01 +03:00
parent 664abf0160
commit ba4743ed52
5 changed files with 150 additions and 21 deletions
+1
View File
@@ -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
+13 -6
View File
@@ -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 не хранится — он чистая функция зерна, и держать его в
+24 -6
View File
@@ -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"]
+18 -9
View File
@@ -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 секунды. Это цена всей цели, а не одной новой строки.
+94
View File
@@ -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}"