📱 Подписаться
IT и цифровая трансформация

От «быстрого JSON» к потоковой обработке данных: смена парадигмы в оценке производительности протоколов

📰 Habr 👁️ 0 просмотров

ihar764 часа назад

От «быстрого JSON» к потоковой обработке данных: смена парадигмы в оценке производительности протоколов

Уровень сложностиСложныйВремя на прочтение7 минОхват и читатели4.4KGo*Высоконагруженные системы*Алгоритмы*МнениеСразу хочу обозначить важный момент: эта статья не столько про пакет SilentJSON, сколько про схему работы с информацией, которую на его примере удалось реализовать и проверить на практике.

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

Сам SilentJSON – это результат многолетней практики работы с передачей и обработкой данных. Многие решения в нём появились не как попытка заранее придумать идеальную архитектуру, а как ответ на реальные ограничения: память, latency, количество аллокаций, размер потоков, стоимость промежуточных преобразований и необходимость обрабатывать данные как можно раньше.

Поэтому ниже я не буду доказывать, что именно SilentJSON является каким-то универсальным решением. Скорее хочу показать подход, который можно взять отдельно от библиотеки и использовать как основу для собственной реализации.

И да – весь этот подход и сам пакет я публикую бесплатно. Мне интереснее поделиться самой идеей и накопленным практическим опытом, чем просто показать очередной benchmark JSON parser’а.

Когда говорят о производительности JSON, обычно измеряют одно и то же: есть готовый буфер, его нужно разобрать, получить структуру в памяти и передать приложению.

получить данные → parse → object → обработкаТак строится большинство benchmark’ов сериализации. Но для системы, работающей с непрерывным потоком, есть более интересный вопрос: насколько рано после получения данных мы можем начать что-то с ними делать?

Если полезное событие находится в начале потока, ожидание полного сообщения становится отдельной задержкой. Поэтому для потоковой обработки важны уже не только MB/s parser’а и количество аллокаций, но и время до первого полезного события, объём памяти и возможность одновременно получать, распознавать и обрабатывать данные.

Свойство JSON, которое привыкли считать недостатком

У JSON структура непосредственно находится в самом потоке:

{ } [ ] , : "Эти символы позволяют находить границы структурных элементов без предварительного чтения отдельного заголовка или метаданных.

Это особенно интересно с точки зрения CPU. После обработки строк, escape-последовательностей и других специальных случаев parser может быстро сканировать вход и находить структурные позиции. SIMD позволяет делать такой поиск большими блоками данных.

Поэтому вопрос производительности протокола стоит рассматривать не только как “сколько байт занимает сообщение”, а как сколько работы требуется CPU, чтобы получить из этих байт информацию, управляющую дальнейшей обработкой.

Например, в protobuf часть этой информации находится в tag и wire type. Parser должен интерпретировать их, обработать varint или длину поля и определить, как перейти к следующему элементу. Это вполне эффективная модель, но она отличается от ситуации, когда структурные признаки непосредственно присутствуют в потоке.

SilentJSON: от decode к потоку

SilentJSON создавался как быстрый JSON encoder/decoder для Go, но потоковый режим позволяет использовать тот же механизм иначе.

сетевой поток
↓
StreamDecoder
↓
очередной объект
↓
обработкаНапример:

registry := silentjson.NewRegistry[Employee]()

decoder := silentjson.NewStreamDecoder(reader, registry)

for {
employee, err := decoder.Next()
if err == io.EOF {
break
}
if err != nil {
return err
}

process(employee)
}Граница Read() здесь не обязана совпадать с границей JSON-объекта. Один объект может прийти несколькими физическими chunks, а один chunk может содержать несколько объектов. Parser сохраняет состояние между чтениями.

Это позволяет отделить транспортный поток от логической структуры данных.

При этом размер логической записи не обязан быть фиксированным: одна запись может занимать сотни байт, следующая – несколько мегабайт. Ограничение рабочего буфера поэтому не означает ограничения размера самого потока или единого размера сообщения.

NextRawBlock: когда декодирование не требуется

Во многих задачах JSON вообще не нужно превращать в Go-структуру. Нужно только определить границы очередного блока и передать его дальше.

Для этого в SilentJSON есть NextRawBlock.

JSON stream
↓
structural scan
↓
raw block
↓
outputВ таком режиме parser не тратит время на mapping полей, создание объектов и последующее преобразование обратно в JSON.

На этой модели построен silent-chunker. Он может разбивать большой JSON-массив на отдельные блоки со скоростью более 4 GB/s. Например:

./silent-chunker -file input.json -count 1000 -out chunk_или:

cat large.json | ./silent-chunker -size 10485760Для файла в 10 GB это уже принципиально другой класс обработки: parser занимается структурой потока, а не построением промежуточного представления каждого объекта.

Registry отвечает уже за другое

Registry не нужен для того, чтобы понять саму структуру JSON. Он нужен, чтобы сопоставить найденное поле с конкретным местом в Go-структуре и его типом.

JSON bytes
↓
структура
↓
поле
↓
Go field / typeВ registry заранее находятся необходимые метаданные: offsets, типы, вложенные registry и другая информация для записи результата.

Таким образом, две задачи разделены: JSON сообщает parser’у, где находятся данные, а registry – что с ними делать в Go.

Это позволяет убрать значительную часть работы reflection-heavy generic path с hot path.

RawJson: данные как вложенный протокол

RawJson позволяет оставить часть JSON в исходном виде и передать её следующему уровню обработки.

Protocol A
↓
RawJson
↓
Protocol B
↓
RawJson
↓
Protocol CЭто позволяет строить вложенные протоколы без обязательного полного декодирования каждого уровня. Текущий слой обрабатывает только ту структуру, которая ему нужна, а вложенный фрагмент передаёт дальше.

Такая модель напоминает обычное разделение протоколов: один уровень знает о своём формате, а содержимое следующего уровня передаётся ему как payload.

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

От сообщений к событиям

Отсюда возникает ещё одна возможность.

Если системе требуется только обнаружить событие, не всегда имеет смысл сначала инициализировать полный pipeline обработки.

Например:

data received ↓ event detected ↓ event registered ↓ normal processingВместо:

initialize → receive → parse → process → detectможно в некоторых сценариях получить:

receive → detect → registerЭто особенно интересно для систем, где само событие важно зафиксировать как можно раньше.

В банковских операциях это может быть регистрация этапов транзакции по мере их появления. В логировании – фиксация критического события до запуска более тяжёлого logging pipeline. В мониторинге – обнаружение роста задержек, retry, проблем с соединением или деградации сервиса до того, как они начнут распространяться по цепочке микросервисов.

Здесь оптимизируется уже не только обработка данных, но и стоимость подготовки системы к обработке.

Как меняется измерение производительности

Для обычного parser’а достаточно спросить:

Сколько MB/s он декодирует?Для потоковой системы этого недостаточно.

Имеют значение:

• время до первого полезного события;
• время до первой обработанной записи;
• throughput непрерывного потока;
• рабочий объём памяти;
• количество аллокаций;
• стоимость обработки структурных событий;
• возможность перекрывать получение и обработку;
• время обнаружения проблемы.Две системы могут обработать одинаковые 10 GB за одинаковое время, но одна начнёт работу только после получения всего объёма, а другая будет обрабатывать первые записи уже во время поступления остальных.

Для throughput это может быть одинаковый результат. Для latency и поведения системы под нагрузкой – совершенно разный.

Результаты SilentJSON

В одном из benchmark’ов StreamDecoder.Next показывает около 614 MB/s при 41 MB аллокаций. NextRawBlock, где отсутствует полноценное mapping в структуры, – около 4009 MB/s при примерно 0.9 MB аллокаций.

NextChan, использующий producer-consumer модель, показывает около 447 MB/s и 41 MB аллокаций.

РежимThroughputAllocationsЧто измеряетсяStreamDecoder.Next~614 MB/s~41 MBдекодирование объектовNextRawBlock~4009 MB/s~0.9 MBструктурное извлечениеNextChan~447 MB/s~41 MBпотоковый producer/consumerВ другом тесте на 100 000 объектов SilentJSON в parallel mode показал около 24 670 MB/s против примерно 644 MB/s для Sonic и 110 MB/s для encoding/json.

Scalar fallback без AVX2 показал около 810 MB/s. Это важно хотя бы потому, что сама модель не сводится к использованию SIMD: SIMD ускоряет структурный поиск, но основная архитектура остаётся работоспособной и без него.

Эти цифры относятся к конкретным тестам, данным и реализации, поэтому их нельзя напрямую переносить на любой JSON workload. Но они показывают разницу между полноценным декодированием и операциями, которым вообще не требуется построение Go-объектов.

Zero-copy

Потоковая обработка также меняет отношение к строкам.

Если возвращаемая строка является представлением исходного JSON buffer, её можно получить без отдельного копирования. Но lifetime такой строки связан с lifetime исходных данных.

Если значение нужно сохранить независимо от входного буфера, копирование должно быть явным:

strings.Clone(s)Это позволяет не платить за копирование там, где данные нужны только во время текущей обработки, и явно платить за него там, где требуется долгоживущая копия.

От быстрого parser’а к потоковой архитектуре

В итоге SilentJSON интересен не только как попытка сделать Unmarshal быстрее.

Главное изменение – в самой единице обработки.

JSON можно рассматривать не как одно неделимое сообщение, а как поток структурных событий и логических блоков:

byte stream
↓
structural scan
↓
logical block
↓
RawJson / typed object
↓
event / processing
↓
next layerНе каждый уровень обязан декодировать всё сообщение. Если нужны границы – достаточно структурного сканирования. Если нужен объект – используется registry. Если нужен вложенный протокол – можно передать RawJson. Если требуется только событие – его можно зарегистрировать отдельно от полной обработки.

Именно поэтому классический вопрос “насколько быстрый у нас JSON parser?” становится недостаточным.

Для непрерывного потока важнее другое:

как рано мы можем получить полезную информацию из данных и сколько работы требуется до первой реакции системы?

В этом смысле производительность протокола – это не только MB/s.

Это время до первого события, объём промежуточного состояния, стоимость перехода между уровнями и способность получать, распознавать, передавать и обрабатывать данные одновременно.

И тогда JSON становится не просто форматом сериализации. Его структурные байты становятся частью механизма управления потоком.

А SilentJSON – не просто быстрым JSON parser’ом, а экспериментом с другой моделью обработки:

получать
↓
распознавать
↓
реагировать
↓
передавать
↓
обрабатыватьПричём эти этапы не обязаны ждать завершения друг друга.Теги:• json
• потоковая обработка
• парсинг
• go
• highload
• производительность
• оптимизация
• silentjson
• zero-copyХабы:• Go
• Высоконагруженные системы
• Алгоритмы

Получайте больше инсайтов о систематизации бизнеса

Подписывайтесь на Telegram-канал Business Operations — ежедневные материалы о бизнес-процессах, операционном управлении и повышении эффективности

💬 Подписаться на канал