feat(generator): добавлена команда слепков заказов

- Зачем:
  - задача #92 добавляет второй источник и показывает слепок как чистую функцию зерна и дня.
- Что:
  - добавлены окно, состояние на границе суток и канонические байты с настоящими миллисекундами.
  - добавлены команда snapshot, сдвиг D → D−1 и переиспользование проигранных дней.
  - опись дополнена хешами слепков и поведенческими тестами.
- Проверка:
  - в generator выполнены make lint, make typecheck и make test: 427 тестов.
This commit is contained in:
2026-08-18 17:12:49 +03:00
parent c11da3b252
commit 8e21381e13
11 changed files with 740 additions and 46 deletions
+31 -16
View File
@@ -1,4 +1,4 @@
"""Интерфейс запуска: `python -m clickstream_generator batch|live`.
"""Интерфейс запуска: `python -m clickstream_generator batch|live|snapshot`.
Зовущий — контейнер (спека генератора, раздел 9), и интерфейс сделан под него:
параметры приходят аргументами или переменными окружения, логи идут в
@@ -17,6 +17,11 @@
различие за числом, а команда называет его словом. Ограниченная пачка живому
дню не полагается: ждать там нечего — ожидание снимает ускорение.
**Третья команда — `snapshot`** — второй источник стенда: слепок заказов в свой
топик. У неё свой смысл `--day`: это день прогона, а уезжает слепок дня D−1
(правило сдвига живёт в проигрывателе). Прогон дня 0 не отправляет ничего —
на старте оси вчера нет.
**Приёмник выбирается тем, что для него назвали**: `--file` или `--brokers`.
Оба сразу — ошибка, ни одного — тоже: молча выбранный по умолчанию приёмник
однажды напишет в файл то, чего ждали в Kafka.
@@ -56,14 +61,19 @@ def main(argv: Sequence[str] | None = None) -> int:
_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,
)
if options.command == "snapshot":
player.play_snapshots(
sink, seed=options.seed, first_day=options.day, days=options.days
)
else:
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
@@ -76,9 +86,9 @@ def _parser() -> argparse.ArgumentParser:
prog="python -m clickstream_generator",
description="Проигрыватель модельных дней кликстрима в файл или Kafka.",
)
modes = parser.add_subparsers(dest="mode", required=True)
commands = parser.add_subparsers(dest="command", required=True)
batch = modes.add_parser(
batch = commands.add_parser(
"batch", help="пачкой, без пауз: заливка снимка и переигровка дня"
)
_common(batch)
@@ -91,7 +101,7 @@ def _parser() -> argparse.ArgumentParser:
help="взять не больше N событий на весь прогон, а не день целиком",
)
live = modes.add_parser("live", help="живой день: темп модельного времени")
live = commands.add_parser("live", help="живой день: темп модельного времени")
_common(live)
live.set_defaults(limit=None)
live.add_argument(
@@ -101,11 +111,16 @@ def _parser() -> argparse.ArgumentParser:
metavar="X",
help=f"ускорение модельного времени (по умолчанию ×{DEFAULT_SPEED:.0f})",
)
snapshot = commands.add_parser(
"snapshot", help="слепок заказов: прогон дня D отправляет слепок дня D−1"
)
_common(snapshot)
return parser
def _common(parser: argparse.ArgumentParser) -> None:
"""Что спрашивают у обоих режимов: какой мир, какие дни и куда."""
"""Что спрашивают у всех команд: какой мир, какие дни и куда."""
parser.add_argument(
"--seed",
type=int,
@@ -118,8 +133,8 @@ def _common(parser: argparse.ArgumentParser) -> None:
type=int,
default=_env_int("GENERATOR_DAY", None),
metavar="D",
help="номер дня на оси мира (D0 — первый); умолчания нет —"
" позицию ведёт зовущий",
help="номер дня на оси мира (D0 — первый); у snapshot это день прогона,"
" а уезжает слепок дня D−1; умолчания нет — позицию ведёт зовущий",
)
parser.add_argument(
"--days",
@@ -133,7 +148,7 @@ def _common(parser: argparse.ArgumentParser) -> None:
type=Path,
default=_env("GENERATOR_FILE", None, Path),
metavar="ПУТЬ",
help="приёмник — файл: одно событие в строке",
help="приёмник — файл: одна запись в строке",
)
parser.add_argument(
"--brokers",
@@ -151,6 +151,10 @@ class Purchases:
# та самая, что уехала в событие как `purchaseRevenue`. Бэкенд назовёт
# её `items_total` — у двух источников свои имена одному числу.
revenue: NDArray[np.int64]
# Когда покупка подтверждена — та же абсолютная метка, что уехала в
# событие. У заказа она станет секундой `created_at`: строка в базе
# источника создаётся синхронно с покупкой.
moment: NDArray[np.datetime64]
def __len__(self) -> int:
return self.person_id.size
@@ -185,6 +189,7 @@ _NO_PURCHASES = Purchases(
quantity=(),
coupon=(),
revenue=np.empty(0, dtype=np.int64),
moment=np.empty(0, dtype="datetime64[s]"),
)
@@ -666,6 +671,7 @@ def _purchases(
quantity=tuple(draws.quantity[group] for group in bought),
coupon=tuple(coupons[mine].tolist()),
revenue=revenue[mine],
moment=rows["UTCEventTime"][here],
)
@@ -26,6 +26,11 @@
пересчитывается он и обычным `sha256sum` по сыгранному в файл дню (как
именно в README репозитория).
Слепок заказов второй артефакт мира, и хешами дней он не покрыт ни при каком
раскладе подпотоков, поэтому у каждого отправленного слепка своя строка. Их на
один меньше, чем дней: прогон дня 0 отправлять ещё нечего. Счёта строк у
слепка нет опись описывает мир, а не доставку.
Собирается опись из каталога генератора целью `make inventory`, а свежесть её
сторожит тест как и у «описания выгрузки».
"""
@@ -39,6 +44,7 @@ from pathlib import Path
from typing import Any
from clickstream_generator import day as day_module
from clickstream_generator import orders as orders_module
from clickstream_generator import serialize, world
from clickstream_generator.catalog import CATALOG_PATH
from clickstream_generator.seeds import CANONICAL_SEED
@@ -50,12 +56,25 @@ STARTING_DAYS = 8
def build() -> dict[str, Any]:
"""Опись целиком: паспорт мира и по строке на каждый его день."""
"""Опись целиком: паспорт мира, строка на день и строка на слепок.
Дни играются по одному и отпускаются: заказы дня остаются, потому что из
них собираются слепки, а полсотни тысяч событий восьми дней сразу в память
не нужны.
"""
days = []
orders = []
for number in range(STARTING_DAYS):
today = day_module.stream(CANONICAL_SEED, number)
days.append(_day(number, today))
orders.append(today.orders)
return {
"seed": CANONICAL_SEED,
"generator_version": version("clickstream-generator"),
"catalog_sha256": _digest(CATALOG_PATH.read_bytes()),
"days": [_day(number) for number in range(STARTING_DAYS)],
"days": days,
"snapshots": [_snapshot(number, orders) for number in range(STARTING_DAYS - 1)],
}
@@ -64,7 +83,7 @@ def render() -> str:
return json.dumps(build(), ensure_ascii=False, indent=2) + "\n"
def _day(number: int) -> dict[str, Any]:
def _day(number: int, today: day_module.Day) -> dict[str, Any]:
"""Строка описи: номер дня, его дата, число событий и хеш байтов.
Дата считается от D0 арифметикой, а не берётся из событий: ось модельного
@@ -72,15 +91,34 @@ def _day(number: int) -> dict[str, Any]:
стенда обрамляют счёт в `ods.event`. Соври она подневная сверка это и
покажет, каждый день сразу.
"""
payloads = serialize.events(day_module.stream(CANONICAL_SEED, number))
payloads = serialize.events(today)
return {
"day": number,
"date": (world.ORIGIN + timedelta(days=number)).isoformat(),
"date": _date(number),
"events": len(payloads),
"sha256": _digest(b"".join(payload + b"\n" for payload in payloads)),
}
def _snapshot(number: int, orders: list[orders_module.Orders]) -> dict[str, Any]:
"""Строка описи слепка: какой день снят, его дата и хеш отправленных байтов.
Байты те же, что уезжают в топик `orders`, с переводом строки после
каждого заказа: пересъёмка слепка обязана дать их снова.
"""
window = [orders[born] for born in orders_module.window(number)]
payloads = serialize.orders(window, number)
return {
"day": number,
"date": _date(number),
"sha256": _digest(b"".join(payload + b"\n" for payload in payloads)),
}
def _date(number: int) -> str:
return (world.ORIGIN + timedelta(days=number)).isoformat()
def _digest(payload: bytes) -> str:
return hashlib.sha256(payload).hexdigest()
@@ -29,6 +29,11 @@
его место занимает 1, а не ноль, иначе «оплатили в секунду рождения» было
бы не отличить от «не оплатили вовсе».
**Состояние на границе суток** чтение готовой судьбы, а не накопление:
слепок дня D учитывает моменты не позже границы D|D+1 и по ним называет
статус. Отсюда «дыхание» окна заказ, оплаченный назавтра, стоит в сегодняшнем
слепке как `created`, а завтра как `paid`.
**Дельта суммы вычеркнутая позиция**: товара не оказалось в наличии, и
заказ приезжает на строку короче клиентской корзины, а `items_total` меньше
ровно на её полную стоимость. Момента у дельты нет склад собрал заказ до
@@ -79,6 +84,10 @@ class Orders:
order_id: tuple[str, ...]
# Пользователь магазина: тот же человек, что стоит за купившей кукой.
user_id: NDArray[np.uint64]
# Когда строка заказа создана в базе источника: секунда той самой покупки
# — в модели строка создаётся синхронно с ней — и миллисекунда часов базы.
# От этого момента отсчитываются и моменты судьбы.
created_at: NDArray[np.datetime64]
# Позиции заказа: номера товаров каталога и штуки, ячейка на заказ. У
# заказа с дельтой позиций на одну меньше, чем в корзине клиента.
product: tuple[NDArray[np.int64], ...]
@@ -104,11 +113,13 @@ def of_day(seed: int, day: int, purchases: Purchases) -> Orders:
delivery = _delivery(rng, len(purchases))
outcome, paid_after, cancelled_after = _fate(rng, len(purchases))
product, quantity, items_total = _delta(rng, purchases)
created_at = _created_at(rng, purchases)
discount = _discount(purchases)
return Orders(
day=day,
order_id=purchases.order_id,
user_id=purchases.person_id,
created_at=created_at,
product=product,
quantity=quantity,
items_total=items_total,
@@ -121,6 +132,49 @@ def of_day(seed: int, day: int, purchases: Purchases) -> Orders:
)
def window(day: int) -> range:
"""Дни рождения, чьи заказы несёт слепок дня `day`.
Окно изменяемости константа мира; у начала оси оно усекается само, а не
сторожем: дней до D0 попросту нет.
"""
return range(max(0, day - world.ORDER_WINDOW_DAYS + 1), day + 1)
def at_boundary(rows: Orders, day: int) -> tuple[list[str], NDArray[np.datetime64]]:
"""Статус заказов и момент их последнего изменения на границе `day`|`day+1`.
Судьба решена при рождении, поэтому слепок её только читает: момент позже
границы для него ещё не случился. Отмена перевешивает оплату у дороги
«оплачен и отменён» она поздняя, и заказ на границе уже отменён.
`updated_at` поздний учтённый момент, а без единого заказ показывает своё
рождение: строку с тех пор никто не трогал.
"""
edge = _boundary(day)
paid = rows.created_at + rows.paid_after.astype("timedelta64[s]")
cancelled = rows.created_at + rows.cancelled_after.astype("timedelta64[s]")
# Момента, которого у исхода нет, в данных нет вовсе: там −1, и без маски
# он прикинулся бы моментом за секунду до рождения.
got_paid = (rows.paid_after >= 0) & (paid <= edge)
got_cancelled = (rows.cancelled_after >= 0) & (cancelled <= edge)
status = np.where(got_cancelled, "cancelled", np.where(got_paid, "paid", "created"))
updated = np.where(
got_cancelled, cancelled, np.where(got_paid, paid, rows.created_at)
)
return status.tolist(), updated
def _boundary(day: int) -> np.datetime64:
"""Граница суток `day`|`day+1` абсолютной меткой: полночь пояса счётчика.
Модельные сутки считаются в поясе счётчика, а моменты заказа метки UTC:
между ними ровно смещение пояса.
"""
midnight = np.datetime64(world.ORIGIN, "s") + np.timedelta64(day + 1, "D")
return midnight - np.timedelta64(world.COUNTER_TIMEZONE_MINUTES, "m")
def _delivery(rng: np.random.Generator, orders: int) -> NDArray[np.int64]:
"""Стоимость доставки каждого заказа дня — броском по таблице весов."""
return _DELIVERY_PRICE[pick(rng, _DELIVERY_CUMULATIVE, orders)]
@@ -147,6 +201,22 @@ def _fate(
return outcome, paid, cancelled
def _created_at(
rng: np.random.Generator, purchases: Purchases
) -> NDArray[np.datetime64]:
"""Момент создания строки заказа: секунда покупки и миллисекунда часов базы.
Секунда приходит из события строка создаётся синхронно с покупкой, и
сдвигать её значило бы подделывать аудит источника. Миллисекунду трекер не
видит вовсе: у него своё разрешение, у базы своё, и три дописанных нуля
выдали бы секундную модель за миллисекундную (исследование формата слепка).
"""
millisecond = rng.integers(0, 1000, len(purchases))
return purchases.moment.astype("datetime64[ms]") + millisecond.astype(
"timedelta64[ms]"
)
def _moment(rng: np.random.Generator, orders: int) -> NDArray[np.int64]:
"""Момент внутри окна: час по таблице весов и равномерная секунда в нём.
@@ -16,6 +16,11 @@
(раздел 5): это разные машины разной природы, и сложенные в одно число они
перестают что-либо говорить. Сериализация считается частью генерации: она
рождает те самые байты, которые сторожит опись.
**Слепки заказов идут третьим ходом того же проигрывателя** и живут по правилу
сдвига: прогон дня D отправляет слепок дня D1 ночная выгрузка бэкенда за
вчера. Правило записано здесь одно и целиком, потому что оно и есть разница
между двумя источниками: трекер шлёт сегодняшний день, бэкенд вчерашний.
"""
import logging
@@ -26,6 +31,7 @@ import numpy as np
from numpy.typing import NDArray
from clickstream_generator import day as day_module
from clickstream_generator import orders as orders_module
from clickstream_generator import serialize
from clickstream_generator.sinks import Sink
@@ -103,6 +109,70 @@ def play(
)
def play_snapshots(sink: Sink, seed: int, first_day: int, days: int = 1) -> None:
"""Отправить слепки прогонов `days` дней подряд начиная с `first_day`.
Прогон дня D отправляет слепок дня D1: содержимое слепка чистая функция
зерна и дня, от момента отправки оно не зависит, а сдвиг делает живой день
обычным. На старте оси вчера нет, поэтому прогон дня 0 не отправляет
ничего.
Слепок собирается переигровкой дней своего окна: заказы дня производная
всей воронки, дешёвого пути к ним нет. Окна соседних слепков перекрываются
почти целиком, и внутри одного прогона день играется один раз иначе
стартовый диапазон стоил бы полусотни проигрышей вместо восьми. Между
прогонами не остаётся ничего: кэш на томе был бы состоянием, которого у
проигрывателя нет.
"""
played: dict[int, orders_module.Orders] = {}
sent = 0
generated = 0.0
delivered = 0.0
for number in range(first_day, first_day + days):
taken = number - 1
if taken < 0:
continue
clock = time.monotonic()
window = [
_orders_of(played, seed, born) for born in orders_module.window(taken)
]
payloads = serialize.orders(window, taken)
spent = time.monotonic() - clock
generated += spent
log.info(
"слепок дня %d: заказов %d, генерация %.1f с", taken, len(payloads), spent
)
clock = time.monotonic()
_send(sink, payloads)
sink.flush()
spent = time.monotonic() - clock
delivered += spent
sent += len(payloads)
log.info(
"слепок дня %d: отправлено %d, доставка %.1f с", taken, len(payloads), spent
)
log.info(
"итого отправлено %d заказов: генерация %.1f с, доставка %.1f с",
sent,
generated,
delivered,
)
def _orders_of(
played: dict[int, orders_module.Orders], seed: int, number: int
) -> orders_module.Orders:
"""Заказы дня `number`, сыгранного один раз на весь прогон."""
if number not in played:
played[number] = day_module.stream(seed, number).orders
return played[number]
def _send(sink: Sink, payloads: list[bytes]) -> None:
"""Пакетно: подряд и без пауз."""
for payload in payloads:
@@ -27,16 +27,25 @@
делает это одним вызовом на колонку, и на дне в полсотни тысяч событий разница
заметна. Обратная сторона день лежит в памяти дважды; проигрыватель поэтому
и берёт его днями, а не горизонтом целиком.
**Второй контракт провода слепок заказов** (мастер-спека, раздел 2). Он не
похож на событие: одиннадцать ключей вместо сорока семи, деньги строками, а
не числами, времена с миллисекундами. Общее у них одно, зато главное: байты
рождаются здесь и только здесь. Запись слепка один словарь с вложенным
списком и один `orjson.dumps`.
"""
from collections.abc import Sequence
from datetime import timedelta
from typing import Any
import numpy as np
import orjson
from numpy.typing import NDArray
from clickstream_generator import schema
from clickstream_generator import catalog, schema, world
from clickstream_generator.day import Day
from clickstream_generator.orders import Orders, at_boundary
_ARRAY_PREFIX = "Array("
@@ -60,6 +69,78 @@ def events(day: Day, limit: int | None = None) -> list[bytes]:
]
def orders(window: Sequence[Orders], day: int) -> list[bytes]:
"""Канонические байты слепка дня `day`: по документу JSON на заказ.
`window` заказы дней окна, от раннего дня к позднему: слепок несёт их
подряд, и порядок строк выходит порядком рождения заказов, он же
возрастание `order_id`. Какие это дни, решает `orders.window`.
Деньги уезжают строками с ровно двумя знаками, а не числами: у заказа они
станут `Decimal`, и дробь двоичного числа была бы потерей точности до
всякого разбора. Времена метки UTC с миллисекундами; `snapshot_date`
одинакова во всей выгрузке это дата дня, состояние которого снято.
"""
goods = catalog.catalog()
sku = goods.sku.tolist()
prices = [_money(price) for price in goods.price.tolist()]
snapshot_date = (world.ORIGIN + timedelta(days=day)).isoformat()
payloads = []
for rows in window:
status, updated = at_boundary(rows, day)
created_at = _moments(rows.created_at)
updated_at = _moments(updated)
user_id = rows.user_id.tolist()
items_total = rows.items_total.tolist()
discount = rows.discount.tolist()
delivery = rows.delivery.tolist()
total = rows.total.tolist()
for number, order_id in enumerate(rows.order_id):
payloads.append(
orjson.dumps(
{
"order_id": order_id,
"user_id": user_id[number],
"status": status[number],
"created_at": created_at[number],
"updated_at": updated_at[number],
"items_total": _money(items_total[number]),
"discount": _money(discount[number]),
"delivery": _money(delivery[number]),
"total": _money(total[number]),
"items": [
{"sku": sku[item], "qty": count, "price": prices[item]}
for item, count in zip(
rows.product[number].tolist(),
rows.quantity[number].tolist(),
strict=True,
)
],
"snapshot_date": snapshot_date,
}
)
)
return payloads
def _money(kopecks: int) -> str:
"""Копейки — строкой с ровно двумя знаками: `129990` → `1299.90`."""
return f"{kopecks // 100}.{kopecks % 100:02d}"
def _moments(values: NDArray[np.datetime64]) -> list[str]:
"""Метки времени — строками RFC 3339 в UTC с миллисекундами.
Три знака стоят всегда, в том числе `.000`: одинаковая длина дробной части
и одинаковая зона дают хронологическую сортировку простым сравнением строк,
а разбор в хранилище идёт по точному шаблону.
"""
ms: NDArray[np.datetime64] = values.astype("datetime64[ms]")
return np.datetime_as_string(ms, unit="ms", timezone="UTC").tolist()
def _values(column: schema.Column, values: NDArray[Any]) -> list[Any]:
"""Колонка питоновскими значениями — в той записи, в какой уедет на провод.