Стриминг и обновление данных в Vega

В основе работы Vega лежит реактивный dataflow-граф, в котором визуализация рассматривается как функция от состояния данных. Любое изменение данных приводит к пересчёту зависимых компонентов: шкал, трансформаций, разметки и отрисовки. Эта модель изначально ориентирована на поддержку динамических и потоковых источников информации.

Данные в Vega описываются как набор именованных таблиц (datasets), каждая из которых может обновляться независимо. Важное свойство архитектуры заключается в том, что визуализация не «перерисовывается целиком» при изменении данных, а инкрементально обновляет затронутые участки графа исполнения.

В Vega-Lite потоковая модель частично абстрагирована: спецификация компилируется в Vega, где уже реализуется полный механизм реактивного обновления.


Обновление данных через API View

Центральный объект исполнения — View. Он предоставляет доступ к управлению данными и перезапуску визуализации.

const view = new vega.View(runtime)
  .renderer('canvas')
  .initialize('#vis')
  .run();

Обновление данных выполняется через именованные наборы:

view.change('table')
  .insert([
    {time: 1, value: 42},
    {time: 2, value: 36}
  ])
  .run();

Каждая операция формирует changeset, который применяется к внутреннему состоянию dataset.


Changeset как механизм инкрементальных изменений

Changeset представляет собой пакет операций над данными:

  • вставка новых элементов
  • удаление элементов
  • обновление существующих записей
  • массовая замена
view.change('table')
  .remove(d => d.time < 10)
  .insert([{time: 11, value: 99}])
  .run();

Удаление часто выполняется через предикаты, что позволяет реализовать потоковую очистку старых данных в сценариях временных рядов.


Потоковые сценарии: таймеры и polling

Типичный поток данных в Vega реализуется через внешние источники:

setInterval(() => {
  view.change('stream')
    .insert([{t: Date.now(), v: Math.random()}])
    .run();
}, 1000);

Такой подход используется для:

  • телеметрии
  • мониторинга систем
  • финансовых котировок
  • сенсорных данных IoT

Внутри Vega обновление не блокирует рендер, а ставит изменения в очередь исполнения dataflow.


Интеграция WebSocket и внешних потоков

Реальные потоковые системы чаще используют WebSocket или SSE:

const socket = new WebSocket('wss://example.com/data');

socket.onmess age = (event) => {
  const point = JSON.parse(event.data);

  view.change('stream')
    .insert([point])
    .run();
};

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


Ограничение размера окна (sliding window)

Потоковые данные требуют контроля объёма хранимых точек. Распространённый паттерн — скользящее окно:

view.change('stream')
  .remove(d => d.t < Date.now() - 60000)
  .insert([{t: Date.now(), v: value}])
  .run();

Это поддерживает постоянный объём данных и предотвращает деградацию производительности.


Реактивные сигналы и динамическое поведение

Помимо данных, Vega поддерживает сигналы (signals), которые также участвуют в dataflow-графе. Изменение сигнала автоматически вызывает пересчёт зависимых выражений.

view.signal('threshold', 50).run();

Сигналы часто используются для:

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

Обновление через setState и изменение масштаба данных

Некоторые сценарии требуют пересборки части визуализации при значительном изменении структуры данных:

view.data('table', newData).run();

Такой подход менее инкрементальный, но применяется при:

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

Streaming transforms внутри dataflow

Встроенные трансформации позволяют обрабатывать поток без внешней логики:

  • window — агрегирование по окнам
  • filter — фильтрация событий
  • aggregate — вычисление метрик
  • stack — накопительные значения

Пример оконной агрегации:

{
  "transform": [
    {
      "type": "window",
      "ops": ["mean"],
      "fields": ["value"],
      "frame": [-10, 0]
    }
  ]
}

Такие трансформации выполняются инкрементально при поступлении новых данных.


Vega-Lite и потоковые обновления

В Vega-Lite потоковая модель реализуется через пересоздание spec или частичное обновление данных в vega-embed.

const view = await vegaEmbed('#vis', spec);

setInterval(() => {
  spec.data[0].values.push({x: Date.now(), y: Math.random()});

  view.view
    .change('source')
    .insert([{x: Date.now(), y: Math.random()}])
    .run();
}, 1000);

В отличие от низкоуровневого Vega, Vega-Lite чаще опирается на пересборку спецификации, но при использовании view доступно прямое управление потоками.


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

Инкрементальная модель Vega оптимизирована через:

  • батчинг изменений (несколько изменений за один run)
  • минимизацию перерасчёта шкал
  • ленивую пересборку сцен-графа
  • дифференциальное обновление DOM/canvas

Критический аспект — частота run() вызовов. При высокочастотных потоках применяется агрегация изменений:

let buffer = [];

setInterval(() => {
  if (buffer.length) {
    view.change('stream')
      .insert(buffer)
      .run();

    buffer = [];
  }
}, 200);

Синхронизация нескольких потоков

При работе с несколькими источниками данных применяется единая точка синхронизации изменений:

Promise.all([streamA, streamB]).then(([a, b]) => {
  view.change('a').insert(a);
  view.change('b').insert(b);
  view.run();
});

Такой подход предотвращает частичные обновления визуализации.


Управление состоянием визуализации

Сложные потоковые системы требуют контроля состояния:

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

Vega не навязывает слой состояния, но предоставляет инструменты для его реализации на уровне приложения через dataset-операции и сигналы.


Использование выражений для динамики потоков

В expressions можно реализовать вычисления прямо в спецификации:

{
  "signal": "now()"
}

или

{
  "expr": "datum.value > threshold"
}

Это позволяет переносить часть логики обработки потока внутрь декларативного слоя, снижая нагрузку на внешний код.


Архитектурные паттерны потоковой визуализации

В практических системах Vega используется в нескольких типовых конфигурациях:

  • push-based поток (WebSocket → View)
  • pull-based polling (setInterval → API)
  • hybrid (stream + batch sync)
  • event-sourced визуализация (лог событий как источник truth)

Каждый паттерн влияет на структуру dataset-операций и частоту обновлений графа исполнения.


Поведение при высокочастотных данных

При частотах обновления выше 60–100 событий/сек критическим становится:

  • уменьшение числа run() вызовов
  • агрегация на стороне клиента
  • использование window-трансформаций
  • фильтрация шумовых данных до вставки

Vega сохраняет стабильность за счёт того, что пересчёт графа оптимизируется на уровне зависимостей, а не всей сцены целиком.


Модель консистентности данных

В потоковой среде Vega использует модель eventual consistency внутри одного кадра рендеринга: все изменения, внесённые до run(), применяются атомарно относительно визуального состояния. Это исключает промежуточные неконсистентные состояния сцены.


Реактивная связка данных и визуальных примитивов

Каждое обновление dataset автоматически влияет на:

  • масштабы (scale domains)
  • оси (axes)
  • геометрию (marks)
  • интерактивные элементы (tooltips, selection)

Эта связка обеспечивает полную реактивность визуализации без ручного пересчёта зависимостей.