Мечта о масштабируемых и обогащённых GraphQL-подписках

Впервые опубликовано на Medium

Стилизованная фотография водопада Ягала в Эстонии

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

В прошлый раз я рассказал о пятилетнем пути GraphQL в Pipedrive. Теперь расскажу о десятилетнем пути доставки событий по WebSocket во фронтенд. Возможно, это поможет и вам.

Анимация доставки событий во фронтенд

Зачем

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

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

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

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

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

К счастью, Pipedrive «решил» эту проблему десять лет назад. В 2012 году Андрис, Капп и Таюр разработали сервис socketqueue, который и сегодня доставляет события API во фронтенд с помощью библиотеки SockJS.

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

Как видно из доклада, в основном он посвящён RabbitMQ - брокеру сообщений между PHP-монолитом и socketqueue.

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

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

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

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

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

Идентификатор ячейки в URL WebSocket

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

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

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

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

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

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

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

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

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

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

Как

Идея GraphQL-подписок довольно проста: с помощью аргументов фильтрации вы объявляете только те данные, которые хотите получать. Сервер должен разумно отфильтровать события.

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

Цели миссии

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

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

Риски

Я не знал:

  • как выполнять аутентификацию;
  • какой протокол выбрать: SSE или WS? Что такое Mercure, и есть ли где-то откат к длительному опросу;
  • будут ли соединения WS/TCP привязаны к одному поду;
  • можно ли одновременно держать несколько подписок.

Вопросы о протоколах и подключениях

Варианты архитектуры подписок

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

Прототип

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

Прототип сервиса подписок

Поток событий через прототип

Запуск

Мы вчетвером (я, Павел, Кристьян и Хиро) отправились на всё лето исследовать неизвестность 🛸 и привезти полезный результат обратно на стартовую площадку.

Команда миссии

GraphQL Go от graph-gophers

Сначала мы взяли за основу graph-gophers/graphql-go и за несколько дней воссоздали прототип подписки на события Kafka:

  • Столкнулись с тем, что целочисленный тип из спецификации GraphQL не соответствует int32 в Go.
  • Реализовали фильтрацию по свойствам и аргументам.
  • Выяснили, что для WSS нужен SSL, иначе CORS блокирует запросы.
  • Настроили проксирование запросов WebSocket через Nginx.
  • Обнаружили два транспортных протокола. WebSocket задаёт транспорт довольно свободно, поэтому каждая библиотека может реализовать его по-своему. Gophers реализовали старую версию, совместимую с Apollo и GraphiQL, но не с graphql-ws.
  • Столкнулись с тем, что библиотека Apollo Federation не принимала Subscription при регистрации схемы. Мы поняли, что можно удалить из корня только тип Subscription, продолжая проверять остальные типы на конфликты.

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

Каналы Kafka и WebSocket

Фильтрация событий в Go

Мы столкнулись с серьёзной проблемой в Gophers, связанной с JWT-аутентификацией. Нам не хватало опыта в Go и времени на изменение крупного фреймворка.

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

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

Gqlgen

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

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

Архитектура на gqlgen

Обмен каналами в gqlgen

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

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

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

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

Итоговая архитектура GraphQL-подписок

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

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

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

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

Обогащённое событие подписки

Типы схемы подписок

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

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

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

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

Из статьи Лорина в блоге The Guild мы взяли такую же функцию fetchOrSubscribe. Сначала она выполняет запрос на обогащение, а затем с помощью магии асинхронного итератора подписывается на события Kafka. Отличается схема: мы используем liveQuery как аргумент, а не как корневое поле.

Одно это решение закрывает 99% требований к согласованности в веб-приложениях без перемотки событий, например с помощью Redis Streams. Поэтому «ревалидация запроса по курсору и вычисление различий», которую Бен Ньюман предложил на последнем GraphQL Summit, кажется редким сценарием. Возможность перематывать события нужна скорее как нишевая функция продукта, например в играх.

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

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

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

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

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

Мы также наблюдали, как случайные поды внезапно теряли все соединения. Из-за этого мы пересмотрели ограничения памяти и разобрались с ошибкой max_old_space_size.

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

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

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

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

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

В защиту GraphQL

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

С кодогенерацией, способной оборачивать REST API, GraphQL может казаться сложным. Без неё он выглядит слишком требовательным к переработке сервисов, рядом с REST API - избыточным, с федерацией - многослойным, а с DataLoader, но без наблюдаемости - источником слишком большого количества плохо сгруппированных RPC-вызовов.

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

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

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

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