- Зачем:
- до сих пор генератор умел собирать день, но не умел его отдать: топик
hits наполнялся пробником, а не настоящими данными. Тикет #41 доводит
события до стенда и закрывает форму на проводе, на которую обопрётся
типизированный ODS (#43).
- сериализатор один по решению спеки: второе место, печатающее событие в
JSON, разошлось бы с первым молча.
- Что:
- serialize.py — канонический сериализатор на orjson: единственное место,
где событие целиком становится JSON; 47 ключей всегда, «пусто» это
пустое значение, даты ISO-8601, ecommerce строкой. Вложенный блок
ecommerce в commerce.py вторым сериализатором не считается — правило
про событие, а не про блок внутри него.
- sinks.py — приёмники: файл (одно событие — одна строка) и Kafka (одно
событие — одно сообщение). Ключа у сообщения нет: WatchID уникален,
ключом он был бы ключом лишь на вид.
- player.py, cli.py — проигрыватель и интерфейс запуска: режимы batch и
live (темп ×60), несколько дней одним запуском, ограниченная пачка,
раздельные тайминги генерации и доставки, лаг в логе.
- день на оси и имя топика умолчаний не имеют: параметр, описывающий
среду или позицию, приходит от зовущего, иначе отказ до генерации.
Умолчания зерна, числа дней и темпа остаются — они описывают мир.
- generator/Dockerfile — свой образ: зависимости из uv.lock, база
закреплена до патча, раскладка репозитория сохранена ради каталога
товаров. Образ Airflow не тронут.
- разовая служба compose под профилем, цели generate-batch и
generate-live, .dockerignore, tmp/ в .gitignore.
- решения внесены в спеку (разделы 4, 8, 9), быстрый старт — в README.
- Проверка:
- make test 406 passed, make lint, make typecheck, make config-test.
- побайтовый детерминизм: два прогона дня в независимых процессах дают
один sha256; день в контейнере совпадает с днём на машине.
- на стенде: пакетный день доехал до stg.hits_raw_dist, счёт по
Distributed сошёлся — отправлено 50626, в таблице 50626.
- топик прочитан обеими нодами: clickhouse-01 раздел 0 (26368),
clickhouse-02 раздел 1 (24258).
- живой день: модельное время 01:00 на 60-й секунде, 02:00 на 120-й —
темп ×60, лаг печатается.
- форма на проводе в колонке raw: даты читаются глазами, ecommerce лежит
строкой.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
79 lines
4.9 KiB
Python
79 lines
4.9 KiB
Python
"""Канонический сериализатор: единственное место, где событие целиком → JSON.
|
||
|
||
Правило «сериализатор один» (спека генератора, разделы 4 и 6) — не про
|
||
экономию строк, а про канон: два прогона одного дня обязаны дать те же байты,
|
||
а байты рождаются здесь. Второе место, собирающее событие руками, разошлось бы
|
||
с этим по экранированию, порядку ключей или записи чисел — и разошлось бы
|
||
молча. Граница правила проходит по событию, а не по всякому JSON: вложенный
|
||
блок `ecommerce` собирает `commerce`, и это часть содержимого колонки, а не
|
||
второй сериализатор.
|
||
|
||
Что делает канон:
|
||
|
||
- **Порядок ключей — порядок контракта схемы.** Он берётся из `schema.COLUMNS`
|
||
и нигде не повторяется: два источника порядка разъехались бы при первой же
|
||
вставке колонки.
|
||
- **Все 47 ключей всегда.** Пусто по контракту — пустое значение: пустой
|
||
массив, пустая строка, ноль. Пропавший ключ увёл бы событие в брак целиком:
|
||
строгий приём хранилища сверяет набор ключей (ADR 0005).
|
||
- **Даты и время — ISO-8601** (спека, раздел 4): `EventDate` уезжает как
|
||
`2026-06-01`, `UTCEventTime` — как `2026-06-01T12:34:56Z`. Довод — читаемость
|
||
сырья: менти открывает колонку `raw` обычным клиентом и разбирает событие
|
||
глазами, а число эпохи этот урок убивает.
|
||
- **Одно событие — один документ JSON**, без перевода строки внутри: приёмник
|
||
сам решает, чем их разделить.
|
||
|
||
Колонки переводятся в питоновские значения целиком, а не по строкам: numpy
|
||
делает это одним вызовом на колонку, и на дне в полсотни тысяч событий разница
|
||
заметна. Обратная сторона — день лежит в памяти дважды; проигрыватель поэтому
|
||
и берёт его днями, а не горизонтом целиком.
|
||
"""
|
||
|
||
from typing import Any
|
||
|
||
import numpy as np
|
||
import orjson
|
||
from numpy.typing import NDArray
|
||
|
||
from clickstream_generator import schema
|
||
from clickstream_generator.day import Day
|
||
|
||
_ARRAY_PREFIX = "Array("
|
||
|
||
|
||
def events(day: Day, limit: int | None = None) -> list[bytes]:
|
||
"""Канонические байты событий дня: по документу JSON на событие.
|
||
|
||
`limit` берёт первые события дня и на этом останавливается — срез для
|
||
того, кто смотрит на конвейер и не хочет ждать целый день (спека,
|
||
раздел 9). Ограничение считается до сериализации: платить за то, что не
|
||
поедет, незачем.
|
||
"""
|
||
count = len(day) if limit is None else min(limit, len(day))
|
||
names = tuple(column.name for column in schema.COLUMNS)
|
||
values = [
|
||
_values(column, day.columns[column.name][:count]) for column in schema.COLUMNS
|
||
]
|
||
return [
|
||
orjson.dumps(dict(zip(names, row, strict=True)))
|
||
for row in zip(*values, strict=True)
|
||
]
|
||
|
||
|
||
def _values(column: schema.Column, values: NDArray[Any]) -> list[Any]:
|
||
"""Колонка питоновскими значениями — в той записи, в какой уедет на провод.
|
||
|
||
Массив узнаётся по типу ClickHouse, а не по `numpy_dtype`: у колонки-массива
|
||
там записан тип элемента (`uint32`), и от скалярной колонки её этим не
|
||
отличить.
|
||
"""
|
||
if column.clickhouse_type.startswith(_ARRAY_PREFIX):
|
||
# Колонка-массив: в ячейке лежит свой массив, пустой у события,
|
||
# которому эта колонка не по смыслу.
|
||
return [cell.tolist() for cell in values]
|
||
if column.numpy_dtype == "datetime64[D]":
|
||
return np.datetime_as_string(values, unit="D").tolist()
|
||
if column.numpy_dtype == "datetime64[s]":
|
||
return np.datetime_as_string(values, unit="s", timezone="UTC").tolist()
|
||
return values.tolist()
|