Батчинг сообщений

Суть батчинга в контексте STOMP

Батчинг сообщений — это стратегия группировки нескольких сообщений в одну логическую или физическую отправку для уменьшения накладных расходов на сеть и повышения пропускной способности. В экосистеме STOMP каждое сообщение обычно представлено отдельным фреймом протокола, что означает наличие заголовков, служебных байтов и отдельного цикла обработки на брокере.

STOMP не накладывает встроенного механизма автоматического батчинга на уровне клиента, поэтому агрегация сообщений реализуется на стороне приложения. В STOMP.js это выражается в ручном управлении очередями сообщений и их последующей отправкой через единый канал публикации.


Причины использования батчинга

Передача каждого сообщения отдельно приводит к ряду накладных расходов:

  • сетевые задержки на каждую операцию SEND
  • увеличение количества TCP-пакетов
  • рост нагрузки на брокер сообщений
  • частые переключения контекста обработки сообщений
  • снижение throughput при высокочастотных событиях

Батчинг решает эти проблемы за счёт:

  • уменьшения количества STOMP frame операций
  • оптимизации сетевого взаимодействия
  • снижения нагрузки на брокер (RabbitMQ, ActiveMQ, Artemis)
  • повышения устойчивости при burst-трафике

Ограничения STOMP, влияющие на батчинг

STOMP (Simple Text Oriented Messaging Protocol) работает поверх TCP и использует текстовые или бинарные фреймы. Каждый вызов отправки в STOMP.js обычно формирует отдельный frame:

SEND
destination:/queue/events
content-length:123

{...payload...}

Ключевые ограничения:

  • отсутствие стандартизированного multi-message frame
  • необходимость ручной сериализации батча
  • лимиты брокера на размер frame
  • чувствительность к latency при больших payload

Модели реализации батчинга в STOMP.js

Существует несколько устойчивых подходов к группировке сообщений.


Тайм-аутный батчинг (time-based batching)

Сообщения накапливаются в буфере и отправляются через фиксированный интервал времени.

import { Client } fr om "@stomp/stompjs";

const client = new Client({
  brokerURL: "ws://localhost:15674/ws",
});

const buffer = [];
let timer = null;

function flush() {
  if (buffer.length === 0) return;

  const batch = buffer.splice(0, buffer.length);

  client.publish({
    destination: "/queue/events.batch",
    body: JSON.stringify(batch),
  });
}

function sendMessage(msg) {
  buffer.push(msg);

  if (!timer) {
    timer = setTimeout(() => {
      flush();
      timer = null;
    }, 50);
  }
}

client.activate();

Особенности:

  • предсказуемая задержка
  • эффективен при стабильном потоке событий
  • риск увеличения latency при редком трафике

Батчинг по размеру (size-based batching)

Отправка происходит при достижении лимита сообщений.

const buffer = [];
const MAX_BATCH_SIZE = 100;

function flush() {
  const batch = buffer.splice(0, buffer.length);

  client.publish({
    destination: "/queue/events.batch",
    body: JSON.stringify(batch),
  });
}

function sendMessage(msg) {
  buffer.push(msg);

  if (buffer.length >= MAX_BATCH_SIZE) {
    flush();
  }
}

Особенности:

  • минимальная задержка при высокой нагрузке
  • возможны редкие задержки при низком трафике
  • стабильный размер payload

Гибридный батчинг (time + size)

На практике наиболее распространённая модель объединяет оба подхода.

const buffer = [];
const MAX_SIZE = 50;
const FLUSH_INTERVAL = 100;

let timer = null;

function flush() {
  if (buffer.length === 0) return;

  const batch = buffer.splice(0, buffer.length);

  client.publish({
    destination: "/queue/events.batch",
    body: JSON.stringify(batch),
  });
}

function scheduleFlush() {
  if (timer) return;

  timer = setTimeout(() => {
    flush();
    timer = null;
  }, FLUSH_INTERVAL);
}

function sendMessage(msg) {
  buffer.push(msg);

  if (buffer.length >= MAX_SIZE) {
    flush();
    if (timer) {
      clearTimeout(timer);
      timer = null;
    }
    return;
  }

  scheduleFlush();
}

Батчинг через STOMP транзакции

STOMP поддерживает транзакции, которые позволяют объединять несколько SEND операций в атомарную группу.

const tx = client.begin();

client.publish({
  destination: "/queue/events",
  body: JSON.stringify({ id: 1 }),
  transaction: tx.id,
});

client.publish({
  destination: "/queue/events",
  body: JSON.stringify({ id: 2 }),
  transaction: tx.id,
});

tx.commit();

Особенности транзакционного батчинга:

  • атомарность доставки
  • отсутствие объединения payload на уровне протокола
  • сохранение отдельных сообщений в брокере
  • полезно для consistency, но не для уменьшения frame overhead

Сериализация батчей

Так как STOMP не поддерживает нативные массивы сообщений, батч обычно кодируется в одном сообщении:

JSON-массив
body: JSON.stringify([
  { type: "click", ts: 1710000000 },
  { type: "scroll", ts: 1710000001 }
])

Плюсы:

  • простота обработки
  • совместимость с большинством backend-языков

Минусы:

  • рост размера payload
  • необходимость полной десериализации

Бинарный батч (ArrayBuffer / Uint8Array)

При высокой нагрузке используется бинарная сериализация:

const encoder = new TextEncoder();

function encodeBatch(batch) {
  const json = JSON.stringify(batch);
  return encoder.encode(json);
}

client.publish({
  destination: "/queue/events.batch",
  binaryBody: encodeBatch(buffer),
});

Особенности:

  • меньший overhead по сравнению с текстом
  • сложнее отладка
  • требует согласованной десериализации на сервере

Контроль размера батча

Размер батча критически влияет на стабильность системы.

Основные риски:

  • превышение broker frame lim it
  • увеличение GC pressure на клиенте
  • задержки доставки сообщений
  • возможные timeouts WebSocket соединения

Практический подход — ограничение:

  • по количеству сообщений (например, 50–500)
  • по суммарному размеру payload (например, 64KB–512KB)
function getSize(batch) {
  return new Blob([JSON.stringify(batch)]).size;
}

const MAX_BYTES = 128 * 1024;

function shouldFlush() {
  return getSize(buffer) > MAX_BYTES;
}

Обратное давление (backpressure)

При высокой скорости генерации событий буфер может расти неконтролируемо. В STOMP.js отсутствует встроенный механизм backpressure, поэтому он реализуется на уровне приложения.

Подходы:

  • ограничение размера буфера
  • drop policy (discard old / new messages)
  • приоритетные очереди
const MAX_QUEUE = 1000;

function sendMessage(msg) {
  if (buffer.length >= MAX_QUEUE) {
    buffer.shift();
  }
  buffer.push(msg);
}

Влияние батчинга на порядок сообщений

Батчинг изменяет семантику доставки:

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

Для строгого порядка используется:

  • последовательная очередь
  • single-thread flush loop
  • серверная сортировка по timestamp или sequenceId

Ошибки доставки и повторная отправка батчей

При неудачной отправке одного батча весь payload считается недоставленным. STOMP не гарантирует автоматического retry на уровне клиента.

Стратегии:

  • повторная отправка всего батча
  • разбиение на меньшие части
  • идемпотентные сообщения с unique id
function flushWithRetry(batch, attempt = 0) {
  try {
    client.publish({
      destination: "/queue/events.batch",
      body: JSON.stringify(batch),
    });
  } catch (e) {
    if (attempt < 3) {
      setTimeout(() => flushWithRetry(batch, attempt + 1), 100);
    }
  }
}

Серверная обработка батчей

На стороне брокера и consumer логика должна учитывать, что:

  • один STOMP message содержит массив событий
  • требуется unpacking на уровне consumer
  • возможна частичная обработка при ошибках

Пример обработки (Node.js consumer):

const batch = JSON.parse(message.body);

for (const event of batch) {
  processEvent(event);
}

Производительность и компромиссы

Эффект батчинга зависит от профиля нагрузки:

Высокочастотные события:

  • существенное снижение network overhead
  • рост latency внутри окна batching

Низкочастотные события:

  • минимальный эффект
  • возможное увеличение задержки доставки

Смешанные потоки:

  • гибридные стратегии дают наилучший результат
  • динамический размер батча повышает стабильность

Типовые ошибки реализации

  • отсутствие ограничения размера буфера
  • слишком большие интервалы flush
  • игнорирование брокерных лимитов frame size
  • отсутствие обработки ошибок publish
  • сериализация без учета стоимости JSON.stringify на больших массивах