Перейти к содержимому

26. Событийный конвейер

Предыдущая глава моделировала request/response SaaS-бэкенд, где доминирующий стиль — синхронные вызовы. В этой главе берётся система, устроенная наоборот: событийный конвейер обработки данных, где почти каждое взаимодействие асинхронно, где архитектурные вопросы другие — и где соблазн нарисовать брокер узлом посередине самый сильный.

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

ProductApp.rawEvents ─┐
StoreFront.rawEvents ─┴─► Enricher (consumes both)
Enricher.enrichedEvents ─┬─► Aggregator ─► Metrics.push
└─► Warehouse
Warehouse.landedEvents ─► Dashboard
Every event interface carries aspect broker: EventStream
EventStream (kafka_cluster) is on the infra plane — reached by aspect, not by an edge.

Семь модулей: два эмиттера (ProductApp, StoreFront), один обогатитель, один агрегатор, поверхность метрик, хранилище и дашборд. Четыре событийных интерфейса, каждый принадлежит своему производителю (rawEvents ×2, enrichedEvents, landedEvents). Каждая межмодульная связь асинхронна, и процесс ниже несёт весь поток.

analytics/
├── package.archspace
├── sources.arch # ProductApp, StoreFront (event emitters)
├── pipeline.arch # Enricher, Aggregator, Metrics (transformers)
├── sinks.arch # Warehouse, Dashboard
├── infra.arch # EventStream (shared broker, infra plane)
├── processes.arch # the ingest flow + admin touchpoints
└── views.arch

Манифест:

package: acme.analytics
version: "0.1.0"
// Explicit working set — never `use *` the stdlib.
use frontend, service, database, message_broker, kafka_cluster,
kafka, grpc_unary, rest_read, db_read
from arch.backend
use user from arch.extras

sources.arch:

frontend #pa31 ProductApp {
aspect team: "Product"
repo.url: "https://github.com/acme/product-app"
"Customer-facing product. Emits clickstream and lifecycle events."
aspect {
domain: "Analytics"
security.zone: "Internal"
}
kafka rawEvents {
"page_view, click, purchase, signup, churn"
aspect {
broker: EventStream
topic: "analytics.product.raw.v1"
}
}
}
frontend #sf72 StoreFront {
aspect team: "Storefront"
repo.url: "https://github.com/acme/storefront"
"E-commerce frontend. Emits its own clickstream into the same pipeline."
aspect {
domain: "Analytics"
security.zone: "Internal"
}
kafka rawEvents {
"cart, checkout, fulfillment events"
aspect {
broker: EventStream
topic: "analytics.storefront.raw.v1"
}
}
}

Два эмиттера, каждый владеет своим rawEvents. Событие живёт на производителе, потому что производитель владеет схемой — контрактом. Два события — это разные интерфейсы (ProductApp.rawEvents, StoreFront.rawEvents); потребитель читает каждое по имени. Нет общего узла «RawEvents» и нет обработчика на каждый источник, который нужно объявить на стороне потребителя, — потребитель просто проводит ребро к каждому производителю в процессе.

pipeline.arch:

service #en44 Enricher {
aspect team: "Data"
repo.url: "https://github.com/acme/enricher"
"Joins raw events with user-profile attributes, emits enriched events."
aspect {
domain: "Analytics"
security.zone: "Internal"
}
kafka enrichedEvents {
"Raw event + resolved user_id, session_id, plan_tier."
aspect {
broker: EventStream
topic: "analytics.enriched.v1"
}
}
}
service #ag17 Aggregator {
aspect team: "Data"
repo.url: "https://github.com/acme/aggregator"
"Rolls enriched events into minute / hour / day buckets, pushes to Metrics."
aspect {
domain: "Analytics"
security.zone: "Internal"
}
}
service #me08 Metrics {
aspect team: "Data"
repo.url: "https://github.com/acme/metrics"
"Tier-1 metrics surface — counters and rollups."
aspect {
domain: "Analytics"
security.zone: "Internal"
}
grpc_unary push { "Receive rolled-up metric increments from Aggregator." }
}

Enricher и Aggregator — чистые преобразователи: каждый читает вышестоящее и (для Enricher) эмитит нижестоящее. У них нет входящего интерфейса-обработчика — быть запущенным событием означает ребро процесса к производителю, а не объявленную подписку. Metrics здесь — единственный синхронный приёмник: Aggregator вызывает push напрямую.

sinks.arch:

database #wh90 Warehouse {
aspect team: "Data"
"Long-term analytical store. Append-only, partitioned by event date."
aspect {
domain: "Analytics"
data.classification: "pii"
engine: "clickhouse"
host: "analytics-cluster"
}
catalog.url: "https://datahub.acme.internal/dataset/warehouse"
console.url: "https://clickhouse.acme.internal"
kafka landedEvents {
"Emitted after a successful land. Downstream consumers read here."
aspect {
broker: EventStream
topic: "analytics.landed.v1"
}
}
grpc_unary replay { "Re-emit archived events for a date range. Admin only." }
db_read query { "Ad-hoc analytical reads." }
}
frontend #db55 Dashboard {
aspect team: "Data"
repo.url: "https://github.com/acme/dashboard"
"Real-time analytics dashboard. Refreshes as new data lands."
aspect {
domain: "Analytics"
security.zone: "Internal"
}
rest_read loadDashboard { "Initial page load — render current rollups." }
}

Warehouse и потребляет (путь приземления — это ребро процесса к Enricher.enrichedEvents), и производит (landedEvents, который читает Dashboard). Узел с двумя ролями — принять вышестоящее, транслировать нижестоящему — обычен на стыках конвейера, и событие, принадлежащее производителю, делает это естественным: входящая сторона — просто ребро в процессе, исходящая сторона — событие, которым он владеет.

Обратите внимание на цепочку хостинга на Warehouse, выраженную целиком аспектами: engine: "clickhouse" именует плоскость СУБД, host: "analytics-cluster" именует плоскость сервера/кластера. Каждая — отдельная плоскость существования, поэтому каждая — это аспект, а не вложение, ровно как broker. Вложение неверно сказало бы «часть домена»; аспект верно говорит «развёрнут на». (Используйте голое значение — host: AnalyticsCluster — когда хотите материализовать эту плоскость как настоящий, рисуемый модуль.)

infra.arch:

// The broker every event rides. It is a real module — but it lives on
// the infra plane and is reached through the `broker` aspect, never by
// routing a call through it. The EventBackbone view toggles its aspect.
kafka_cluster #ev21 EventStream {
aspect team: "Data Platform"
aspect engine: "kafka"
"Shared streaming backbone for the analytics pipeline."
console.url: "https://kafka.acme.internal/clusters/analytics"
}

Ошибка, которой стоит избегать. Не маршрутизируйте поток через брокер — Producer → EventStream → Consumer — это антипаттерн. Он хоронит архитектуру, которая важна (кто от кого зависит), под транспортной механикой, и каждый модуль в итоге указывает на один и тот же узел брокера, так что диаграмма не говорит ничего. Брокер не отсутствует в модели — это настоящий модуль kafka_cluster — но он сидит на инфра-плоскости и подключается аспектом broker (aspect broker: EventStream) на каждом событийном интерфейсе. Раскрывайте его по запросу оверлей-проекцией; никогда не помещайте его в путь вызова.

Асинхронный поток данных — не латентная разводка, это процесс. Каждый потребитель проводит ребро к событийному интерфейсу производителя; порядок шагов и есть конвейер:

process #pi63 AnalyticsIngest {
"Producer-owned events. A consumer edges to the producer it reads from;
the broker carries each hop on a separate plane (aspects, not steps)."
Enricher > ProductApp.rawEvents
Enricher > StoreFront.rawEvents
Aggregator > Enricher.enrichedEvents
Aggregator > Metrics.push
Warehouse > Enricher.enrichedEvents
Dashboard > Warehouse.landedEvents
}

Шесть шагов, каждый модуль на доске, поток данных читается сверху вниз. Ребро зависимости указывает потребитель→производитель (Enricher > ProductApp.rawEvents = «Enricher зависит от события ProductApp»); рендерер переворачивает стрелку, чтобы данные читались естественным образом. Этот переворот — забота рендера, а не моделирования.

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

Даже у событийного конвейера есть синхронные моменты — обычно администрирование:

user #u0e1 DataEngineer
user #u0u2 Viewer
process #bf02 BackfillReplay {
DataEngineer > Warehouse.replay
}
process #dl04 DashboardLoad {
Viewer > Dashboard.loadDashboard
}

Два процесса на две синхронные точки соприкосновения — ручной бэкфилл и начальная загрузка дашборда.

view #pd01 PipelineDataflow {
"Full pipeline — sources, transformers, sinks. Async edges follow the ingest process."
show @@domain:"Analytics"
}
view #pi02 PIIScope {
"Every module touching PII data. Used for data-classification audits."
show @@data.classification:"pii"
group by @@team
}
view #eb03 EventBackbone {
"The infra overlay: every event interface that rides the shared broker."
show @@broker:EventStream
}

PIIScope работает благодаря аспекту data.classification: "pii" на Warehouse; каскад протаскивает её ко всему вложенному, а show-фильтр делает остальное. EventBackbone — это вознаграждение за хранение транспорта на аспекте: show @@broker:EventStream раскрывает плоскость брокера без того, чтобы эта плоскость когда-либо засоряла проекцию потока данных.

После валидации:

  • Семь модулей, четыре событийных интерфейса, принадлежащих производителю, один общий брокер на инфра-плоскости.
  • Диаграмма с направленными рёбрами, следующими за процессом приёма, — источники сверху, приёмники снизу, пунктирные рёбра для асинхронных связей.
  • Три проекции: dataflow для инженеров, область PII для governance и оверлей брокера для инфраструктуры.

Там, где SaaS-бэкенд (Глава 25) доминировался синхронным запрос/ответ, этот конвейер доминируется асинхронными событиями — но моделирование в обоих одинаково: события живут на своём производителе, потребители проводят к ним ребро в процессах, а брокер — это аспект, а не путевая точка.

Где живёт брокер. На аспекте (aspect broker: EventStream) — никогда как узел в пути вызова. Он остаётся настоящим модулем, так что вы можете владеть им, ссылаться на него и накладывать оверлеем; он просто не загрязняет поток данных.

Кто владеет событием. Производитель. Событийный интерфейс живёт на модуле, который его эмитит, потому что этот модуль владеет схемой. Добавление потребителя тогда — чистое дополнение на стороне потребителя: файл производителя нетронут, ровно как в настоящем pub/sub.

Когда поток заслуживает собственного узла. Правило запрещает транспортную путевую точку, а не доменный поток. Если канал несёт идентичность, которую архитектура учитывает, — межкомандный событийный позвоночник со своей схемой и SLA, которым владеет платформенная команда, — этот поток является полноправным узлом, к которому и публикующие, и подписчики проводят ребро. Дискриминатор — «есть ли у канала собственная идентичность», то же суждение, которое вы выносите для базы данных: тупое хранилище вкладывается или получает аспект; продукт данных получает узел.

Цепочка хостинга. Вложенное или отдельное хранилище — не дно стека. Warehouse → engine (СУБД) → host (кластер) сцепляется аспектами, каждый прыжок — отдельная плоскость. Голые значения материализуют плоскость как узел; строковые значения держат её как аспектный оверлей.

Backpressure и повторные попытки. Не в модели. Это операционные заботы. Модель фиксирует кто что потребляет; насколько надёжно — это другой артефакт (документ SLO, runbook по реакции на инциденты).

  • События принадлежат своему производителю; потребитель проводит ребро к событийному интерфейсу производителя в процессе. Никакой косвенности subscribes:.
  • Асинхронный поток — это процесс: порядок шагов и есть конвейер, и каждый модуль доказывает, что он не мёртвый код, появляясь в нём.
  • Брокер — это настоящий модуль на инфра-плоскости, привязанный через aspect broker. Никогда не маршрутизируйте путь вызова через него.
  • Стек хостинга (БД → СУБД → сервер) сцепляется аспектами, каждый прыжок — отдельная плоскость.
  • Проекции для потока данных, governance (PII) и оверлея брокера — все вытекают из каскадирующих аспектов.

Глава 27: Внешняя интеграция → — ваша система плюс сторонний API плюс их webhooks, точно моделируем границу.