feat(generator): сериализатор, приёмники, проигрыватель и запуск контейнером
- Зачем:
- до сих пор генератор умел собирать день, но не умел его отдать: топик
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>
This commit is contained in:
@@ -0,0 +1,9 @@
|
||||
"""Точка входа пакета: `python -m clickstream_generator`.
|
||||
|
||||
Код возврата уходит наружу как есть — по нему судит о прогоне и контейнер, и
|
||||
даг, который его запустит (спека генератора, раздел 9).
|
||||
"""
|
||||
|
||||
from clickstream_generator.cli import main
|
||||
|
||||
raise SystemExit(main())
|
||||
@@ -0,0 +1,207 @@
|
||||
"""Интерфейс запуска: `python -m clickstream_generator batch|live`.
|
||||
|
||||
Зовущий — контейнер (спека генератора, раздел 9), и интерфейс сделан под него:
|
||||
параметры приходят аргументами или переменными окружения, логи идут в
|
||||
стандартный вывод, итог виден кодом возврата. Ни файла настроек, ни состояния
|
||||
между запусками нет — позиция на оси принадлежит тому, кто зовёт.
|
||||
|
||||
Зовущих трое, и все трое видны в форме команд:
|
||||
|
||||
- даги `world_init` и `next_day` этапа 5 — по дню за запуск, приёмник Kafka;
|
||||
- заливка зернового мира (#42) — восемь дней подряд одним запуском: `--days`;
|
||||
- проверки хранилища (#43) — ограниченная пачка в файл: `--limit` и `--file`.
|
||||
|
||||
**Режимы разведены командами, а не флагом**, потому что различаются не темпом
|
||||
в числе, а тем, что у них разное: у пакетного есть `--limit` и нет ожидания, у
|
||||
живого есть `--speed` и нет пачки. Один флаг `--speed 0` прятал бы это
|
||||
различие за числом, а команда называет его словом. Ограниченная пачка живому
|
||||
дню не полагается: ждать там нечего — ожидание снимает ускорение.
|
||||
|
||||
**Приёмник выбирается тем, что для него назвали**: `--file` или `--brokers`.
|
||||
Оба сразу — ошибка, ни одного — тоже: молча выбранный по умолчанию приёмник
|
||||
однажды напишет в файл то, чего ждали в Kafka.
|
||||
|
||||
**Умолчания есть только у того, что описывает мир, а не окружение.** Зерно,
|
||||
число дней и темп живого дня — свойства модели, они заданы спекой и тикетом,
|
||||
и умолчание здесь — тот же самый ответ. Позиция на оси (`--day`), адрес
|
||||
брокера, имя топика и путь файла — факты стенда, на котором нас запустили:
|
||||
генератор их не знает и знать не должен, он отдельная и переносимая сущность.
|
||||
Подставь он своё умолчание — забытый параметр превратился бы в неверный ответ,
|
||||
поданный как успех.
|
||||
"""
|
||||
|
||||
import argparse
|
||||
import logging
|
||||
import os
|
||||
import sys
|
||||
from collections.abc import Callable, Sequence
|
||||
from contextlib import closing
|
||||
from pathlib import Path
|
||||
|
||||
from clickstream_generator import player
|
||||
from clickstream_generator.seeds import CANONICAL_SEED
|
||||
from clickstream_generator.sinks import FileSink, KafkaSink, Sink
|
||||
|
||||
DEFAULT_SPEED = 60.0
|
||||
|
||||
|
||||
def main(argv: Sequence[str] | None = None) -> int:
|
||||
"""Разобрать аргументы, проиграть дни, вернуть код возврата."""
|
||||
parser = _parser()
|
||||
options = parser.parse_args(argv)
|
||||
logging.basicConfig(
|
||||
level=logging.INFO, format="%(message)s", stream=sys.stdout, force=True
|
||||
)
|
||||
|
||||
_check(parser, options)
|
||||
try:
|
||||
with closing(_sink(parser, options)) as sink:
|
||||
player.play(
|
||||
sink,
|
||||
seed=options.seed,
|
||||
first_day=options.day,
|
||||
days=options.days,
|
||||
limit=options.limit,
|
||||
speed=options.speed,
|
||||
)
|
||||
except (OSError, RuntimeError) as failure:
|
||||
logging.error("прогон не удался: %s", failure)
|
||||
return 1
|
||||
return 0
|
||||
|
||||
|
||||
def _parser() -> argparse.ArgumentParser:
|
||||
"""Разбор командной строки; умолчания приходят из окружения."""
|
||||
parser = argparse.ArgumentParser(
|
||||
prog="python -m clickstream_generator",
|
||||
description="Проигрыватель модельных дней кликстрима в файл или Kafka.",
|
||||
)
|
||||
modes = parser.add_subparsers(dest="mode", required=True)
|
||||
|
||||
batch = modes.add_parser(
|
||||
"batch", help="пачкой, без пауз: заливка снимка и переигровка дня"
|
||||
)
|
||||
_common(batch)
|
||||
batch.set_defaults(speed=None)
|
||||
batch.add_argument(
|
||||
"--limit",
|
||||
type=int,
|
||||
default=_env_int("GENERATOR_LIMIT", None),
|
||||
metavar="N",
|
||||
help="взять не больше N событий на весь прогон, а не день целиком",
|
||||
)
|
||||
|
||||
live = modes.add_parser("live", help="живой день: темп модельного времени")
|
||||
_common(live)
|
||||
live.set_defaults(limit=None)
|
||||
live.add_argument(
|
||||
"--speed",
|
||||
type=float,
|
||||
default=_env_float("GENERATOR_SPEED", DEFAULT_SPEED),
|
||||
metavar="X",
|
||||
help=f"ускорение модельного времени (по умолчанию ×{DEFAULT_SPEED:.0f})",
|
||||
)
|
||||
return parser
|
||||
|
||||
|
||||
def _common(parser: argparse.ArgumentParser) -> None:
|
||||
"""Что спрашивают у обоих режимов: какой мир, какие дни и куда."""
|
||||
parser.add_argument(
|
||||
"--seed",
|
||||
type=int,
|
||||
default=_env_int("GENERATOR_SEED", CANONICAL_SEED),
|
||||
metavar="N",
|
||||
help="зерно мира (по умолчанию каноническое)",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--day",
|
||||
type=int,
|
||||
default=_env_int("GENERATOR_DAY", None),
|
||||
metavar="D",
|
||||
help="номер дня на оси мира (D0 — первый); умолчания нет —"
|
||||
" позицию ведёт зовущий",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--days",
|
||||
type=int,
|
||||
default=_env_int("GENERATOR_DAYS", 1),
|
||||
metavar="N",
|
||||
help="сколько дней подряд проиграть одним запуском",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--file",
|
||||
type=Path,
|
||||
default=_env("GENERATOR_FILE", None, Path),
|
||||
metavar="ПУТЬ",
|
||||
help="приёмник — файл: одно событие в строке",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--brokers",
|
||||
default=os.environ.get("KAFKA_BOOTSTRAP_SERVERS"),
|
||||
metavar="АДРЕС",
|
||||
help="приёмник — Kafka: адреса брокеров через запятую",
|
||||
)
|
||||
parser.add_argument(
|
||||
"--topic",
|
||||
default=os.environ.get("KAFKA_TOPIC"),
|
||||
metavar="ИМЯ",
|
||||
help="топик Kafka; умолчания нет — имя топика принадлежит стенду, а не пакету",
|
||||
)
|
||||
|
||||
|
||||
def _check(parser: argparse.ArgumentParser, options: argparse.Namespace) -> None:
|
||||
"""День на оси обязан быть назван — иначе прогон не начинается."""
|
||||
if options.day is None:
|
||||
parser.error(
|
||||
"день на оси не назван: передайте --day или GENERATOR_DAY."
|
||||
" Позицию ведёт зовущий — проигрыватель её не помнит"
|
||||
)
|
||||
|
||||
|
||||
def _sink(parser: argparse.ArgumentParser, options: argparse.Namespace) -> Sink:
|
||||
"""Приёмник по названному: файл либо Kafka, но не оба и не ничего."""
|
||||
if options.file and options.brokers:
|
||||
parser.error("названы оба приёмника: оставьте --file либо --brokers")
|
||||
if options.file:
|
||||
return FileSink(options.file)
|
||||
if options.brokers:
|
||||
if not options.topic:
|
||||
parser.error(
|
||||
"топик не назван: передайте --topic или KAFKA_TOPIC."
|
||||
" Имя топика — факт стенда, генератор его не знает"
|
||||
)
|
||||
try:
|
||||
return KafkaSink(options.brokers, options.topic)
|
||||
except ImportError:
|
||||
# Ловим там, где возникает: обёрнутый вокруг всего прогона,
|
||||
# этот перехват однажды объявил бы «нет клиента Kafka» о чужой
|
||||
# сорванной загрузке модуля.
|
||||
parser.error(
|
||||
"приёмник Kafka требует клиента: поставьте пакет с группой"
|
||||
" зависимостей kafka (`uv sync --extra kafka`)"
|
||||
)
|
||||
parser.error("приёмник не назван: нужен --file либо --brokers")
|
||||
|
||||
|
||||
def _env[T](name: str, fallback: T, kind: Callable[[str], T]) -> T:
|
||||
"""Значение переменной окружения нужного типа; нет переменной — умолчание.
|
||||
|
||||
Пустая строка считается отсутствием: в compose так выглядит переменная,
|
||||
которую не задали, — `${GENERATOR_DAYS:-}`. Мусор в переменной называется
|
||||
вместе с её именем: из контейнера иначе не видно, чьё это значение.
|
||||
"""
|
||||
value = os.environ.get(name)
|
||||
if not value:
|
||||
return fallback
|
||||
try:
|
||||
return kind(value)
|
||||
except ValueError:
|
||||
raise SystemExit(f"переменная {name} не разбирается: {value!r}") from None
|
||||
|
||||
|
||||
def _env_int(name: str, fallback: int | None) -> int | None:
|
||||
return _env(name, fallback, int)
|
||||
|
||||
|
||||
def _env_float(name: str, fallback: float) -> float:
|
||||
return _env(name, fallback, float)
|
||||
@@ -0,0 +1,157 @@
|
||||
"""Проигрыватель: гонит дни мира в приёмник — пачкой или с темпом живого дня.
|
||||
|
||||
Состояния у него нет (спека генератора, раздел 9). Зерно и номер дня приходят
|
||||
параметрами, позицию на оси он не хранит и из данных не выводит: её ведёт тот,
|
||||
кто зовёт, — на этапе 5 это переменная Airflow у дага `next_day`. Отсюда и
|
||||
переигровка обрыва: позвали тот же день заново — получили те же `WatchID`, и
|
||||
дедупликация склеила повтор.
|
||||
|
||||
Два режима отличаются только темпом. Пакетный шлёт события подряд, без пауз, —
|
||||
это заливка снимка и переигровка дня. Живой держит модельное время: событие
|
||||
уезжает тогда, когда до него дошли модельные часы, поделённые на ускорение.
|
||||
По умолчанию ускорение ×60 — модельные сутки за 24 реальные минуты (спека,
|
||||
раздел 5): суточная волна разворачивается на глазах.
|
||||
|
||||
Тайминги печатаются раздельно — генерация и доставка, как требует спека
|
||||
(раздел 5): это разные машины разной природы, и сложенные в одно число они
|
||||
перестают что-либо говорить. Сериализация считается частью генерации: она
|
||||
рождает те самые байты, которые сторожит манифест.
|
||||
"""
|
||||
|
||||
import logging
|
||||
import time
|
||||
from dataclasses import dataclass
|
||||
|
||||
import numpy as np
|
||||
from numpy.typing import NDArray
|
||||
|
||||
from clickstream_generator import day as day_module
|
||||
from clickstream_generator import serialize
|
||||
from clickstream_generator.sinks import Sink
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
||||
# Как часто живой режим отчитывается о ходе дня. Минута реального времени — это
|
||||
# час модельного при ×60: отчёт на каждый модельный час.
|
||||
REPORT_SECONDS = 60.0
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class Played:
|
||||
"""Итог прогона: сколько уехало и за сколько."""
|
||||
|
||||
events: int
|
||||
generated_seconds: float
|
||||
delivered_seconds: float
|
||||
|
||||
|
||||
def play(
|
||||
sink: Sink,
|
||||
seed: int,
|
||||
first_day: int,
|
||||
days: int = 1,
|
||||
limit: int | None = None,
|
||||
speed: float | None = None,
|
||||
) -> Played:
|
||||
"""Проиграть `days` дней подряд начиная с `first_day` в приёмник `sink`.
|
||||
|
||||
`limit` — потолок событий на весь прогон: срез для того, кто смотрит на
|
||||
конвейер и не хочет ждать целый день. `speed` — ускорение живого режима;
|
||||
`None` означает пакетный, то есть без пауз вовсе.
|
||||
"""
|
||||
events = 0
|
||||
generated = 0.0
|
||||
delivered = 0.0
|
||||
|
||||
for number in range(first_day, first_day + days):
|
||||
left = None if limit is None else limit - events
|
||||
if left is not None and left <= 0:
|
||||
break
|
||||
|
||||
clock = time.monotonic()
|
||||
today = day_module.stream(seed, number)
|
||||
payloads = serialize.events(today, limit=left)
|
||||
seconds = _event_seconds(today, len(payloads))
|
||||
spent = time.monotonic() - clock
|
||||
generated += spent
|
||||
log.info("день %d: событий %d, генерация %.1f с", number, len(payloads), spent)
|
||||
|
||||
clock = time.monotonic()
|
||||
if speed is None:
|
||||
_send(sink, payloads)
|
||||
else:
|
||||
_send_paced(sink, payloads, seconds, speed)
|
||||
# Рубеж дня: доставка асинхронна, и без него напечатанное время
|
||||
# означало бы только «события легли в очередь отправителя».
|
||||
sink.flush()
|
||||
spent = time.monotonic() - clock
|
||||
delivered += spent
|
||||
|
||||
events += len(payloads)
|
||||
log.info(
|
||||
"день %d: отправлено %d, доставка %.1f с", number, len(payloads), spent
|
||||
)
|
||||
|
||||
log.info(
|
||||
"итого отправлено %d событий: генерация %.1f с, доставка %.1f с",
|
||||
events,
|
||||
generated,
|
||||
delivered,
|
||||
)
|
||||
return Played(
|
||||
events=events, generated_seconds=generated, delivered_seconds=delivered
|
||||
)
|
||||
|
||||
|
||||
def _send(sink: Sink, payloads: list[bytes]) -> None:
|
||||
"""Пакетно: подряд и без пауз."""
|
||||
for payload in payloads:
|
||||
sink.send(payload)
|
||||
|
||||
|
||||
def _send_paced(
|
||||
sink: Sink, payloads: list[bytes], seconds: NDArray[np.int64], speed: float
|
||||
) -> None:
|
||||
"""С темпом: событие уезжает, когда до него дошло модельное время.
|
||||
|
||||
Отставание не догоняется рывком и не прячется: спешить некуда — событие
|
||||
всё равно уедет, — а вот увидеть отставание в логе нужно, иначе живой
|
||||
режим врёт про темп. Обгонять модельное время нельзя, отставать можно, и
|
||||
именно это печатает отчёт.
|
||||
"""
|
||||
started = time.monotonic()
|
||||
origin = int(seconds[0]) if seconds.size else 0
|
||||
reported = started
|
||||
lag = 0.0
|
||||
|
||||
for number, payload in enumerate(payloads):
|
||||
due = started + (int(seconds[number]) - origin) / speed
|
||||
now = time.monotonic()
|
||||
if now < due:
|
||||
time.sleep(due - now)
|
||||
else:
|
||||
lag = max(lag, now - due)
|
||||
sink.send(payload)
|
||||
|
||||
now = time.monotonic()
|
||||
if now - reported >= REPORT_SECONDS:
|
||||
log.info(
|
||||
"проиграно %d из %d, модельное время %s, лаг %.1f с",
|
||||
number + 1,
|
||||
len(payloads),
|
||||
_model_time(int(seconds[number]) - origin),
|
||||
lag,
|
||||
)
|
||||
reported = now
|
||||
lag = 0.0
|
||||
|
||||
|
||||
def _event_seconds(today: day_module.Day, count: int) -> NDArray[np.int64]:
|
||||
"""Секунды событий абсолютной меткой — по ним живой режим держит темп."""
|
||||
times: NDArray[np.datetime64] = today.columns["UTCEventTime"][:count]
|
||||
return times.astype("datetime64[s]").astype(np.int64)
|
||||
|
||||
|
||||
def _model_time(elapsed: int) -> str:
|
||||
"""Прожитое модельное время дня в виде `ЧЧ:ММ` — от первого события."""
|
||||
return f"{elapsed // 3600:02d}:{elapsed % 3600 // 60:02d}"
|
||||
@@ -0,0 +1,78 @@
|
||||
"""Канонический сериализатор: единственное место, где событие целиком → 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()
|
||||
@@ -0,0 +1,137 @@
|
||||
"""Приёмники: куда уезжают канонические байты. Про содержимое они не знают.
|
||||
|
||||
Приёмник глуп по замыслу (спека генератора, раздел 4): он берёт готовый байт
|
||||
события и доставляет его. Ни формы, ни темпа он не решает — форму задал
|
||||
сериализатор, темп задаёт проигрыватель. Отсюда и весь их интерфейс: принять
|
||||
событие, дождаться принятого, закрыться.
|
||||
|
||||
Приёмников два, а режимов три: «Kafka пачкой» и «Kafka с темпом ×60»
|
||||
различаются не приёмником, а тем, кто его зовёт. Ускорять доставку, зная о
|
||||
времени события, значило бы вернуть приёмнику знание о содержимом — ровно то,
|
||||
чего правило не хочет.
|
||||
|
||||
**Ключа у сообщения Kafka нет** — решение тикета #41, и вот довод. `WatchID`
|
||||
уникален у каждого события, поэтому ключом он не был бы ключом: обещание Kafka
|
||||
про ключ — «сообщения одного ключа лежат в одном разделе и сохраняют порядок», а
|
||||
у ряда, где ключи не повторяются, обещать нечего. Читатель же прочёл бы такой
|
||||
ключ как смысл, которого в нём нет. Настоящий ключ здесь — `ClientID`
|
||||
(события одной куки по порядку), и он тикетом не назначен: раскладку по
|
||||
разделам стенд намеренно оставляет транспорту (спека, раздел 2), а лаба
|
||||
«какая нода читала топик» живёт как раз тем, что раскладка не предрешена.
|
||||
|
||||
Раздел для сообщения без ключа librdkafka выбирает случайно, но подряд идущие
|
||||
сообщения на десяток миллисекунд липнут к одному — измеренные умолчания и их
|
||||
следствия записаны в спеке (раздел 9).
|
||||
"""
|
||||
|
||||
from pathlib import Path
|
||||
from typing import Any, Protocol
|
||||
|
||||
# Сколько ждать разбора очереди отправителя, когда она заполнилась, и сколько —
|
||||
# доставки остатка при закрытии. Оба числа — потолок ожидания, а не пауза:
|
||||
# обычно ждать не приходится вовсе.
|
||||
QUEUE_WAIT_SECONDS = 1.0
|
||||
FLUSH_WAIT_SECONDS = 60.0
|
||||
|
||||
|
||||
class Sink(Protocol):
|
||||
"""Приёмник: принимает байты события и доводит их до места."""
|
||||
|
||||
def send(self, payload: bytes) -> None:
|
||||
"""Принять одно событие."""
|
||||
|
||||
def flush(self) -> None:
|
||||
"""Дождаться, пока принятое дойдёт до места.
|
||||
|
||||
Нужно не только при закрытии: доставка асинхронна, и без этого рубежа
|
||||
«доставка 0,1 с» в логе означала бы лишь то, что события успели лечь в
|
||||
очередь отправителя. Проигрыватель ставит рубеж в конце каждого дня.
|
||||
"""
|
||||
|
||||
def close(self) -> None:
|
||||
"""Закрыть приёмник, дождавшись всего принятого."""
|
||||
|
||||
|
||||
class FileSink:
|
||||
"""Файл: одно событие — одна строка.
|
||||
|
||||
Построчность — контракт файла, а не удобство: файл читают построчно, и
|
||||
склейка двух событий в строку сломала бы разбор целиком. Файл — кэш чистой
|
||||
функции (спека, раздел 4): потерял — пересчитал, поэтому места в
|
||||
репозитории ему не отведено.
|
||||
"""
|
||||
|
||||
def __init__(self, path: Path) -> None:
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
self._file = path.open("wb")
|
||||
|
||||
def send(self, payload: bytes) -> None:
|
||||
self._file.write(payload + b"\n")
|
||||
|
||||
def flush(self) -> None:
|
||||
self._file.flush()
|
||||
|
||||
def close(self) -> None:
|
||||
self._file.close()
|
||||
|
||||
|
||||
class KafkaSink:
|
||||
"""Kafka: одно событие — одно сообщение.
|
||||
|
||||
«Пачкой» относится к темпу отправки, а не к упаковке: события идут подряд
|
||||
без пауз, но каждое отдельным сообщением (спека, раздел 4). Отправка
|
||||
асинхронная — отправитель копит сообщения в своей очереди и шлёт их
|
||||
пачками сам; отсюда `poll`, который отдаёт нам отчёты о доставке, и
|
||||
`flush`, без которого хвост очереди уехал бы в никуда вместе с процессом.
|
||||
"""
|
||||
|
||||
def __init__(self, brokers: str, topic: str) -> None:
|
||||
# Импорт внутри: клиент — необязательная часть пакета, и без него
|
||||
# генератор пишет в файл (спека, раздел 9).
|
||||
from confluent_kafka import Producer
|
||||
|
||||
self._topic = topic
|
||||
self._failure: str | None = None
|
||||
self._producer = Producer({"bootstrap.servers": brokers})
|
||||
|
||||
def send(self, payload: bytes) -> None:
|
||||
while True:
|
||||
try:
|
||||
self._producer.produce(
|
||||
self._topic, value=payload, on_delivery=self._report
|
||||
)
|
||||
break
|
||||
except BufferError:
|
||||
# Очередь отправителя полна: ждём, пока брокер её разберёт.
|
||||
# `poll` здесь и работа, и пауза — он же отдаёт отчёты.
|
||||
self._producer.poll(QUEUE_WAIT_SECONDS)
|
||||
self._producer.poll(0)
|
||||
self._raise_failure()
|
||||
|
||||
def flush(self) -> None:
|
||||
remaining = self._producer.flush(FLUSH_WAIT_SECONDS)
|
||||
if remaining:
|
||||
raise RuntimeError(
|
||||
f"Kafka не приняла {remaining} сообщений за"
|
||||
f" {FLUSH_WAIT_SECONDS:.0f} с: брокер недоступен или не успевает"
|
||||
)
|
||||
self._raise_failure()
|
||||
|
||||
def close(self) -> None:
|
||||
# Закрывать у отправителя нечего — важно лишь не бросить хвост
|
||||
# очереди: он уехал бы в никуда вместе с процессом.
|
||||
self.flush()
|
||||
|
||||
def _report(self, error: Any, message: Any) -> None:
|
||||
"""Отчёт о доставке: первую неудачу запоминаем, остальные не важны."""
|
||||
if error is not None and self._failure is None:
|
||||
self._failure = str(error)
|
||||
|
||||
def _raise_failure(self) -> None:
|
||||
"""Неудачную доставку превращаем в остановку прогона.
|
||||
|
||||
Молча потерянное сообщение — худший исход: счёт в хранилище разойдётся
|
||||
с числом отправленного, а причина будет забыта.
|
||||
"""
|
||||
if self._failure is not None:
|
||||
raise RuntimeError(f"Kafka не приняла сообщение: {self._failure}")
|
||||
Reference in New Issue
Block a user