refactor(orders): SQL приёма вынесен из DAG

- Зачем:
  - преобразования хранилища должны читаться рядом с целевым слоем, а DAG должен показывать оркестрацию.
- Что:
  - запросы забора и разбора заказов разложены по каталогам STG и ODS.
  - общий контракт провода подключён в обе ветви штатным шаблонизатором Airflow.
  - SQL смонтирован во все службы Airflow, а правило раскладки записано в архитектуре.
- Проверка:
  - make lint config-test smoke check-services.
  - airflow tasks render для pull_batch и parse_batch; airflow tasks test для parse_batch.
This commit is contained in:
2026-08-18 23:55:04 +03:00
parent bdcbe48b59
commit a5508bd9ec
8 changed files with 195 additions and 165 deletions
+4
View File
@@ -71,6 +71,9 @@ x-airflow-common: &airflow-common
KAFKA_ORDERS_TOPIC: orders KAFKA_ORDERS_TOPIC: orders
volumes: volumes:
- ./dags:/opt/airflow/dags:ro - ./dags:/opt/airflow/dags:ro
# Запросы дагов: даг называет файл, а текст читает Airflow при исполнении
# задачи. Каталог тот же, что применяет DDL, — вся SQL стенда живёт в sql/.
- ./sql:/opt/airflow/sql:ro
- ./infra/airflow/init.sh:/opt/airflow/init.sh:ro - ./infra/airflow/init.sh:/opt/airflow/init.sh:ro
- airflow_logs:/opt/airflow/logs - airflow_logs:/opt/airflow/logs
- airflow_auth:/opt/airflow/auth - airflow_auth:/opt/airflow/auth
@@ -391,6 +394,7 @@ services:
volumes: volumes:
- /var/run/docker.sock:/var/run/docker.sock - /var/run/docker.sock:/var/run/docker.sock
- ./dags:/opt/airflow/dags:ro - ./dags:/opt/airflow/dags:ro
- ./sql:/opt/airflow/sql:ro
- ./infra/airflow/init.sh:/opt/airflow/init.sh:ro - ./infra/airflow/init.sh:/opt/airflow/init.sh:ro
- airflow_logs:/opt/airflow/logs - airflow_logs:/opt/airflow/logs
- airflow_auth:/opt/airflow/auth - airflow_auth:/opt/airflow/auth
+19 -161
View File
@@ -24,159 +24,9 @@ from airflow.sdk import Connection, dag, get_current_context, task
START_DATE = datetime.datetime(2026, 1, 1, tzinfo=datetime.UTC) START_DATE = datetime.datetime(2026, 1, 1, tzinfo=datetime.UTC)
# Забор — один прямой SELECT, без цикла до пустоты: одна порция ClickHouse # Запросы лежат файлами в sql/ репозитория, разложенные по слоям хранилища;
# берёт десятки тысяч сообщений, а слепок дня — порядка полутора тысяч строк. # сюда их монтирует compose — всем службам Airflow сразу.
# Короткая порция оставит хвост до следующего прогона, а отказ после чтения SQL_ROOT = "/opt/airflow/sql"
# унесёт прочитанное с собой: офсеты коммитятся в момент чтения. Граница
# целиком — ADR 0008, «Следствия».
#
# stream_like_engine_allow_direct_select разрешает читать чтеца запросом; вторая
# половина пары объявлена на самой таблице (sql/ddl/10-stg-tables.sql).
# distributed_foreground_insert = 1 — конвенция ETL-вставок стенда: задача не
# должна зеленеть раньше, чем строки легли на шарды.
TAKE_ONE_BATCH = """
INSERT INTO stg.orders_raw_dist
SELECT
raw,
_topic AS kafka_topic,
_partition AS kafka_partition,
_offset AS kafka_offset,
_timestamp_ms AS kafka_timestamp,
hostName() AS consumer_host,
{load_id:String} AS _load_id,
now64(3) AS _load_ts
FROM stg.orders_raw_kafka
SETTINGS
stream_like_engine_allow_direct_select = 1,
distributed_foreground_insert = 1
"""
# Строгий приём: проверяется форма провода и ничего сверх неё.
#
# Ключей одиннадцать, и проверяются все: десять скалярных значений прямо
# образуют типизированную строку заказа. Внутрь items приём не смотрит —
# содержимое позиций, переходы статуса и равенства сумм остаются ниже границы
# (docs/architecture/orders/ingestion.md, «Граница строгого приёма»).
#
# Блок общий у обеих вставок нарочно: годная ветвь берёт row_is_valid, брак —
# буквальное NOT. Разойдись условия хоть на символ — строка либо задвоится,
# либо исчезнет молча.
#
# Отсюда два запрета на выражения предиката, и оба серьёзные. Первый: ни одно
# не возвращает NULL — трёхзначная логика дала бы строку, которую не берёт ни
# условие, ни его отрицание. Поэтому сравнения дают 0 или 1, а обнуляемый
# разбор заканчивается IS NOT NULL. Второй: ни одно не бросает исключений на
# произвольном raw — иначе одна грязная строка роняет весь переход, ради
# отсутствия чего таблица брака и заведена.
WIRE_CONTRACT = r"""
WITH
-- Каноническое время провода: UTC, ровно три знака долей секунды.
'yyyy-MM-dd\'T\'HH:mm:ss.SSS\'Z\'' AS ts_mask,
JSONType(raw) = 'Object' AS is_object,
arraySort(JSONExtractKeys(raw)) = arraySort([
'order_id', 'user_id', 'status', 'created_at', 'updated_at',
'items_total', 'discount', 'delivery', 'total', 'items',
'snapshot_date'
]) AS keys_match,
JSONType(raw, 'order_id') = 'String'
AND JSONType(raw, 'status') = 'String'
-- Целое у JSONType зовётся двумя именами, и нужны оба: с одним
-- Int64 законный идентификатор за 2^63 уехал бы в брак.
AND JSONType(raw, 'user_id') IN ('Int64', 'UInt64')
AND JSONExtract(raw, 'user_id', 'Nullable(UInt64)') IS NOT NULL
AND JSONType(raw, 'items') = 'Array'
-- Форму держит регулярка, реальность — разбор, и порознь они дырявы:
-- маска берёт «2026-6-3T…» без ведущих нулей, а регулярка пропускает
-- 30 февраля. Обе половины измерены — docs/architecture/storage.md,
-- «Что проверено».
AND arrayAll(k ->
JSONType(raw, k) = 'String'
AND match(JSONExtractString(raw, k),
'^\\d{4}-\\d{2}-\\d{2}T\\d{2}:\\d{2}:\\d{2}\\.\\d{3}Z$')
AND parseDateTime64InJodaSyntaxOrNull(
JSONExtractString(raw, k), ts_mask, 'UTC') IS NOT NULL,
['created_at', 'updated_at'])
-- Деньги — строка с ровно двумя знаками после точки. Регулярка держит
-- форму, разбор — вместимость: у строки из двадцати девяток форма
-- хороша, а в Decimal(18, 2) она не влезает.
AND arrayAll(k ->
JSONType(raw, k) = 'String'
AND match(JSONExtractString(raw, k), '^\\d+\\.\\d{2}$')
AND toDecimal64OrNull(JSONExtractString(raw, k), 2) IS NOT NULL,
['items_total', 'discount', 'delivery', 'total'])
AND JSONType(raw, 'snapshot_date') = 'String'
AND match(JSONExtractString(raw, 'snapshot_date'), '^\\d{4}-\\d{2}-\\d{2}$')
AND parseDateTimeInJodaSyntaxOrNull(
JSONExtractString(raw, 'snapshot_date'), 'yyyy-MM-dd', 'UTC') IS NOT NULL
AS fields_valid,
is_object AND keys_match AND fields_valid AS row_is_valid
"""
# Годные версии заказа. Срез читается по _load_id — тому же, что проставил
# забор: разбор идёт по неизменной порции, а не по «всему, что появилось».
#
# Разобранные значения берёт assumeNotNull: обнуляемый разбор стоит за
# предикатом, который NULL уже отсёк, и приведение здесь не может упасть.
PARSE_GOOD_ROWS = (
"INSERT INTO ods.order_snapshot_dist"
+ WIRE_CONTRACT
+ r"""
SELECT
JSONExtractString(raw, 'order_id') AS order_id,
JSONExtract(raw, 'user_id', 'UInt64') AS user_id,
JSONExtractString(raw, 'status') AS status,
assumeNotNull(parseDateTime64InJodaSyntaxOrNull(
JSONExtractString(raw, 'created_at'), ts_mask, 'UTC')) AS created_at,
assumeNotNull(parseDateTime64InJodaSyntaxOrNull(
JSONExtractString(raw, 'updated_at'), ts_mask, 'UTC')) AS updated_at,
assumeNotNull(toDecimal64OrNull(
JSONExtractString(raw, 'items_total'), 2)) AS items_total,
assumeNotNull(toDecimal64OrNull(JSONExtractString(raw, 'discount'), 2)) AS discount,
assumeNotNull(toDecimal64OrNull(JSONExtractString(raw, 'delivery'), 2)) AS delivery,
assumeNotNull(toDecimal64OrNull(JSONExtractString(raw, 'total'), 2)) AS total,
-- items кладётся сырым фрагментом JSON, а не разобранной структурой.
JSONExtractRaw(raw, 'items') AS items,
toDate(assumeNotNull(parseDateTimeInJodaSyntaxOrNull(
JSONExtractString(raw, 'snapshot_date'),
'yyyy-MM-dd', 'UTC'))) AS snapshot_date,
-- Метки запуска и прибытия переносятся как есть. Поставь здесь now64(3) —
-- и _load_ts молча ответила бы на другой вопрос: «когда разобрали».
_load_id,
_load_ts
FROM stg.orders_raw_dist
WHERE _load_id = {load_id:String} AND row_is_valid
SETTINGS distributed_foreground_insert = 1
"""
)
# Брак: тот же срез и буквальное отрицание того же предиката.
#
# Классы перекрываются, поэтому проверяются по порядку, а в error_class идёт
# первый совпавший: скаляр проваливает и проверку на объект, и сверку ключей —
# без объявленного порядка он попал бы то в один класс, то в другой.
PARSE_BAD_ROWS = (
"INSERT INTO ods.order_snapshot_errors_dist"
+ WIRE_CONTRACT
+ r"""
SELECT
raw,
multiIf(
NOT is_object, 'not_an_object',
NOT keys_match, 'keyset_mismatch',
'field_invalid'
) AS error_class,
kafka_topic,
kafka_partition,
kafka_offset,
kafka_timestamp,
consumer_host,
_load_id,
_load_ts
FROM stg.orders_raw_dist
WHERE _load_id = {load_id:String} AND NOT row_is_valid
SETTINGS distributed_foreground_insert = 1
"""
)
@dag( @dag(
@@ -188,6 +38,7 @@ SETTINGS distributed_foreground_insert = 1
# дрались бы за неё, а слепок разъехался бы по двум _load_id; второй # дрались бы за неё, а слепок разъехался бы по двум _load_id; второй
# прогон подождёт своей очереди. # прогон подождёт своей очереди.
max_active_runs=1, max_active_runs=1,
template_searchpath=SQL_ROOT,
tags=["заказы"], tags=["заказы"],
) )
def orders_ingest(): def orders_ingest():
@@ -210,12 +61,16 @@ def orders_ingest():
send_receive_timeout=30, send_receive_timeout=30,
) )
@task # Аргумент с расширением из templates_exts Airflow подменяет текстом файла:
def pull_batch() -> None: # берёт его из template_searchpath и прогоняет через Jinja при исполнении
# задачи, а не при разборе дага. В задачу приезжает готовый запрос — вместе
# с тем, что файл подключил через include.
@task(templates_exts=(".sql",))
def pull_batch(sql: str) -> None:
load_id = get_current_context()["run_id"] load_id = get_current_context()["run_id"]
client = clickhouse_client() client = clickhouse_client()
try: try:
summary = client.command(TAKE_ONE_BATCH, parameters={"load_id": load_id}) summary = client.command(sql, parameters={"load_id": load_id})
finally: finally:
client.close() client.close()
# Размер порции — read_rows: written_rows у вставки в Distributed # Размер порции — read_rows: written_rows у вставки в Distributed
@@ -226,8 +81,8 @@ def orders_ingest():
load_id, load_id,
) )
@task @task(templates_exts=(".sql",))
def parse_batch() -> None: def parse_batch(good_rows_sql: str, bad_rows_sql: str) -> None:
"""Разобрать срез сырья в версии заказов и в брак. """Разобрать срез сырья в версии заказов и в брак.
Обе вставки в одном task_id: транзакции между ними ClickHouse не даёт, Обе вставки в одном task_id: транзакции между ними ClickHouse не даёт,
@@ -237,15 +92,18 @@ def orders_ingest():
load_id = get_current_context()["run_id"] load_id = get_current_context()["run_id"]
client = clickhouse_client() client = clickhouse_client()
try: try:
client.command(PARSE_GOOD_ROWS, parameters={"load_id": load_id}) client.command(good_rows_sql, parameters={"load_id": load_id})
client.command(PARSE_BAD_ROWS, parameters={"load_id": load_id}) client.command(bad_rows_sql, parameters={"load_id": load_id})
finally: finally:
client.close() client.close()
# Счётчиков строк нет: у запроса с WHERE read_rows считает прочитанное # Счётчиков строк нет: у запроса с WHERE read_rows считает прочитанное
# с диска, а не подошедшее (storage.md, «Что проверено»). # с диска, а не подошедшее (storage.md, «Что проверено»).
logging.info("срез разобран: _load_id %s", load_id) logging.info("срез разобран: _load_id %s", load_id)
pull_batch() >> parse_batch() pull_batch("stg/orders_raw_load.sql") >> parse_batch(
"ods/order_snapshot_load.sql",
"ods/order_snapshot_errors_load.sql",
)
orders_ingest() orders_ingest()
+8 -2
View File
@@ -38,19 +38,25 @@
дёргает его тот, кто положил слепок в топик, — работник пульта мира, — и ждёт дёргает его тот, кто положил слепок в топик, — работник пульта мира, — и ждёт
конца прогона. Все строки получают `_load_id`, равный `run_id` Airflow. конца прогона. Все строки получают `_load_id`, равный `run_id` Airflow.
`_load_ts` вычисляется при этой записи и дальше переносится без пересчёта. `_load_ts` вычисляется при этой записи и дальше переносится без пересчёта.
Запрос забора лежит в [`sql/stg/orders_raw_load.sql`](../../../sql/stg/orders_raw_load.sql):
даг задаёт порядок и параметры, а преобразование остаётся в SQL своего слоя.
Один следующий `task_id` отвечает за весь переход STG → ODS. Внутри него два Один следующий `task_id` отвечает за весь переход STG → ODS. Внутри него два
последовательных `INSERT SELECT` читают неизменный срез по `_load_id`: первый последовательных `INSERT SELECT` читают неизменный срез по `_load_id`: первый
пишет годные строки в `ods.order_snapshot`, второй — брак в пишет годные строки в `ods.order_snapshot`, второй — брак в
`ods.order_snapshot_errors`. Транзакции между запросами нет. При частичном сбое `ods.order_snapshot_errors`. Транзакции между запросами нет. При частичном сбое
Airflow повторяет весь `task_id`; одинаковые исходные строки и служебные метки Airflow повторяет весь `task_id`; одинаковые исходные строки и служебные метки
не вычисляются заново. не вычисляются заново. Запросы лежат рядом с целями:
[`order_snapshot_load.sql`](../../../sql/ods/order_snapshot_load.sql) и
[`order_snapshot_errors_load.sql`](../../../sql/ods/order_snapshot_errors_load.sql).
Условия запросов взаимодополняющие: один общий предикат определяет брак, а Условия запросов взаимодополняющие: один общий предикат определяет брак, а
годная ветвь использует его буквальное отрицание. Все функции предиката годная ветвь использует его буквальное отрицание. Все функции предиката
возвращают результат без исключения, а сам предикат всегда заканчивается в возвращают результат без исключения, а сам предикат всегда заканчивается в
`true` или `false`, не в `NULL`. Постоянный классификатор между STG и ODS для `true` или `false`, не в `NULL`. Постоянный классификатор между STG и ODS для
этого не нужен. этого не нужен. Обе ветви включают один файл
[`_order_wire_contract.sql`](../../../sql/ods/_order_wire_contract.sql), поэтому
предикат нельзя случайно исправить только в одной из них.
## Граница строгого приёма ## Граница строгого приёма
+18 -2
View File
@@ -414,7 +414,22 @@ kafka_offset)`: смотрят такую таблицу от класса, а
разрастается до имени отдельного поля. Точная граница приёма — в разрастается до имени отдельного поля. Точная граница приёма — в
[спецификации заказов](orders/ingestion.md). [спецификации заказов](orders/ingestion.md).
## Раскладка DDL ## Раскладка SQL
Исполняемые дагами запросы лежат по правилу
`sql/<слой>/<объект>_<роль>.sql`. Слой — слой цели запроса: забор заказов в
сырьё живёт в `sql/stg/orders_raw_load.sql`, две ветви разбора — в `sql/ods/`.
Даг называет файлы, передаёт параметры и задаёт порядок выполнения, но не
хранит текст преобразований. Общий фрагмент начинается с подчёркивания: файл
`sql/ods/_order_wire_contract.sql` сам не исполняется, его включают обе ветви
разбора.
Airflow читает и собирает эти файлы штатным шаблонизатором при исполнении
задачи. Каталог `sql/` смонтирован во все его службы только для чтения. Поэтому
обработчик дагов не зависит от наличия SQL на машине при разборе Python, а
планировщик видит те же файлы при исполнении задачи.
### DDL
Файлы лежат в `sql/ddl/` и применяются по порядку имён. Сначала все статичные Файлы лежат в `sql/ddl/` и применяются по порядку имён. Сначала все статичные
объекты, потом матвью — тогда к моменту создания матвью её цель уже существует. объекты, потом матвью — тогда к моменту создания матвью её цель уже существует.
@@ -488,7 +503,8 @@ ODS. Второе: матвью приёма создаётся последне
| ODS | `ods.order_snapshot_errors_rep` / `_dist` | строки слепка, не прошедшие строгий приём | | ODS | `ods.order_snapshot_errors_rep` / `_dist` | строки слепка, не прошедшие строгий приём |
Матвью разбора у заказов нет: срез сырья раскладывают по этим двум целям два Матвью разбора у заказов нет: срез сырья раскладывают по этим двум целям два
`INSERT SELECT` шага `parse_batch` в даге `orders_ingest`. `INSERT SELECT` шага `parse_batch` из файлов `sql/ods/order_snapshot_load.sql`
и `sql/ods/order_snapshot_errors_load.sql`.
Слои DDS и DM появляются на следующих этапах; их состав задан разделом 7 Слои DDS и DM появляются на следующих этапах; их состав задан разделом 7
мастер-спеки и переносится сюда по мере постройки. мастер-спеки и переносится сюда по мере постройки.
+61
View File
@@ -0,0 +1,61 @@
-- Форма провода одной строки слепка: общая часть обеих вставок ODS.
--
-- Строгий приём: проверяется форма провода и ничего сверх неё.
--
-- Ключей одиннадцать, и проверяются все: десять скалярных значений прямо
-- образуют типизированную строку заказа. Внутрь items приём не смотрит —
-- содержимое позиций, переходы статуса и равенства сумм остаются ниже границы
-- (docs/architecture/orders/ingestion.md, «Граница строгого приёма»).
--
-- Файл включают обе вставки нарочно: годная ветвь берёт row_is_valid, брак —
-- буквальное NOT. Разойдись условия хоть на символ — строка либо задвоится,
-- либо исчезнет молча.
--
-- Отсюда два запрета на выражения предиката, и оба серьёзные. Первый: ни одно
-- не возвращает NULL — трёхзначная логика дала бы строку, которую не берёт ни
-- условие, ни его отрицание. Поэтому сравнения дают 0 или 1, а обнуляемый
-- разбор заканчивается IS NOT NULL. Второй: ни одно не бросает исключений на
-- произвольном raw — иначе одна грязная строка роняет весь переход, ради
-- отсутствия чего таблица брака и заведена.
WITH
-- Каноническое время провода: UTC, ровно три знака долей секунды.
'yyyy-MM-dd\'T\'HH:mm:ss.SSS\'Z\'' AS ts_mask,
JSONType(raw) = 'Object' AS is_object,
arraySort(JSONExtractKeys(raw)) = arraySort([
'order_id', 'user_id', 'status', 'created_at', 'updated_at',
'items_total', 'discount', 'delivery', 'total', 'items',
'snapshot_date'
]) AS keys_match,
JSONType(raw, 'order_id') = 'String'
AND JSONType(raw, 'status') = 'String'
-- Целое у JSONType зовётся двумя именами, и нужны оба: с одним
-- Int64 законный идентификатор за 2^63 уехал бы в брак.
AND JSONType(raw, 'user_id') IN ('Int64', 'UInt64')
AND JSONExtract(raw, 'user_id', 'Nullable(UInt64)') IS NOT NULL
AND JSONType(raw, 'items') = 'Array'
-- Форму держит регулярка, реальность — разбор, и порознь они дырявы:
-- маска берёт «2026-6-3T…» без ведущих нулей, а регулярка пропускает
-- 30 февраля. Обе половины измерены — docs/architecture/storage.md,
-- «Что проверено».
AND arrayAll(k ->
JSONType(raw, k) = 'String'
AND match(JSONExtractString(raw, k),
'^\\d{4}-\\d{2}-\\d{2}T\\d{2}:\\d{2}:\\d{2}\\.\\d{3}Z$')
AND parseDateTime64InJodaSyntaxOrNull(
JSONExtractString(raw, k), ts_mask, 'UTC') IS NOT NULL,
['created_at', 'updated_at'])
-- Деньги — строка с ровно двумя знаками после точки. Регулярка держит
-- форму, разбор — вместимость: у строки из двадцати девяток форма
-- хороша, а в Decimal(18, 2) она не влезает.
AND arrayAll(k ->
JSONType(raw, k) = 'String'
AND match(JSONExtractString(raw, k), '^\\d+\\.\\d{2}$')
AND toDecimal64OrNull(JSONExtractString(raw, k), 2) IS NOT NULL,
['items_total', 'discount', 'delivery', 'total'])
AND JSONType(raw, 'snapshot_date') = 'String'
AND match(JSONExtractString(raw, 'snapshot_date'), '^\\d{4}-\\d{2}-\\d{2}$')
AND parseDateTimeInJodaSyntaxOrNull(
JSONExtractString(raw, 'snapshot_date'), 'yyyy-MM-dd', 'UTC') IS NOT NULL
AS fields_valid,
is_object AND keys_match AND fields_valid AS row_is_valid
+25
View File
@@ -0,0 +1,25 @@
-- Брак: тот же срез и буквальное отрицание того же предиката.
--
-- Классы перекрываются, поэтому проверяются по порядку, а в error_class идёт
-- первый совпавший: скаляр проваливает и проверку на объект, и сверку ключей —
-- без объявленного порядка он попал бы то в один класс, то в другой.
INSERT INTO ods.order_snapshot_errors_dist
{% include "ods/_order_wire_contract.sql" %}
SELECT
raw,
multiIf(
NOT is_object, 'not_an_object',
NOT keys_match, 'keyset_mismatch',
'field_invalid'
) AS error_class,
kafka_topic,
kafka_partition,
kafka_offset,
kafka_timestamp,
consumer_host,
_load_id,
_load_ts
FROM stg.orders_raw_dist
WHERE _load_id = {load_id:String} AND NOT row_is_valid
SETTINGS distributed_foreground_insert = 1
+33
View File
@@ -0,0 +1,33 @@
-- Годные версии заказа. Срез читается по _load_id — тому же, что проставил
-- забор: разбор идёт по неизменной порции, а не по «всему, что появилось».
--
-- Разобранные значения берёт assumeNotNull: обнуляемый разбор стоит за
-- предикатом, который NULL уже отсёк, и приведение здесь не может упасть.
INSERT INTO ods.order_snapshot_dist
{% include "ods/_order_wire_contract.sql" %}
SELECT
JSONExtractString(raw, 'order_id') AS order_id,
JSONExtract(raw, 'user_id', 'UInt64') AS user_id,
JSONExtractString(raw, 'status') AS status,
assumeNotNull(parseDateTime64InJodaSyntaxOrNull(
JSONExtractString(raw, 'created_at'), ts_mask, 'UTC')) AS created_at,
assumeNotNull(parseDateTime64InJodaSyntaxOrNull(
JSONExtractString(raw, 'updated_at'), ts_mask, 'UTC')) AS updated_at,
assumeNotNull(toDecimal64OrNull(
JSONExtractString(raw, 'items_total'), 2)) AS items_total,
assumeNotNull(toDecimal64OrNull(JSONExtractString(raw, 'discount'), 2)) AS discount,
assumeNotNull(toDecimal64OrNull(JSONExtractString(raw, 'delivery'), 2)) AS delivery,
assumeNotNull(toDecimal64OrNull(JSONExtractString(raw, 'total'), 2)) AS total,
-- items кладётся сырым фрагментом JSON, а не разобранной структурой.
JSONExtractRaw(raw, 'items') AS items,
toDate(assumeNotNull(parseDateTimeInJodaSyntaxOrNull(
JSONExtractString(raw, 'snapshot_date'),
'yyyy-MM-dd', 'UTC'))) AS snapshot_date,
-- Метки запуска и прибытия переносятся как есть. Поставь здесь now64(3) —
-- и _load_ts молча ответила бы на другой вопрос: «когда разобрали».
_load_id,
_load_ts
FROM stg.orders_raw_dist
WHERE _load_id = {load_id:String} AND row_is_valid
SETTINGS distributed_foreground_insert = 1
+27
View File
@@ -0,0 +1,27 @@
-- Забор порции слепка заказов из Kafka в сырьё STG.
--
-- Забор — один прямой SELECT, без цикла до пустоты: одна порция ClickHouse
-- берёт десятки тысяч сообщений, а слепок дня — порядка полутора тысяч строк.
-- Короткая порция оставит хвост до следующего прогона, а отказ после чтения
-- унесёт прочитанное с собой: офсеты коммитятся в момент чтения. Граница
-- целиком — ADR 0008, «Следствия».
--
-- stream_like_engine_allow_direct_select разрешает читать чтеца запросом; вторая
-- половина пары объявлена на самой таблице (sql/ddl/10-stg-tables.sql).
-- distributed_foreground_insert = 1 — конвенция ETL-вставок стенда: задача не
-- должна зеленеть раньше, чем строки легли на шарды.
INSERT INTO stg.orders_raw_dist
SELECT
raw,
_topic AS kafka_topic,
_partition AS kafka_partition,
_offset AS kafka_offset,
_timestamp_ms AS kafka_timestamp,
hostName() AS consumer_host,
{load_id:String} AS _load_id,
now64(3) AS _load_ts
FROM stg.orders_raw_kafka
SETTINGS
stream_like_engine_allow_direct_select = 1,
distributed_foreground_insert = 1