Учебная дата-платформа кликстрима

Стенд для работы с кликстримом. Преемник учебного стенда clickstream-ch-kafka-superset-demo.

Статус

Репозиторий строится по спеке «Боевой реализм стенда (v2)». Сейчас работают кластер ClickHouse из двух шардов, отдельный clickhouse-keeper, односерверная Kafka в режиме KRaft, Airflow 3.3, Superset, Prometheus, Grafana и общая база Postgres для метаданных. Генератор кликстрима в generator/ собран целиком со стороны клиента: у него есть контракт схемы события, из которого собрано описание выгрузки, модельный мир и проигрыватель, отправляющий дни в Kafka или в файл. События доезжают до типизированного ods.event, а make up наполняет стенд стартовым миром — первой неделей модельного времени.

Быстрый старт

Нужны Docker с Compose и около 8 ГБ памяти, доступной Docker. Это не объём ноутбука, а то, что отдано самому Docker: в Docker Desktop и WSL2 он живёт внутри виртуальной машины и получает лишь часть памяти хозяина. Сколько выдано сейчас, в байтах, покажет docker info --format '{{.MemTotal}}'. В WSL2 это поднимается параметром memory в файле .wslconfig домашнего каталога пользователя Windows; после правки нужен wsl --shutdown. Если своей машины не хватает, стенд одинаково хорошо живёт на недорогом VPS.

Проверкам на поднятом стенде также нужны curl, jq, awk, grep, sed, tail, sleep и timeout. jq нужен и самому make up: им читается опись мира, по которой он ждёт заливки. По умолчанию должны быть свободны порты 23000, 28080, 28088, 28123, 28124, 29000, 29001, 29090 и 29092. Проверкам без стенда — make config-test и make lint в корне, всем целям генератора в generator/ — нужен uv.

Сначала скопируйте настройки стенда:

cp .env.example .env
make up
make smoke
make check-clickhouse
make check-services

Первому знакомству нужны все три проверки на стенде: make smoke говорит, что стенд собран, make check-clickhouse — что кластер работает кластером, а make check-services — что Airflow запускает DAG, а Superset ходит в базу. Дальше, в рабочей петле, обычно хватает make smoke.

.env.example — образец обязательной локальной настройки. В нём живут имя экземпляра, внешние порты, учебные учётные данные и ключи. Одно значение образца верно не везде — DOCKER_GID: у группы docker на каждой машине свой номер, и свой подскажет getent group docker. Пока он не сойдётся, стенд работает весь, кроме пульта мира — о нём ниже. compose.yaml не дублирует их умолчаниями: если значения нет в .env, Compose сразу предложит скопировать образец. Внутренние адреса, имена топиков и версии образов описывают сам стенд, поэтому записаны литералами в compose.yaml и Dockerfile.

Учётные данные Postgres и Grafana применяются при создании их томов. После первого запуска меняйте их только вместе с make clean: команда удалит все локальные данные стенда, а следующий make up создаст их с новыми значениями.

make up собирает локальные образы Airflow и Superset, поднимает весь стенд и ждёт здорового состояния долгоживущих контейнеров. В образ Airflow добавлены закреплённые клиенты ClickHouse и Kafka. Одноразовые airflow-init и superset-init завершаются с кодом 0. Первый обновляет схему Airflow, подготавливает администратора и подключение к clickhouse-01. Второй обновляет Superset, создаёт администратора и импортирует подключение к clickhouse-02. С нуля подъём занимает около трёх минут, на живом стенде — около минуты.

make smoke за секунды спрашивает, собран ли стенд: зависимости машины, здоровье контейнеров, устройство keeper, ответ Kafka с машины через отображённый порт, три цели Prometheus, источник Grafana, компоненты Airflow и подготовленное подключение к clickhouse-01. В конце проверка спрашивает у Docker, не убивало ли ядро что-нибудь в долгоживущих контейнерах за нехватку памяти и не включалась ли политика перезапуска: убитый контейнер Docker поднимает сам, и проверка состояния об этом промолчит.

Одиннадцать проверок здоровья сразу после make up --wait повторяют то, чего Compose уже дождался: у каждой долгоживущей службы есть своя healthcheck. Оставлены они потому, что первый вопрос к стенду всё равно «всё ли живо», а ответ на него стоит меньше секунды. Устройство keeper — другое дело: он должен работать от пользователя clickhouse, с пределом в 262144 открытых файла и со своим томом под данные. Всё это объявлено в compose.yaml, но здоровым keeper выглядит и без этого, поэтому смоук спрашивает у живого контейнера, дошли ли объявленные настройки до процесса.

Kafka по той же причине спрашивают снаружи. Её healthcheck обращается к брокеру изнутри контейнера и по внутреннему слушателю, поэтому здоровой Kafka остаётся и тогда, когда объявленный наружу адрес ведёт не туда. Клиент с машины в этом случае подключается, получает метаданные и молча виснет на адресе, которого с его стороны не существует. Смоук идёт тем же путём, что и будущий генератор: стучится в отображённый порт и просит список топиков.

make check-services проверяет то, ради чего приходится ждать службу: работу с топиком Kafka с машины — создан, найден, удалён, — ручной запуск пробников test_clickhouse и test_kafka, вход в Superset, его метаданные и подключение к clickhouse-02. Первый пробник создаёт таблицы на обеих нодах и читает через Distributed на ноде 2 строку из локальной таблицы ноды 1. Второй пишет в Kafka и читает свой маркер. Временный топик проверки с машины и запуски DAG удаляются; постоянный топик пробника сохраняется, а старые записи чистит Kafka.

make check-clickhouse запускает отдельную глубокую проверку ClickHouse: описание кластера, макросы, связь с keeper, ReplicatedMergeTree, Distributed, очередь распределённых DDL и очистку временных таблиц. Последняя из девяти проверок — единственная на настоящих данных: она подневно сверяет события стартового мира с описью и при расхождении говорит, где искать — в событиях или в браке.

make config-test проверяет Compose, синтаксис файлов DAG и пробельные ошибки в diff без запуска стенда.

Обычный рабочий цикл — make config-test и make smoke; остальные цели гоняют тогда, когда правка их касается. Что утверждает каждая проверка, нужен ли ей поднятый стенд, сколько стоит прогон и куда класть новую — в карте проверок.

Остановить контейнеры без удаления данных можно командой make down. Для полного сброса с удалением всех именованных томов используйте make clean. Повторный make up безопасен: одноразовая подготовка приложений идемпотентна.

Стартовый мир

Стенд поднимается не пустым: разовая служба world-init играет в топик hits первые восемь дней модельного времени — понедельник по понедельник, 401 185 событий. Дальше их обычным путём разбирает хранилище, и к концу make up они лежат в ods.event. Так у всякой лабы есть данные, и всегда одни и те же.

Ждать приходится дольше, чем работает заливка: приём асинхронный, поэтому вторым шагом make up зовёт scripts/wait-for-world.sh — тот опрашивает ClickHouse, пока мир не доедет. Повторный make up заливает мир заново; это не ошибка, а свойство: номера событий те же, и повтор схлопнет ReplacingMergeTree.

Сам мир в git не хранится — он чистая функция зерна, и держать его в репозитории значило бы держать там кэш. Вместо него лежит опись мира, data/world-inventory.json: зерно, версия генератора, хеш каталога товаров и по строке на каждый день — дата, число событий и хеш его байтов. Опись отвечает на единственный вопрос: тот ли это мир, что был вчера.

Спрашивают её двое. make test в generator/ сверяет опись с тем, что собирается из кода сегодня: правка генератора меняет мир, и опись надо пересобрать — make inventory там же. make check-clickhouse в корне сверяет с описью то, что доехало до ods.event, подневно. Хеш дня можно пересчитать и руками — это обычный sha256sum файла, который пишет файловый приёмник:

uv run --project generator python -m clickstream_generator batch \
  --day 0 --file tmp/day0.jsonl
sha256sum tmp/day0.jsonl

Как позвать генератор

События производит генератор из generator/. Модельный день — функция зерна и номера дня, поэтому один и тот же день всегда даёт те же события: обрыв лечится повтором. Позицию на оси генератор не помнит — какой день играть, решает зовущий.

Режима два, и различаются они темпом. Пакетный гонит день подряд, без пауз: так заливается мир и так переигрывается день после обрыва. Живой держит темп модельного времени — по умолчанию ×60: модельные сутки за 24 реальные минуты, суточная волна разворачивается на глазах, отставание видно в логе.

На стенде генератор ходит своим образом — разовой службой, которую поднимают и убирают на один прогон. На каждый режим по цели:

# день D0 целиком в топик hits
make generate-batch
# первая сотня событий дня D3
make generate-batch GENERATOR_DAY=3 GENERATOR_LIMIT=100
# день D3 живьём, ускорение ×1000
make generate-live GENERATOR_DAY=3 GENERATOR_SPEED=1000
Переменная Что задаёт Умолчание
GENERATOR_DAY номер дня на оси мира (D0 — первый) 0, задано в Makefile
GENERATOR_LIMIT потолок событий на прогон; только generate-batch нет: день целиком
GENERATOR_SPEED ускорение модельного времени; только generate-live ×60, задано в генераторе

Куда отправлять, обеим целям задаёт сам стенд: внутри сети Compose это брокер kafka:9092 и топик hits.

Образ цели не пересобирают: нет образа — Compose соберёт его сам, есть — возьмёт как есть. Пересобрать намеренно, после правки generator/Dockerfile или зависимостей, — отдельной командой:

docker compose --profile generator build generator

Тот же образ несёт заливка стартового мира, поэтому свежим его держит и обычный make up: он собирает образы всего стенда, и генератор теперь среди них.

Разведено это нарочно: собрать образ и запустить контейнер — разные действия, и цель запуска, молча пересобирающая образ, стирает между ними границу. Что образ устарел, видно по собственному прогону — это обратная связь, а не ловушка.

Посмотреть на события, не поднимая стенд, помогает файловый приёмник: одно событие — одна строка. Флаг --limit берёт начало дня вместо целого дня — тому, кто смотрит на конвейер, ждать полсотни тысяч событий незачем.

uv run --project generator python -m clickstream_generator batch \
  --day 0 --limit 100 --file tmp/day0.jsonl

Несколько дней подряд играет один запуск: --days 8. Всё, что принимается флагами, принимается и переменными окружения — так генератор позовёт даг этапа 5; полный список у --help.

Умолчания есть не у всего, и это нарочно. Зерно, число дней и темп живого дня описывают модель мира — их генератор знает сам. День на оси, приёмник и имя топика описывают стенд, на котором его запустили: не назвали — отказ и ненулевой код возврата, ещё до первого события. Поэтому умолчание дня и живёт в Makefile: день называет тот, кто запускает, а не тот, кого запускают.

Пульт мира

Считать дни самому необязательно. В Airflow живёт пульт мира — три дага, которые зовут тот же образ генератора и ведут позицию на оси за вас.

Даг Что делает
world_next_day играет следующие дни пачкой; сколько — параметр запуска, по умолчанию один
world_live_day играет следующий день в темпе модельного времени, около двадцати четырёх минут
world_live выключатель: пока снят с паузы, дёргает world_live_day день за днём

Позицию хранит переменная Airflow world_position — номер первого несыгранного дня. Ставит её работник, сыгравший день, и только по успеху: оборванный прогон позицию не двигает, и следующий запуск играет тот же день заново. Нет переменной — мир в стартовом состоянии.

Выключатель создаётся на паузе. Снимите — мир поедет сам; поставите обратно — встанет на границе модельных суток, доиграв начатый день. Форма пульта и доводы целиком — ADR 0009.

Плата названа вслух: планировщику Airflow отдан сокет докера — иначе контейнер генератора ему не поднять. Доступ к сокету равен праву root на машине; для локального учебного стенда размен принят, но знать о нём надо. Открывает дверь не монтирование, а членство в группе: GID группы docker у каждой машины свой, живёт в .env и подсматривается командой getent group docker.

Состав и доступ

  • clickhouse-01 — инициатор DDL и точка подключения Airflow;
  • clickhouse-02 — точка подключения Superset;
  • clickhouse-keeper — координатор кластера;
  • kafka — один брокер KRaft;
  • postgres-metadata — один Postgres с отдельными базами и пользователями airflow и superset;
  • airflow-apiserver, airflow-scheduler и airflow-dag-processor — Airflow 3.3 с LocalExecutor, без triggerer;
  • superset — интерфейс и подготовленное подключение ClickHouse;
  • prometheus и grafana — сбор и просмотр встроенных метрик ClickHouse.

С настройками из .env.example порты доступны только с локальной машины:

  • нода 1 — http://127.0.0.1:28123, нативный порт 29000;
  • нода 2 — http://127.0.0.1:28124, нативный порт 29001;
  • Kafka — 127.0.0.1:29092;
  • Airflow — http://127.0.0.1:28080, пользователь admin, пароль airflow;
  • Superset — http://127.0.0.1:28088, пользователь admin, пароль superset;
  • Prometheus — http://127.0.0.1:29090;
  • Grafana — http://127.0.0.1:23000, пользователь admin, пароль admin.

В ClickHouse четыре пользователя:

  • etl применяет DDL и подключает Airflow; роль etl_writer читает и пишет слои хранилища;
  • bi подключает Superset; роль bi_reader читает будущие слои DDS и DM;
  • analyst предназначен для подключения человека и читает все слои;
  • default остаётся служебным: им ходят проверки здоровья и скрипты внутри контейнеров, но не приложения.

Решение и его доводы — в ADR 0007. Пользователи и роли объявлены в infra/clickhouse/users.d/access.xml. Пароли и общий секрет нод живут в .env и передаются в конфигурацию через окружение; Superset получает пароль bi тем же путём через штатную функцию настройки. Значения для локального стенда есть в .env.example.

Пароли ClickHouse и интерфейсов, пароли Postgres, ключи Airflow и Superset, отсутствие проверки доступа у Kafka и Prometheus — намеренно простые и явно ненастоящие настройки локального учебного стенда. Это не пример настройки защиты: не копируйте значения из .env.example в рабочую среду. Все опубликованные порты привязаны только к 127.0.0.1; Postgres наружу не опубликован.

В бою перед репликами ClickHouse обычно был бы балансировщик. Здесь в каждом шарде одна реплика, поэтому балансировать нечего. Балансировщик и топология 2×2 намеренно не входят в стенд.

Конфигурация сверена 30 июля 2026 года с официальной документацией ClickHouse: настройками сервера, Keeper, ReplicatedMergeTree, ON CLUSTER и Distributed. Описание кластера задаётся через remote_servers, макросы — через macros, подключение к keeper — через zookeeper; путь ReplicatedMergeTree содержит {shard} и {replica}, а Distributed получает имя кластера, базу, локальную таблицу и ключ шардирования. Макросы выбраны, чтобы один DDL через ON CLUSTER создавал отдельный путь каждого шарда без вписанных вручную значений. Для образа зафиксирован точный текущий LTS-выпуск 26.3.17.56; серверы и keeper используют один образ. Настройки Kafka 4.3.1 сверены с примером односерверного KRaft. Секция метрик взята из конфигурации закреплённого образа ClickHouse и проверена на серверах и keeper. Подготовка источника Grafana сверена с официальным описанием автоматической настройки. Prometheus собирает только встроенные метрики двух серверов и keeper; внешних сборщиков, панелей и правил оповещения пока нет. После изменения infra/clickhouse/config.d/prometheus.xml выполните docker compose restart clickhouse-01 clickhouse-02: обычный make up не перезапускает уже созданные серверы и они продолжают работать со старой конфигурацией.

Airflow закреплён на 3.3.0. Состав обязательных процессов, LocalExecutor, публичный airflow.sdk, API здоровья и SimpleAuthManager сверены с архитектурой Airflow 3.3, публичным интерфейсом и описанием здоровья. Для пробников проверены публичный Connection.get из airflow.sdk и клиент clickhouse-connect. Официальный провайдер Kafka сам использует confluent-kafka; отдельное подключение и обёртки провайдера здесь не нужны, поэтому прямой клиент оставляет образ и пример короче. Superset закреплён на 6.1.0; драйвер clickhouse-connect, форма clickhousedb:// и драйвер Postgres сверены с документацией подключений Superset и настройкой базы метаданных.

Что здесь будет

  • одно широкое событие кликстрима по образцу выгрузки Яндекс Метрики вместо четырёх топиков;
  • второй источник — заказы бэкенда, ежедневным слепком в ту же Kafka;
  • сверка клиентской покупки с заказом бэкенда: деньги считаем по бэкенду, поведение и атрибуцию — по трекеру;
  • анонимный кликстрим и склейка кука↔пользователь через покупки;
  • ClickHouse кластером как единственным режимом.

Чем отличается от предшественника

Предшественник остаётся стабильным учебным стендом и заморожен для новых фич: там событие разрезано на четыре топика, есть только просмотры страниц, посетители опознаны по email, ClickHouse — одна нода. Развитие идёт здесь.

Документация

  • docs/specs/ — спеки: источник истины о задуманном.
  • docs/adr/ — принятые решения с доводами и отвергнутыми вариантами: почему сделано так, а не иначе.
  • docs/architecture/ — рабочие справочники по зонам ответственности: хранилище — конвенции имён, раскладка по шардам, приём событий, карта таблиц; проверки — что утверждает каждая цель make и куда класть новую проверку.
  • docs/research/ — исследования; сейчас это формат кликстрима Яндекса, по которому строится модель события.
  • docs/formats/ — описания форматов источников: по ним пишется сторона хранилища. Собираются из кода командой make docs в generator/, руками не правятся.
  • AGENTS.md — контракт работы в репозитории.
S
Description
No description provided
Readme
1,006 KiB
Languages
Python 83.3%
Shell 14.9%
Dockerfile 1.2%
Makefile 0.6%