---
title: Мечта о масштабируемых и обогащённых GraphQL-подписках
date: 2022-11-23T15:22
tags: [ai, frontend, go, graphql, php]
---

Впервые опубликовано на [Medium](https://medium.com/pipedrive-engineering/a-dream-of-scalable-and-enriched-graphql-subscriptions-724284448e65)

![Стилизованная фотография водопада Ягала в Эстонии](https://miro.medium.com/max/1400/1*Rvb94EOQA-BzsmYJ4TklGg.jpeg)

Стилизованная фотография водопада Ягала в Эстонии. Оригинал: Александр Абросимов, Wikimedia Commons

[В прошлый раз](/ru/blog/tech/backend/put-k-federativnomu-graphql/) я рассказал о пятилетнем пути GraphQL в Pipedrive. Теперь расскажу о десятилетнем пути доставки событий по WebSocket во фронтенд. Возможно, это поможет и вам.

![Анимация доставки событий во фронтенд](https://miro.medium.com/max/1400/1*xffiu2fshUMbBFm-JlSAGw.gif)

<!-- truncate -->

# **Зачем**

**Потребность продукта** в асинхронных событиях возникает каждый раз, когда серверу нужно о чём-то сообщить пользователю:

- вам пришло мгновенное сообщение;
- завершилось изменение размера загруженного изображения;
- коллега сменил аватар;
- массовое изменение 10 000 сделок выполнено на 99%.

Таким образом, асинхронные события **делают интерфейс информативнее и интерактивнее**.

Но важнее всего то, что они решают проблему **несогласованных данных**, отображаемых или хранящихся в браузере. Без таких обновлений другой пользователь не увидит _переименование сделки в Pipedrive_, а другая вкладка браузера не узнает об _удалённых активностях_. Даже представления в одной вкладке могут плохо взаимодействовать друг с другом, если у них нет общего хранилища, и показывать рассинхронизированные данные.

# **Экскурс в историю**

К счастью, Pipedrive «решил» эту проблему десять лет назад. В 2012 году [Андрис](https://www.linkedin.com/in/andris-reinman/), [Капп](https://www.linkedin.com/in/martinkapp/) и [Таюр](https://www.linkedin.com/in/martintajur/) разработали сервис **socketqueue**, который и сегодня доставляет события API во фронтенд с помощью библиотеки [SockJS](https://www.npmjs.com/package/sockjs).

Доклад Таюра в 2016 году стал одной из причин, по которым я воссоздал похожую систему, задумался, как управлять пользовательскими подключениями с [Socket.IO](https://socket.io/) и масштабировать их за пределы одного сервера, а затем присоединился к Pipedrive, чтобы найти ответ. Меня восхищали потоковая обработка событий и возможности, которые она открывает.

Как видно из доклада, в основном он посвящён [RabbitMQ](https://www.rabbitmq.com/) - брокеру сообщений между PHP-монолитом и socketqueue.

# **Плюсы и минусы**

Очереди в памяти позволяли Pipedrive выдерживать всплески событий от внешних интеграций и внутренних операций, например импорта сделок из XLS-файла или массового редактирования.

Событие API, отправляемое в RabbitMQ и далее во фронтенд, **денормализовано** и содержит разные сущности. Это хорошо: вы получаете **согласованный снимок** нескольких объектов в один момент, а получателям **не приходится запрашивать данные повторно**.

Но это крайне **неэффективно**: все вкладки браузера получают все возможные изменения независимо от того, нужны они им или нет. Для одного клиента трафик иногда достигал 80 ГБ в день, а веб-приложение на его ноутбуке выполняло всю тяжёлую фильтрацию.

Поэтому нельзя бесконечно раздувать схему события, делая его ещё тяжелее.

![Идентификатор ячейки в URL WebSocket](https://miro.medium.com/max/1400/1*SEK3NXeBwKWHvw4fgcbVBg.png)

Логика ячеек: номер ячейки в URL WebSocket

Вторая проблема - **шумный сосед**. Масштабирование, как я уже упоминал, реализовано с помощью «логики ячеек», то есть, в хорошем случае, изоляции арендаторов: серверы приложения socketqueue разделены по идентификатору компании. Это также означает, что для предсказуемой маршрутизации необходимо фиксированное количество контейнеров и очередей.

Проблема возникает, когда одна компания генерирует тысячи событий. Они не распределяются между всеми серверами, а вызывают **скачок нагрузки на CPU** одного узла вплоть до провала проверки состояния. Вертикально нарастить одноядерный CPU уже некуда, поэтому остаётся лишь удвоить количество серверов, что ведёт к **растрате инфраструктурных ресурсов**.

![Скачок CPU отдельного пода из-за шумного соседа](https://miro.medium.com/max/1400/1*2pcMt6Bjn_pgBLO9CLjjOA.png)

Скачок CPU отдельного пода из-за шумного соседа

Третья проблема - **наблюдаемость и трассировка**. Как и с REST API, без хорошей документации, будь то Swagger или определения типов, трудно понять, что может находиться в схеме события. Ещё важнее знать, кто и какие данные использует, если вы хотите внести несовместимое изменение. Возникает порочный круг: никто ничего не удаляет из события, и оно раздувается ещё сильнее.

Наконец, если соединение WebSocket оборвалось из-за того, что пользователь заехал в туннель или просто закрыл ноутбук, он может потерять события. Может ли это привести к потере _сделки_ и прибыли?

За годы мы выжали из socketqueue почти всё. Сервис:

- работает в нескольких потоках, обслуживая очереди и WS-соединения;
- сжимает трафик;
- проверяет права доступа и видимость.

Но что, если можно сделать лучше? 🤔

# **Как**

Идея [GraphQL-подписок](https://spec.graphql.org/June2018/#sec-Subscription-Operation-Definitions) довольно проста: с помощью аргументов фильтрации вы объявляете только те данные, которые хотите получать. Сервер должен разумно отфильтровать события.

Я не смогу объяснить базовую настройку лучше, чем это сделал [Бен](https://twitter.com/benawad) 😁. В его примере есть один экземпляр сервера, где pub/sub просто маршрутизирует события в памяти прямо из мутации. Но в нашем случае мутаций нет: реальные изменения генерирует далёкая-далёкая устаревшая система на PHP.

# **Цели миссии**

Я хотел добиться следующего:

- продемонстрировать GraphQL-подписки в боевой среде с явно заданной схемой;
- проверить горизонтальное масштабирование подов;
- повысить надёжность, используя Kafka вместо RabbitMQ;
- сохранить простоту благодаря небольшим нормализованным доменным событиям;
- удержать задержку доставки событий ниже двух секунд, чтобы не уступать socketqueue, которую мы хотели заменить.

## **Риски**

Я не знал:

- как выполнять аутентификацию;
- какой протокол выбрать: SSE или WS? Что такое [Mercure](https://github.com/dunglas/mercure), и есть ли где-то откат к длительному опросу;
- будут ли соединения WS/TCP привязаны к одному поду;
- можно ли одновременно держать несколько подписок.

![Вопросы о протоколах и подключениях](https://miro.medium.com/max/3920/1*K-aB3FKXH0jXeG4HiLRdVw.png)

![Варианты архитектуры подписок](https://miro.medium.com/max/4004/1*5q23Z56at20SsvRkOB_6Hw.png)

- Что использовать для хранения или передачи: БД, Redis Streams, Kafka, Redis Pub/Sub или KTable?
- Можно ли перематывать события назад, если пользователь отключился? Нужно ли хранить идентификатор курсора для каждой сущности? И как работает live query?
- Как фильтровать события по компании, пользователю, сессии и сущности? Есть ли стандарт для полей фильтрации подписок?
- Можно ли повторно использовать федеративную схему GraphQL и обогащать события? Какими должны быть QoS и логика повторных попыток, если шлюз недоступен? Ли Байрон подробно разобрал это [в своём докладе более пяти лет назад](https://youtu.be/rapO30fpREg?t=6675).
- На сколько объектов можно подписаться?
- Как отписываться или отключать пользователей после выхода из системы?

# **Прототип**

До начала миссии я написал простой сервис на Node.js, который подключался к Kafka и проксировал события без фильтрации. Результат уже выглядел многообещающе. Но во время миссии мне хотелось попробовать [**Go**](https://go.dev/), чтобы задействовать все процессорные ядра и добиться максимальной эффективности. Node.js-прототип оставался запасным вариантом, и позже это оказалось очень полезно.

![Прототип сервиса подписок](https://miro.medium.com/max/2256/1*8i540taseWXr_TAV2ysBZw.png)

![Поток событий через прототип](https://miro.medium.com/max/4952/1*z39G_eYPNR7Zuu_Hl1dK7g.png)

# **Запуск**

Мы вчетвером ([я](https://www.linkedin.com/in/kurapov/), [Павел](https://www.linkedin.com/in/pavel-nikolajev-84184a64/), [Кристьян](https://www.linkedin.com/in/kristjan-luik-89bb03122/) и [Хиро](https://www.linkedin.com/in/abhishek-goswami-591541b1/)) отправились на всё лето исследовать неизвестность 🛸 и привезти полезный результат обратно на стартовую площадку.

![Команда миссии](https://miro.medium.com/max/1400/1*QlBXCk8EjRtp5n_Fxr2xTw.png)

## **GraphQL Go от graph-gophers**

Сначала мы взяли за основу [graph-gophers/graphql-go](https://github.com/graph-gophers/graphql-go) и за несколько дней воссоздали прототип подписки на события Kafka:

- Столкнулись с тем, что целочисленный тип из спецификации GraphQL [не соответствует int32 в Go](https://github.com/graphql/graphql-js/issues/292#issuecomment-186702912).
- Реализовали фильтрацию по свойствам и аргументам.
- Выяснили, что для WSS нужен SSL, иначе CORS блокирует запросы.
- Настроили [проксирование запросов WebSocket через Nginx](http://nginx.org/en/docs/http/websocket.html).
- Обнаружили два транспортных протокола. WebSocket задаёт транспорт довольно свободно, поэтому каждая библиотека может реализовать его по-своему. [Gophers реализовали](https://github.com/graph-gophers/graphql-transport-ws) старую версию, совместимую с Apollo и GraphiQL, но не с graphql-ws.
- Столкнулись с тем, что библиотека Apollo Federation не принимала `Subscription` при [регистрации схемы](https://github.com/pipedrive/graphql-schema-registry). Мы поняли, что можно удалить из корня только тип Subscription, продолжая проверять остальные типы на конфликты.

Базовая фильтрация оказалась довольно простой. Нужно соединить два потока, или канала: Kafka и WebSocket, который Gophers использует в резолверах схемы. Между ними работает goroutine, преобразующая данные.

![Каналы Kafka и WebSocket](https://miro.medium.com/max/3712/1*zXxbPFcX5Xqu7zGz8k3kSw.png)

![Фильтрация событий в Go](https://miro.medium.com/max/3668/1*0dXQgXJcVTET0AswerSujA.png)

Мы столкнулись с серьёзной [проблемой в Gophers](https://github.com/graph-gophers/graphql-transport-ws/pull/9), связанной с JWT-аутентификацией. Нам не хватало опыта в _Go_ и времени на изменение крупного фреймворка.

Кроме того, мы поняли, что обогащение событий пока не в приоритете. Нужно было вывести код в боевую среду 🚀 и проверить масштабирование. Таков Agile-подход.

Вторым серьёзным препятствием стала проверка видимости сущностей. При чтении доменных событий из Kafka у нас нет дополнительной информации, например о том, кто может видеть конкретную сделку. Эти сведения нужно запрашивать у другого сервиса. Но запрос на каждое событие каждой компании легко может привести к DoS 🤦‍♂️

## **Gqlgen**

Мы отказались от Gophers и перешли на gqlgen. Логика обмена каналами осталась прежней, а JWT-аутентификация заработала. Теперь события подписок автоматически генерировались на основе `schema.graphql`.

Затем я обнаружил недостаток exchange, то есть внутрипроцессного pub/sub в Go: для начала потоковой передачи нужна как минимум одна подписка. Новая версия стала немного лучше и позволяла иметь отдельный канал для каждой компании.

![Архитектура на gqlgen](https://miro.medium.com/max/2532/1*_tIlshM1g5F8tYgIN4qAfQ.png)

![Обмен каналами в gqlgen](https://miro.medium.com/max/2176/1*w6B9z69_upq05WXXb2Vt5Q.png)

Мы также начали продумывать схему, которая масштабируется на разные сущности, поддерживает фильтры, несколько действий над сущностями и несколько идентификаторов.

Прошёл месяц. Интерфейс работал с неидеальной, но производительной проверкой прав. Мы ушли в отпуск остыть 🌴 И только тогда я увидел поразительный 🤯 доклад [Мэнди Уайз](https://twitter.com/mandiwise) из Apollo, хотя к тому моменту уже пересмотрел почти все видео на YouTube.

## **Итоговая архитектура**

Вернувшись, мы снова начали с нуля и отказались от gqlgen. Было больно, но таков Agile. Мы разделили сервис на **два**, поместив между ними Redis Pub/Sub с репликацией и Sentinel.

![Итоговая архитектура GraphQL-подписок](https://miro.medium.com/max/1400/1*JDm6Oo8YfE1-r7lzttAyNQ.png)

- graphql-subscription-workers остался на Go, но стал намного проще. Он хорошо масштабируется для чтения всех партиций Kafka. Воркеры отфильтровывают события, параллельно выполняя до десяти проверок прав, и при необходимости преобразуют схему события.
- graphql-subscriptions обслуживает соединения и масштабируется горизонтально настолько, насколько нужно. Сервис использует отличную библиотеку [graphql-ws](https://github.com/enisdenjo/graphql-ws) со множеством хуков. Спасибо её автору [Денису Бадурине](https://github.com/enisdenjo).
- Redis Pub/Sub служит брокером обмена и выполняет основную фильтрацию и шардирование по компаниям, пользователям и идентификаторам сущностей.

Например, воркер отправляет сообщение в канал `<companyId>.deal.<dealId>`. Это очень похоже на маршрутизацию RabbitMQ, где подписчики тоже могут сопоставлять адреса по шаблону. Теоретически мы можем подписаться на `<companyId>.*` и при необходимости обрабатывать все события компании.

## **Типы в схеме подписок**

Как видно на схеме, мы последовали совету Мэнди и реализовали обогащение событий запросами к шлюзу. Для этого пришлось периодически получать федеративную схему и жёстко задавать, какой запрос должна выполнять конкретная сущность. Зато фронтенд получил огромную гибкость.

![Обогащённое событие подписки](https://miro.medium.com/max/2376/1*zBOSyHn9wocWYbXR7r80gA.gif)

![Типы схемы подписок](https://miro.medium.com/max/3504/1*dauu7v8ZBwU667SdCCBBnA.png)

Нам пришлось принять несколько решений о схеме:

- Мы не объединяем исходное событие с обогащённым, хотя могли бы. Нет гарантии, что GraphQL-сервис ответит вовремя, а «сырые» данные события Kafka гораздо надёжнее. При использовании ссылок схемы также немного различаются: `deal.stageId` в исходном событии и [deal.stage.id](http://deal.stage.id) в обогащённом.
- Мы добавили `delta` в формате JSON. Фронтенд хотел использовать её, чтобы хранилище могло обновлять только отдельные поля. Однако для этого событие Kafka должно иметь определённую структуру.
- Наряду с событиями для конкретных действий (`dealAdded`) у нас есть универсальный тип события (`dealEvent`). Он сохраняет правильный порядок событий одной сущности (`added -> changed -> deleted`), который не гарантирован хронологически при подписке на отдельные типы.

## **Живые запросы (live queries)**

Это отличная концепция, которая элегантно решает проблему, особенно беспокоившую меня, - потерю соединения.

Из [статьи Лорина](https://the-guild.dev/blog/subscriptions-and-live-queries-real-time-with-graphql) в блоге The Guild мы взяли такую же функцию **fetchOrSubscribe**. Сначала она выполняет запрос на обогащение, а затем с помощью магии асинхронного итератора подписывается на события Kafka. Отличается схема: мы используем `liveQuery` как аргумент, а не как корневое поле.

Одно это решение закрывает 99% требований к согласованности в веб-приложениях без перемотки событий, например с помощью [Redis Streams](https://redis.io/docs/data-types/streams/). Поэтому «ревалидация запроса по курсору и вычисление различий», которую [Бен Ньюман предложил на последнем GraphQL Summit](https://twitter.com/benjamn/status/1579891243798401026), кажется редким сценарием. Возможность перематывать события нужна скорее как нишевая функция продукта, например в играх.

# **Тестирование производительности и безопасности**

В конце миссии мы проверили, насколько хорошо решение масштабируется в боевой среде, постепенно включая его для реальных клиентов.

![Тестирование производительности в боевой среде](https://miro.medium.com/max/1400/1*mHE_vMsNTPr25y7IkePG2Q.png)

Тестирование производительности в боевой среде

Тесты выявили интересную особенность. Если подписываться на _сделки_ по идентификаторам в представлении _воронки_, где пользователь может быстро прокручивать список, возникают волны подписок. Поэтому лучше не блокировать `onSubscribe` проверками прав и добавить ограничение частоты запросов во фронтенде.

Мы также наблюдали, как случайные поды внезапно теряли все соединения. Из-за этого мы пересмотрели ограничения памяти и разобрались с [ошибкой max_old_space_size](https://github.com/nodejs/node/issues/35573).

Мы добавили различные ограничения безопасности на количество соединений и подписок, тайм-ауты, выходы из системы, изменения прав доступа и другие случаи. Причина не только во внешних рисках: выяснилось, что Redis расходует много CPU на сопоставление **pub**lishers и **sub**scribers.

# **Следующие шаги**

Теоретически при проблемах с производительностью можно разделить экземпляры Redis по типам сущностей.

Чтобы ускорить отказ от socketqueue, можно было бы _временно_ добавить события API в качестве универсальной подписки и проксировать все данные во фронтенд без фильтрации. Мы пожертвовали бы эффективностью ради единого транспортного слоя.

Сервис работает уже больше года. Сейчас он поддерживает около десяти сущностей, связанных с разными топиками Kafka.

## **В защиту GraphQL**

Внедрение GraphQL, особенно подписок, считают медленным, хотя оно и добровольно.

С [кодогенерацией](https://youtu.be/UZWe5Usun7I?t=4130), способной оборачивать REST API, GraphQL может казаться сложным. Без неё он выглядит слишком требовательным к переработке сервисов, рядом с REST API - избыточным, с [федерацией](https://www.apollographql.com/docs/federation/) - многослойным, а с [DataLoader](https://github.com/graphql/dataloader), но без наблюдаемости - источником [слишком большого количества плохо сгруппированных RPC-вызовов](https://twitter.com/Popeska/status/1592177670212943873).

Возврат инвестиций не виден, если не измерять _ценность_.

GraphQL взял на себя большую ответственность: сделать интернет прозрачнее, предсказуемее, эффективнее и упростить повторное использование данных. Я призываю разработчиков 🌞 проявлять профессиональное терпение, обучать коллег и показывать руководителям **объективные метрики**, доказывающие преимущества GraphQL.

Я мечтаю о дне, когда разработчики заменят события WS без схемы и REST API на GraphQL и подписки, доступные любому клиенту. Клиенты - наше будущее. Таков путь.

## Связанные материалы

- [Видео-презентация: Масштабирование GraphQL-подписок](/ru/talks/scaling-graphql-subscriptions/)
- [README.md: как документировать репозиторий](/en/blog/tech/readmemd-how-to-document-your-repo/)
- [Путь к Федеративному GraphQL](/ru/blog/tech/backend/put-k-federativnomu-graphql/)
- [Заметки со встречи об AI (26 ноября 2024 года)](/ru/blog/events/2024-11-26-ai-meetup-notes/)
