← назад к разделу

Настоящее приложение почти никогда не держит данные в одном месте. Рядом с основной базой обычно живут кеш, поисковый индекс, аналитическое хранилище, денормализованные таблицы, read-модели. В этом зоопарке легко утонуть и перестать понимать, кто кому хозяин и какой копии верить при расхождении.

Есть простая рамка, которая сразу наводит порядок: разделить все хранилища на два вида — систему записи и производные данные.

пишем только в источник — производные выводим из него запись: order-1 → PAID система записи order-1 NEW PAID терять нельзя поток изменений #1 #2 #3 производные данные — можно выбросить кеш поисковый индекс read-модель NEW NEW NEW PAID PAID PAID пусто PAID одно изменение — в одном порядке во всех копиях индекс потеряли — источник цел пересобрали из источника — снова сходится

Приложение пишет в одно место — в систему записи. Дальше изменение расходится по производным наборам единым потоком, поэтому у кеша, индекса и read-модели один и тот же порядок обновлений. Потеря производного набора не теряет данных: его выбрасывают и выводят из источника заново. Потеря самого источника — уже потеря данных.

Обязательно

Источник правды и его отражения

Система записи (source of truth, источник правды) — это хранилище, где факт лежит ровно один раз и в авторитетном виде. Обычно это ваша нормализованная база: пользователь оформил заказ — заказ записан сюда. Важная оговорка: одного источника правды «на всю компанию» не бывает. Когда сервисов несколько, у каждого факта свой хозяин — заказ живёт в сервисе заказов, платёж в платёжном, профиль в профильном, — и источник правды ищут не «где-то в компании», а для каждого факта отдельно. И если данные в источнике вдруг разошлись с данными в любом другом хранилище, правым по определению считается источник.

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

  • кеш — быстрая копия части данных;
  • поисковый индекс (Elasticsearch, GIN) — те же данные, переупакованные под поиск по тексту;
  • материализованное представление или денормализованное поле — заранее посчитанный результат;
  • read-модель в CQRS (Java, Go, Node, Python) — данные, разложенные ровно под конкретный экран;
  • аналитическое хранилище (ClickHouse) — те же события, но под сканы.

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

Главное свойство: производное можно выбросить

Из этого определения следует вывод, который меняет отношение к данным: производные данные одноразовы. Потеряли кеш, индекс или read-модель — не катастрофа: пересобрали из источника, и система как новенькая. «Не катастрофа» не значит «бесплатно»: переиндексация полусотни миллионов документов идёт часами и в ночное окно может не влезть, так что время пересборки надо знать заранее, а не выяснять в момент аварии. А вот потеряли саму систему записи — вот это настоящая беда, восстанавливать неоткуда.

Отсюда и остальное: бэкапить и защищать в первую очередь нужно источник правды. У каждого производного набора должна быть процедура пересборки с нуля (переиндексировать, прогреть кеш, перестроить проекцию).

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

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

  • Online (сервис). Ждёт запрос и отвечает как можно быстрее; главное — время отклика и доступность. Так работает само приложение. Но для получения больших производных наборов online не подходит.
  • Batch (пакетная обработка). Берёт большой набор данных на вход, молотит его минуты или часы, выдаёт результат; главное тут — пропускная способность, а запускается это по расписанию (ночная переиндексация, пересчёт витрины). Классика больших данных — MapReduce/Spark; в обычном бэкенде проще — фоновый обработчик на очереди.
  • Stream (потоковая обработка). Нечто среднее: реагирует на события по мере их поступления, с низкой задержкой, но обрабатывает не весь набор сразу, а поток изменений. Так производные данные держат свежими почти в реальном времени.

Дисциплина: деривировать, а не дублировать записью

Главная ошибка при работе с производными данными — писать в них напрямую, в обход источника. Классический пример — двойная запись: приложение пишет и в базу, и в кеш (или в поисковый индекс) двумя отдельными операциями. Рано или поздно одна операция пройдёт, а вторая упадёт — и данные молча разъедутся.

приложение база: PAID запись прошла кеш: NEW запись упала

Приложение пишет двумя отдельными операциями: в базу заказ лёг как PAID, вторая запись упала, и в кеше остался NEW.

Правильный принцип: писать только в систему записи, а всё производное выводить из неё — единым потоком изменений. Тогда у всех производных наборов один общий порядок обновлений, они согласованы между собой и легко пересобираются. Технически это делают двумя способами. Первый — события и таблица исходящих (outbox): событие кладут в ту же транзакцию, что и саму запись, а отправляет его отдельный процесс. Второй — CDC (change data capture): база назначается ведущим узлом, а производные системы читают её журнал изменений и применяют их ровно в том же порядке, в каком они легли в базу. На том же принципе стоит event sourcing: источник правды — это поток событий, а всё остальное (проекции, read-модели) выводится из него.

Разница видна на десятке строк: источник — карта заказов, поток изменений — список записей. Сначала в кеш пишут напрямую, потом тот же кеш выбрасывают и выводят из потока заново.

живой пример

import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;

public class DerivedData {
    record Change(String orderId, String status) {}

    static final Map<String, String> source = new LinkedHashMap<>();
    static final List<Change> changes = new ArrayList<>();

    public static void main(String[] args) {
        write("order-1", "NEW");
        write("order-2", "NEW");
        write("order-1", "PAID");

        Map<String, String> cache = derive();
        cache.put("order-1", "SHIPPED");
        System.out.println("источник:    " + source);
        System.out.println("кеш:         " + cache);
        System.out.println("сходятся: " + cache.equals(source));

        Map<String, String> rebuilt = derive();
        System.out.println("пересобрали: " + rebuilt);
        System.out.println("сходятся: " + rebuilt.equals(source));
    }

    static void write(String orderId, String status) {
        source.put(orderId, status);
        changes.add(new Change(orderId, status));
    }

    static Map<String, String> derive() {
        Map<String, String> view = new LinkedHashMap<>();
        for (Change c : changes) {
            view.put(c.orderId(), c.status());
        }
        return view;
    }
}
Запустить

Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →

живой пример

package main

import (
	"fmt"
	"strings"
)

type change struct{ orderID, status string }

type ordered struct {
	keys []string
	vals map[string]string
}

func newOrdered() *ordered { return &ordered{vals: map[string]string{}} }

func (o *ordered) put(key, value string) {
	if _, ok := o.vals[key]; !ok {
		o.keys = append(o.keys, key)
	}
	o.vals[key] = value
}

func (o *ordered) equals(other *ordered) bool {
	if len(o.vals) != len(other.vals) {
		return false
	}
	for key, value := range o.vals {
		if other.vals[key] != value {
			return false
		}
	}
	return true
}

func (o *ordered) String() string {
	parts := make([]string, len(o.keys))
	for i, key := range o.keys {
		parts[i] = key + "=" + o.vals[key]
	}
	return "{" + strings.Join(parts, ", ") + "}"
}

var source = newOrdered()
var changes []change

func write(orderID, status string) {
	source.put(orderID, status)
	changes = append(changes, change{orderID, status})
}

func derive() *ordered {
	view := newOrdered()
	for _, c := range changes {
		view.put(c.orderID, c.status)
	}
	return view
}

func main() {
	write("order-1", "NEW")
	write("order-2", "NEW")
	write("order-1", "PAID")

	cache := derive()
	cache.put("order-1", "SHIPPED")
	fmt.Println("источник:   ", source)
	fmt.Println("кеш:        ", cache)
	fmt.Println("сходятся:", cache.equals(source))

	rebuilt := derive()
	fmt.Println("пересобрали:", rebuilt)
	fmt.Println("сходятся:", rebuilt.equals(source))
}
Запустить

Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →

живой пример

const source = new Map();
const changes = [];

const show = (m) => '{' + [...m].map(([k, v]) => `${k}=${v}`).join(', ') + '}';
const same = (a, b) => a.size === b.size && [...a].every(([k, v]) => b.get(k) === v);

function write(orderId, status) {
  source.set(orderId, status);
  changes.push({ orderId, status });
}

function derive() {
  const view = new Map();
  for (const c of changes) view.set(c.orderId, c.status);
  return view;
}

write('order-1', 'NEW');
write('order-2', 'NEW');
write('order-1', 'PAID');

const cache = derive();
cache.set('order-1', 'SHIPPED');
console.log('источник:    ' + show(source));
console.log('кеш:         ' + show(cache));
console.log('сходятся: ' + same(cache, source));

const rebuilt = derive();
console.log('пересобрали: ' + show(rebuilt));
console.log('сходятся: ' + same(rebuilt, source));
Запустить

Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →

живой пример

from dataclasses import dataclass


@dataclass(frozen=True)
class Change:
    order_id: str
    status: str


source: dict[str, str] = {}
changes: list[Change] = []


def show(m: dict) -> str:
    return "{" + ", ".join(f"{k}={v}" for k, v in m.items()) + "}"


def write(order_id: str, status: str) -> None:
    source[order_id] = status
    changes.append(Change(order_id, status))


def derive() -> dict[str, str]:
    view: dict[str, str] = {}
    for c in changes:
        view[c.order_id] = c.status
    return view


write("order-1", "NEW")
write("order-2", "NEW")
write("order-1", "PAID")

cache = derive()
cache["order-1"] = "SHIPPED"
print("источник:    " + show(source))
print("кеш:         " + show(cache))
print("сходятся: " + str(cache == source).lower())

rebuilt = derive()
print("пересобрали: " + show(rebuilt))
print("сходятся: " + str(rebuilt == source).lower())
Запустить

Запуск примеров доступен в платном доступе. Там этот же код выполняется прямо в статье: редактор, запуск и проверка рядом с абзацем. Три дня бесплатно →

Запись мимо источника нигде не упала — расхождение видно только по выводу: кеш говорит SHIPPED, источник — PAID. Пересобранная из потока копия совпадает с источником всегда.

Отставание производных данных

Производное всегда отстаёт от источника — вопрос только на сколько и что с этим делать в интерфейсе.

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

Что показывать, пока не догнало. Три честных варианта. Читать источник для того конкретного случая, где свежесть обязательна (пользователь только что создал объект — покажем его из основной базы, а не из индекса). Показать данные с явной отметкой («данные на 12:05»), и тогда отставание становится видимым, а не загадочным. Отключить функцию, которая на проекцию опирается, если отставание вышло за порог, — лучше сказать «поиск временно недоступен», чем молча показывать вчерашний каталог.

Что нельзя — делать вид, что проекция и источник согласованы. Самая дорогая ошибка тут — принимать решения по производным данным: списывать товар по индексу поиска, считать деньги по аналитической витрине.

Практический минимум: метрика отставания на каждую проекцию, порог в оповещении и решение про интерфейс, записанное заранее.

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

Правило «деривировать, а не дублировать записью» звучит абсолютно, а исключения есть — они просто не отменяют источник правды.

Кэш со сквозной записью. Записали в базу и тут же положили значение в кэш — это не второй источник правды, а способ не ждать первого чтения. Работает потому, что кэш в любой момент можно выбросить целиком.

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

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

Общий признак законного случая: производное всё ещё можно выбросить и пересобрать, а его расхождение с источником не меняет решений. Как только запись в производное становится единственным местом, где что-то хранится (счётчик, который никто не умеет пересчитать), правило нарушено — и однажды это обнаружится расхождением, которое нечем починить.

Как выбирать между тремя режимами

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

По запросу (online). Считаем в момент обращения. Дёшево в эксплуатации (ничего не надо поддерживать), данные всегда свежие, но каждый запрос платит временем. Годится, пока расчёт укладывается в бюджет ответа и запросов немного.

Пакетно (batch). Пересчитываем целиком по расписанию. Самый простой и надёжный режим: одна программа, воспроизводимый результат, легко перезапустить. Задержка равна периоду, стоимость — полный проход по данным. Берут, когда задержка в час-сутки допустима.

Потоком (stream). Обновляем по событиям. Задержка секунды, полного прохода нет; зато появляется состояние обработчика, порядок событий, повторы, догон после сбоя — то есть самая дорогая эксплуатация из трёх.

Признаки выбора: свежесть измеряется в секундах — поток; в часах — пакет; расчёт дешёвый и данные редко просят — по запросу. И частый рабочий вариант — пара «поток плюс ночной пакетный пересчёт»: поток даёт свежесть, пакет чинит накопленные расхождения. Как устроен второй — в разделе про пакетную обработку ниже, как первый — в статье про обработку потоков.

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

«Производное можно выбросить» — обещание, которое либо проверено, либо нет. Непроверенное почти всегда оказывается ложным: скрипт пересборки написан год назад, с тех пор схема изменилась, и в момент аварии он падает на третьей минуте.

Поэтому пересборку выполняют регулярно и не в аварии: на тестовом стенде по расписанию (полный пересчёт раз в неделю с проверкой результата) или в проде для одной небольшой части данных. Заодно измеряют время — «пересобрать поисковый индекс за четыре часа» это факт, который нужен для решения в момент инцидента, а не предположение.

Где это применяется

Как только рядом с базой появляется кеш, поисковый индекс или отдельная витрина — вы уже управляете производными данными, даже если не называли их так. Вот на чём спотыкаются:

  • Двойная запись в базу и в кеш/индекс — две операции без общей транзакции рано или поздно разойдутся.
  • Нет процедуры пересборки производного набора — и «переиндексировать с нуля» вдруг оказывается невозможным, а потеря индекса становится инцидентом на ровном месте.
  • Считают read-модель или денормализованное поле «второй правдой» — и при расхождении начинают гадать, кто прав.
Дополнительно: при первом чтении можно пропустить

Глубже: пакетная обработка: ночная пересборка и повторный прогонрасширенное

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

Идемпотентность повторного прогона. Задача упала на середине, её перезапустили, и витрина получила половину строк дважды. Поэтому пакет никогда не «дописывает», а заменяет целиком свою единицу работы: витрину за день считают во временную таблицу и подменяют секцию за этот день одной операцией (ALTER TABLE ... ATTACH PARTITION в PostgreSQL, REPLACE PARTITION в ClickHouse), либо пишут через INSERT ... ON CONFLICT DO UPDATE по ключу строки. Правило проверки: прогон той же задачи за тот же день дважды даёт тот же результат.

Единица работы это интервал по времени события, а не «всё, что накопилось». Задача «пересчитать за 23 сентября» может быть запущена 24-го ночью, повторена днём после исправления ошибки и запущена за прошлый месяц при пересборке истории. Задача «обработать всё новое» этого не умеет, и после сбоя не знает, откуда продолжать. Отсюда параметр запуска, дата или интервал, и планировщик, который умеет запустить задачу за прошедшие интервалы по очереди (в Airflow это и есть основа модели, в Spring Batch параметр задания, в cron придётся писать самим).

Опоздавшие данные. Событие за 23-е пришло 25-го, а витрина за 23-е уже посчитана. Пакетный ответ на это: пересчитывать не один последний день, а окно в несколько дней, и принять, что цифры за вчера чуть поменяются к послезавтра. Ширину окна выбирают по доле опозданий, а границу «после этого не пересчитываем» пишут в описание отчёта.

Соединение больших наборов. Пакет соединяет таблицу событий за день с таблицей клиентов целиком, и делать это по одной строке через сервис нельзя: миллион запросов к базе за ночь. Пакет живёт там, где данные: INSERT ... SELECT с JOIN в самой базе, запрос в ClickHouse, задание Spark над файлами. Сервис в пакетной обработке обычно лишь запускает и следит, а не перекладывает строки.

И наблюдаемость: у витрины есть отметка «данные по такое-то время», и сигнал тревоги срабатывает на её отставание, а не на падение задачи, потому что задача может «успешно» посчитать ноль строк. Как пакет устроен в коде сервиса, показывает статья про фоновую пакетную обработку, здесь важна модель, которая от языка не зависит.

Глубже: переезд на второе хранилище: перезаливка, догон, сверка и переключениерасширенное

Каждая развилка «PostgreSQL или X» заканчивается словами «нужен конвейер». Вот как выглядит переезд существующих данных в новое хранилище (в поисковый индекс, в ClickHouse, в кэш, в новую базу), когда старое продолжает работать.

Перезаливка (backfill). Берут согласованный снимок источника, для PostgreSQL это транзакция с REPEATABLE READ или реплика с известной позицией WAL, и переливают его целиком пакетами по первичному ключу. Запоминают позицию, на которой снимок был сделан: номер LSN, отметку времени, offset в Kafka. Тысячи строк в секунду это часы для миллионов и дни для миллиардов, и на это время нужна отдельная реплика для чтения, чтобы не положить прод.

Догон. Пока шла перезаливка, источник менялся. Изменения с запомненной позиции читают через CDC (Debezium) или из таблицы outbox и применяют к новому хранилищу в том же порядке. Применение обязано быть идемпотентным: строка из снимка и та же строка из потока изменений дадут одно состояние. Догон догоняет, когда отставание меньше секунд, и с этого момента новое хранилище живёт в режиме «стрим» из статьи про производные данные.

Двойное чтение и сверка. Прежде чем кому-то показать новое хранилище, из него читают параллельно со старым, молча, и сравнивают: доля запросов, где ответы разошлись, и какие именно. Сверку данных делают не «по всей таблице», а по контрольным суммам на диапазон ключей: сумма и число строк за каждый день или каждые сто тысяч идентификаторов в обоих хранилищах, расхождения ищут внутри разошедшегося диапазона. Сверка идёт до тех пор, пока расхождения не объяснены (обычно это опоздавшие изменения или разная обработка NULL).

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

Это тот же набор шагов при переезде из Oracle в PostgreSQL, из MongoDB в PostgreSQL и из монолитной базы в базу нового сервиса; меняются инструменты, не порядок.

Коротко

  • Система записи — хранилище, где факт лежит ровно один раз; при расхождении права она.
  • Производные данные (кеш, индекс, материализованное представление, read-модель, витрина) избыточны: они выводятся из источника.
  • Производное одноразово: выбросил и пересобрал. Незаменим только источник — его и бэкапят в первую очередь.
  • У каждого производного набора должна быть процедура пересборки с нуля — её выполняют регулярно и измеряют время: непроверенная пересборка в аварии не работает.
  • Три режима выбирают по требуемой свежести: секунды — поток, часы — пакет, дёшево и редко — по запросу; частая пара в жизни — поток плюс ночной пересчёт.
  • Пишем только в источник; производное обновляем из единого потока изменений (события/outbox/CDC), а не второй операцией рядом.
  • Пакет заменяет свою единицу работы целиком (секция за день), принимает интервал параметром, пересчитывает окно в несколько дней и живёт там, где данные; тревога по отставанию витрины, а не по падению задачи.
  • Переезд в новое хранилище: перезаливка со снимка с запомненной позицией, догон изменений через CDC, сверка контрольными суммами по диапазонам, переключение флагом с обратным потоком для отката.
  • У каждой проекции записывают допустимое отставание, метрику и решение для интерфейса: читать источник, показывать отметку времени или отключать функцию.
  • Писать в производное напрямую законно, пока его можно выбросить и пересобрать (кэш со сквозной записью, кэш запроса, оптимистичный интерфейс).

Что почитать дальше

  • read-модель в CQRS (Java, Go, Node, Python) — производные данные под конкретный экран.
  • Потоковая обработка — CDC и время в потоках: как устроен сам поток изменений.
  • Инвалидация кеша (Java, Go, Node, Python) — как держать производную копию свежей.
  • Event sourcing — поток событий как источник правды.