Обработчик, который под всплесками теряет и переставляет обновления стакана, — обычно два разных отказа за одним симптомом: работа с последовательностью на приёме и слой раздачи, где одна медленная сессия меняет то, что получают все остальные. Разбираем, как устроены последовательность, восстановление, fan-out и backpressure и где конфляция перестаёт быть честной.
Обработчик, который под всплесками теряет и переставляет обновления L2, — обычно два разных отказа за одним симптомом. Первый живёт на приёме, где отслеживается последовательность фида и пропуск нужно обнаружить, а не проглотить. Второй — на раздаче, где один нормализованный стакан расходится по множеству сессий и единственный медленный читатель меняет то, что получают остальные.
Чинятся они по-разному, и не то лечение сдвигает симптом вместо того, чтобы его убрать. Дальше — устройство конвейера, когда последовательность обязана пережить всплеск: что такое пропуск на самом деле, как восстановление стыкует снимок с живым потоком, как fan-out решает вопрос порядка один раз и где конфляция честна.
Короткий ответ структурный. Порядок решается ровно в одном месте, выше по потоку от всех сессий: один писатель на инструмент сворачивает нумерованный фид в стакан, а сессии получают представления, производные от этой свёртки, — сами они ничего не переупорядочивают. Что amBrain может подтвердить публично: построенная нами мини-биржа работает в проде на колокации MOEX, мы построили торговый терминал Spectre Trade, а публикуемая нами задержка рыночных данных замерена — менее 5 мс на путях, которые строим мы. Эта цифра описывает наши пути, а не бенчмарк схемы, описанной ниже.
Номер последовательности обещает порядок, а не доставку
Фиды нумеруют свои обновления, и этот номер — единственный источник порядка, который у вас есть. Время прихода таким источником не является: multicast-пути переставляют пакеты, один инструмент идёт по нескольким каналам, приёмные очереди разложены по ядрам, а всплеск растягивает всё это. Обработчик, упорядочивающий по приходу, корректен, только пока сеть спокойна, — в том самом условии, за которое никто и не переживал.
Шесть свойств фида нужно знать до того, как написана логика восстановления. Каждое меняет то, что значит пропуск.
- Единица, которую покрывает последовательность: канал, инструмент или стакан. Номер по каналу не скажет, какой инструмент потерял обновление, а номер по инструменту не скажет, что канал встал
- Правило приращения: строго подряд внутри единицы или возрастание с допустимыми дырами. Встречается и то и другое, и чтение второго как первого рождает восстановления, которые были не нужны
- Начинается ли нумерация заново на границе сессии и чем эта граница помечена: перезапуск, прочитанный как пропуск, отправляет на восстановление все инструменты одновременно
- Несут ли heartbeat'ы текущий номер последовательности. Без этого мёртвое соединение и тихий инструмент выглядят одинаково
- Есть ли ретрансмиссия и на каком окне. Если её нет, восстановление по снимку — единственный путь назад, и оно обязано быть достаточно дешёвым, чтобы пользоваться им часто
- С каким номером последовательности согласован снимок. Без него снимок вообще не состыковать с живым потоком
Там, где свойство действительно неизвестно, замерьте его, а не зашивайте догадку. Каждое из шести становится ветвью в пути восстановления, и неверное допущение там находится позже — как стакан, который тихо расходится с площадкой.
Пропуск и перестановка первые миллисекунды выглядят одинаково
Оба начинаются одинаково: следующее обновление несёт не тот номер, которого вы ждали. Разница — во времени, поэтому классификация делается не в момент прихода, а когда истекает ограниченное ожидание.
- Не по порядку: ждали N, пришёл N+2, а N+1 приходит, пока ожидание ещё открыто. Ничего не потеряно, и вся цена — это ожидание
- Дубль или ретрансмиссия: номер не выше последнего применённого. Отбрасывается, не трогая стакан, и считается, потому что рост доли дублей что-то говорит о пути
- Пропуск: ожидание истекло, а N+1 так и не пришёл. Стакан не может двигаться дальше дыры, и этот инструмент уходит на восстановление
- Устаревшее: номер правильный, но пришёл слишком поздно, чтобы пригодиться. Байты дошли, а ниже по потоку это потеря
Одно правило не даёт порче стать тихой: обновление применяется, только если его номер ровно тот, которого ждут. Всё остальное уходит в буфер ожидания или на восстановление. Стакан, принявший дельту не по порядку, продолжает отдавать цены и выглядит здоровым — расхождение с площадкой находят позже, находит его клиент, на исполнении, которое не сошлось.
Ожидание — ограниченная структура, а не растущая очередь. Она держит обновления, обогнавшие ожидаемый номер, с ключом по номеру, поэтому выпуск из неё — поиск, а не сортировка.
- Выпуск из буфера — цикл: примените ожидаемый номер, затем применяйте то, что уже лежит в буфере, пока номера идут подряд
- Срок ожидания задаётся во времени, а не только числом отложенных обновлений: всплеск заполняет окно, считающее штуки, куда раньше, чем задумывалось
- Сколько ждёт буфер, столько ждёт каждый потребитель. Задавайте срок ожидания по перестановкам, замеренным на вашем собственном пути, а не по числу, которое показалось безопасным
- Переполнение буфера само по себе объявляет пропуск: ожидание ограничено и по памяти, и по времени
- Ожидание ведётся по инструменту или по каналу, но никогда глобально. Один тихий инструмент не должен тормозить всё вокруг
Восстановление — это снимок, состыкованный с потоком, который вы уже буферизовали
Ломается именно стыковка. Снимок — это стакан на некоторый номер последовательности, и он устаревает в момент своего создания; пригодным его делает инкрементальный поток, буферизованный, пока снимок ехал.
- Буферизуйте инкрементальный поток до того, как запрошен снимок. Снимок без живого потока за спиной отстаёт от рынка уже в момент прихода
- Читайте номер последовательности, которому соответствует снимок. Если фид его не публикует, на практике фид работает только снимками, и проект обязан сказать это вслух
- Отбросьте буферизованные обновления с номером снимка и ниже, остальные примените по порядку. Если первое из них — не обновление сразу за снимком, стыковка не удалась и восстановление начинается заново
- Если буфер заполнился до прихода снимка, начните восстановление заново, а не применяйте его частично: частично применённое восстановление неотличимо от здорового стакана
- На время восстановления публикуйте инструмент как деградировавший — явным состоянием в потоке. Стакан с дырой, отданный как актуальный, хуже, чем отсутствие стакана
- Проверьте после стыковки: контрольную сумму, которую публикует фид, если он её публикует, или совпадение вашего свёрнутого стакана со следующим снимком
Восстановление — штатное событие, а не инцидент, и его цена принадлежит плану мощностей: сколько занимает получение снимка, сколько потока буферизуется тем временем и сколько инструментов могут восстанавливаться одновременно, прежде чем сервис снимков станет узким местом.
Fan-out: нормализуем раз, кодируем раз, отправляем многим
Сотни сессий терминала хотят один и тот же стакан. Ошибка, которая множится под всплеском, — делать на каждую сессию работу, которая по своей природе не сессионная: пересобирать стакан для каждого подписчика или сериализовать одно и то же обновление по разу на сокет.
- Стаканом владеет один писатель на шард инструментов. Читатели его никогда не меняют, что снимает и блокировку, и вопрос о том, чья версия главная
- Писатель публикует версионированные обновления в кольцевой буфер, за которым читатели следуют в своём темпе, поэтому отставший читатель никого не тормозит
- Каждое обновление кодируется один раз на каждый формат передачи и раздаётся сессиям по ссылке. Своими у сессии остаются только фрейминг и управление потоком
- У каждой сессии свой исходящий номер последовательности, поэтому клиент обнаруживает собственные потери, ничего не зная о фиде выше по потоку
- Порядок гарантируется в пределах инструмента, потому что именно на эту гарантию опираются клиенты. Порядок между инструментами либо обещан явно и реализован, либо не обещан вовсе
- За пределами одного процесса fan-out превращается в ярус ретрансляторов: каждый берёт одну подписку выше по потоку и обслуживает свою долю сессий, поэтому работа писателя остаётся постоянной
Стоимость fan-out определяется тем, сколько раз обновление преобразуется, а не тем, сколько сокетов его получают. Закодировать один раз и передать ссылку — масштабируется по числу сессий; пересобирать стакан на каждую сессию — нет.
Медленный потребитель — выбранная вами политика, а не случайность
Где-то есть сессия на плохой сети или терминал с застрявшим циклом отрисовки, и его исходящий буфер заполняется. Возможных поведений четыре, и два из них выбирают только по случайности.
- Блокировать писателя, пока медленная сессия не разгрузится, — никогда. Так одно плохое соединение превращается в скачок задержки для всех на шарде
- Растить очередь без предела — медленный потребитель превращается в исчерпание памяти, а затем в отказ, никак не связанный с исходной сессией
- Ограниченная очередь с конфляцией — верно для состояния стакана, где клиенту нужна текущая картина, а не каждый промежуточный шаг
- Ограниченная очередь с разрывом на верхней отметке — верно для потоков, к которым конфляция неприменима, где выброшенный элемент выбрасывает смысл
- Какой бы ни была политика, очередь ведётся по сессии, а отставание меряется постоянно: глубина очереди и разрыв между опубликованным номером и номером, записанным в сокет
- Разрыв соединения называет свою причину. Необъяснённое закрытие переоткрывают в цикле; после объяснённого клиент переподписывается
Backpressure — место, где сходятся обе стороны. Если исходящий путь может давить назад на писателя стакана, медленный терминал в итоге тормозит свёртку фида, и обнаружение пропусков начинает срабатывать по причинам, никак не связанным с площадкой. Ограниченное кольцо между ними разрывает эту цепочку.
Конфляция честна для состояния и неверна для событий
Стакан — это состояние: клиенту нужны текущие уровни, а уже заменённое значение само по себе смысла не несёт. Лента сделок — журнал событий, где каждый элемент есть случившийся факт, и свернуть его в сводку нельзя.
- Конфляция уместна: обновления ценовых уровней, вершина стакана, агрегированная глубина и производная статистика вроде последней цены или объёма сессии
- Конфляция недопустима: сделки и принты, отчёты по заявкам и исполнениям, аукционы и смены фаз, а также всё, что клиент агрегирует во времени, — лента, собранная из потока с конфляцией, это неверное число, в которое верят
- Конфляция по ключу, а не по потоку. Хранить последнее обновление на каждый ценовой уровень — значит сохранить стакан; хранить последнее обновление в целом — выбросить все уровни, которые изменились не последними
- Обновление после конфляции несёт номер последовательности того состояния, которое представляет, поэтому клиент знает, какому моменту оно соответствует
- Интервал конфляции — часть той задержки, о которой вы отчитываетесь. Поток с конфляцией на интервале не описывается задержкой, замеренной на потоке без неё
- Клиент, которому нужны все промежуточные состояния — бэктест, комплаенс-запись, — берёт поток без конфляции и платит трафиком
Конфляция — смена формы, а не настройка сжатия. Как только к потоку применена конфляция, клиент не восстановит, что произошло между двумя обновлениями, и говорить ему, что поток полон, нельзя. Публиковать оба — поток стакана с конфляцией и поток событий без неё — и есть то, что держит правыми обе категории клиентов.
Переподключение — это ресинхронизация, и приходят они все разом
Когда сессия возвращается, удерживаемый ею стакан ничего не стоит, пока сервер не докажет непрерывность. По умолчанию — свежий снимок на каждую подписку, со своим номером последовательности, применяемый к клиенту, который сначала сбросил локальное состояние.
- Возобновление с номера последовательности предлагается только там, где есть ограниченный буфер реплея. Когда запрошенный номер уже вытеснен, сервер говорит об этом и откатывается к снимку, а не отдаёт поток с дырой
- Состояние сессии при переподключении — явное решение: либо сервер держит подписки ограниченное время по токену сессии, либо клиент заявляет их заново при подключении. Работают оба; неявная смесь — нет
- Дубли после возобновления — ожидаемое поведение, клиент отбрасывает их по номеру. At-least-once плюс нумерация последовательности реализуются корректно проще, чем exactly-once
- Переподключения приходят вместе: то, что оборвало одну сессию, обычно оборвало многие. Backoff с джиттером на клиенте и контроль допуска на сервере не дают восстановлению стать вторым отказом
- Снимки для этой толпы берутся из кеша по инструментам, обновляемого с заданной периодичностью, поэтому писатель сериализует снимок по расписанию, а не по разу на каждую переподключающуюся сессию
- Стакан на стороне клиента пересобирается, а не патчится. Терминал, сохранивший старые уровни и накладывающий новые дельты сверху, переносит ошибку, бывшую до обрыва, в стакан, который теперь выглядит свежим
Проектировать стоит не под одно переподключение, а под сетевое событие, которое возвращает сотни сессий в одну и ту же секунду, и каждая просит снимок по каждому инструменту, за которым следила, — пока приёмная сторона восстанавливается после пропуска, вызванного тем же событием.
Что amBrain может подтвердить публично: мы строим low latency торговые платформы, матчинг-движки и системы real-time bidding на Rust из Еревана, Армения, а публикуемая нами задержка рыночных данных — менее 5 мс — замерена на путях, которые строим мы. Если ваш обработчик теряет последовательность под всплесками, разговор, который стоит провести, — тот, что разделяет приём и раздачу до того, как переписывать любую из сторон.