Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ba4743ed52 |
@@ -13,6 +13,7 @@ GENERATOR_SPEED ?=
|
|||||||
up:
|
up:
|
||||||
$(COMPOSE) up --detach --build --wait --wait-timeout 600
|
$(COMPOSE) up --detach --build --wait --wait-timeout 600
|
||||||
COMPOSE_BIN="$(COMPOSE)" ./scripts/wait-for-world.sh
|
COMPOSE_BIN="$(COMPOSE)" ./scripts/wait-for-world.sh
|
||||||
|
COMPOSE_BIN="$(COMPOSE)" ./scripts/ingest-starting-orders.sh
|
||||||
|
|
||||||
down:
|
down:
|
||||||
$(COMPOSE) down --remove-orphans
|
$(COMPOSE) down --remove-orphans
|
||||||
|
|||||||
@@ -126,13 +126,20 @@ Kafka по той же причине спрашивают снаружи. Её
|
|||||||
|
|
||||||
Стенд поднимается не пустым: разовая служба `world-init` играет в топик `hits`
|
Стенд поднимается не пустым: разовая служба `world-init` играет в топик `hits`
|
||||||
первые восемь дней модельного времени — понедельник по понедельник, 401 185
|
первые восемь дней модельного времени — понедельник по понедельник, 401 185
|
||||||
событий. Дальше их обычным путём разбирает хранилище, и к концу `make up` они
|
событий. Следом второй запуск того же генератора отправляет в топик `orders`
|
||||||
лежат в `ods.event`. Так у всякой лабы есть данные, и всегда одни и те же.
|
семь слепков — за дни 0…6. Дальше оба источника разбирает хранилище, и к концу
|
||||||
|
`make up` события лежат в `ods.event`, а заказы — в `ods.order_v`. Так у всякой
|
||||||
|
лабы есть данные, и всегда одни и те же.
|
||||||
|
|
||||||
Ждать приходится дольше, чем работает заливка: приём асинхронный, поэтому
|
Последний день намеренно неполон: события дня 7 уже приехали, а его слепок
|
||||||
вторым шагом `make up` зовёт `scripts/wait-for-world.sh` — тот опрашивает
|
уедет только с первым ходом мира. Поэтому на свежем стенде правый край графика
|
||||||
ClickHouse, пока мир не доедет. Повторный `make up` заливает мир заново; это
|
выручки отстаёт от кликстрима и дозаполняется по мере движения мира — так
|
||||||
не ошибка, а свойство: номера событий те же, и повтор схлопнет
|
виден разный темп потокового трекера и ночной выгрузки магазина.
|
||||||
|
|
||||||
|
Ждать приходится дольше, чем работает заливка. `make up` сначала опрашивает
|
||||||
|
ClickHouse, пока события не доедут, затем разово запускает даг `orders_ingest`
|
||||||
|
и ждёт его конца. Повторный `make up` заливает мир заново; это не ошибка, а
|
||||||
|
свойство: номера событий и версии заказов те же, и повторы схлопнут
|
||||||
`ReplacingMergeTree`.
|
`ReplacingMergeTree`.
|
||||||
|
|
||||||
Сам мир в git не хранится — он чистая функция зерна, и держать его в
|
Сам мир в git не хранится — он чистая функция зерна, и держать его в
|
||||||
|
|||||||
+24
-6
@@ -36,11 +36,19 @@ x-clickhouse-common: &clickhouse-common
|
|||||||
# пульт мира. Названы по одному разу — иначе однажды разойдутся, и заметит это
|
# пульт мира. Названы по одному разу — иначе однажды разойдутся, и заметит это
|
||||||
# не проверка, а менти с пустым топиком.
|
# не проверка, а менти с пустым топиком.
|
||||||
x-generator-image: &generator-image clickstream-generator:local
|
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
|
x-kafka-target: &kafka-target
|
||||||
KAFKA_BOOTSTRAP_SERVERS: kafka:9092
|
<<: *kafka-common
|
||||||
KAFKA_TOPIC: hits
|
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-world-starting-days: &world-starting-days "8"
|
||||||
|
|
||||||
x-airflow-common: &airflow-common
|
x-airflow-common: &airflow-common
|
||||||
@@ -71,7 +79,7 @@ x-airflow-common: &airflow-common
|
|||||||
GENERATOR_IMAGE: *generator-image
|
GENERATOR_IMAGE: *generator-image
|
||||||
STAND_NETWORK: ${COMPOSE_PROJECT_NAME}_default
|
STAND_NETWORK: ${COMPOSE_PROJECT_NAME}_default
|
||||||
WORLD_STARTING_DAYS: *world-starting-days
|
WORLD_STARTING_DAYS: *world-starting-days
|
||||||
KAFKA_ORDERS_TOPIC: orders
|
KAFKA_ORDERS_TOPIC: *orders-topic
|
||||||
volumes:
|
volumes:
|
||||||
- ./dags:/opt/airflow/dags:ro
|
- ./dags:/opt/airflow/dags:ro
|
||||||
# Запросы дагов: даг называет файл, а текст читает Airflow при исполнении
|
# Запросы дагов: даг называет файл, а текст читает Airflow при исполнении
|
||||||
@@ -307,6 +315,17 @@ services:
|
|||||||
<<: *generator-common
|
<<: *generator-common
|
||||||
command: ["batch", "--day", "0", "--days", *world-starting-days]
|
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`.
|
# Тот же образ для ручных прогонов: `make generate-batch`, `make generate-live`.
|
||||||
# Под профилем — чтобы обычный подъём стенда её не трогал.
|
# Под профилем — чтобы обычный подъём стенда её не трогал.
|
||||||
generator:
|
generator:
|
||||||
@@ -355,10 +374,9 @@ services:
|
|||||||
# упавшим для --wait.
|
# упавшим для --wait.
|
||||||
clickhouse-init:
|
clickhouse-init:
|
||||||
condition: service_completed_successfully
|
condition: service_completed_successfully
|
||||||
# По той же причине. Заливка стартового мира идёт последней в цепи
|
# По той же причине. Слепки идут за событиями и замыкают цепь разовых
|
||||||
# разовых служб, и зависимого ей взять больше негде: долгоживущие службы
|
# служб; долгоживущие службы данных её не ждут, а `airflow-init` ждёт.
|
||||||
# данных не ждут, а `airflow-init` эту цепь и так замыкает.
|
world-snapshots:
|
||||||
world-init:
|
|
||||||
condition: service_completed_successfully
|
condition: service_completed_successfully
|
||||||
entrypoint: ["/bin/bash"]
|
entrypoint: ["/bin/bash"]
|
||||||
command: ["/opt/airflow/init.sh"]
|
command: ["/opt/airflow/init.sh"]
|
||||||
|
|||||||
@@ -71,9 +71,9 @@
|
|||||||
тогда окружение окупится, и условие возврата — именно это, а не «стало
|
тогда окружение окупится, и условие возврата — именно это, а не «стало
|
||||||
неудобно».
|
неудобно».
|
||||||
|
|
||||||
Цена самого `make up` — 2 м 50 с с нуля (`make clean` перед ним — ещё 23 с) и
|
Цена самого `make up` — 8 м 18 с с нуля и 3 м 35 с на живом стенде. В цену
|
||||||
1 м 5 с на живом стенде. В цену целей она не входит, но записана здесь по той
|
целей она не входит, но записана здесь по той же причине: подъём пересчитывает
|
||||||
же причине: с #42 подъём стенда заливает стартовый мир и стал заметно дороже.
|
оба источника стартового мира и ждёт их приёма.
|
||||||
|
|
||||||
## Правило: смоук обязан оставаться быстрым
|
## Правило: смоук обязан оставаться быстрым
|
||||||
|
|
||||||
@@ -194,12 +194,12 @@ Kafka → STG → ODS, и живёт он в `make check-clickhouse`.
|
|||||||
| `make check-services` | `scripts/stand-services.sh` |
|
| `make check-services` | `scripts/stand-services.sh` |
|
||||||
| `make check-clickhouse` | `scripts/check-clickhouse.sh` |
|
| `make check-clickhouse` | `scripts/check-clickhouse.sh` |
|
||||||
|
|
||||||
Особняком — `scripts/wait-for-world.sh`: он не проверка, а вторая половина
|
Особняком — два сценария продолжения `make up`, оба не проверки.
|
||||||
`make up`. `docker compose up --wait` дожидается служб, в том числе успешно
|
`scripts/wait-for-world.sh` ждёт, пока события доедут из Kafka в ODS:
|
||||||
отработавшей заливки, но «заливка кончилась» значит «события в топике», а не
|
`docker compose up --wait` видит конец отправки, но приём идёт асинхронно.
|
||||||
«события в хранилище»: приём асинхронный. Скрипт ждёт доезда ограниченным
|
Затем `scripts/ingest-starting-orders.sh` разово запускает единственный даг
|
||||||
циклом опроса. Без него `make check-clickhouse` следом краснел бы по
|
приёма заказов `orders_ingest` и ждёт его конца. Оба ждут ограниченным циклом,
|
||||||
устройству, а не по поломке.
|
а не паузой наугад.
|
||||||
|
|
||||||
Общее у смоука и `check-services` — счёт проверок, обращение к Compose и две
|
Общее у смоука и `check-services` — счёт проверок, обращение к Compose и две
|
||||||
проверки — вынесено в `scripts/stand-common.sh`; сам он не запускается.
|
проверки — вынесено в `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.** После добавления десятой
|
**Перезамер 19 августа 2026 года при исполнении #95.** После добавления десятой
|
||||||
проверки словаря `make check-clickhouse` на чисто поднятом стенде прошёл за
|
проверки словаря `make check-clickhouse` на чисто поднятом стенде прошёл за
|
||||||
24 секунды. Это цена всей цели, а не одной новой строки.
|
24 секунды. Это цена всей цели, а не одной новой строки.
|
||||||
|
|||||||
Executable
+94
@@ -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}"
|
||||||
Reference in New Issue
Block a user