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.analyticsversion: "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.backenduse user from arch.extrassources.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 DataEngineeruser #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, точно моделируем границу.