Настоящее приложение почти никогда не держит данные в одном месте. Рядом с основной базой обычно живут кеш, поисковый индекс, аналитическое хранилище, денормализованные таблицы, read-модели. В этом зоопарке легко утонуть и перестать понимать, кто кому хозяин и какой копии верить при расхождении.
Есть простая рамка, которая сразу наводит порядок: разделить все хранилища на два вида — систему записи и производные данные.
Приложение пишет в одно место — в систему записи. Дальше изменение расходится по производным наборам единым потоком, поэтому у кеша, индекса и read-модели один и тот же порядок обновлений. Потеря производного набора не теряет данных: его выбрасывают и выводят из источника заново. Потеря самого источника — уже потеря данных.
Источник правды и его отражения
Система записи (source of truth, источник правды) — это хранилище, где факт лежит ровно один раз и в авторитетном виде. Обычно это ваша нормализованная база: пользователь оформил заказ — заказ записан сюда. Важная оговорка: одного источника правды «на всю компанию» не бывает. Когда сервисов несколько, у каждого факта свой хозяин — заказ живёт в сервисе заказов, платёж в платёжном, профиль в профильном, — и источник правды ищут не «где-то в компании», а для каждого факта отдельно. И если данные в источнике вдруг разошлись с данными в любом другом хранилище, правым по определению считается источник.
Производные данные — это результат преобразования данных из источника, который всегда можно пересчитать заново. Сюда попадает почти всё остальное:
- кеш — быстрая копия части данных;
- поисковый индекс (Elasticsearch, GIN) — те же данные, переупакованные под поиск по тексту;
- материализованное представление или денормализованное поле — заранее посчитанный результат;
- read-модель в CQRS (Java, Go, Node, Python) — данные, разложенные ровно под конкретный экран;
- аналитическое хранилище (ClickHouse) — те же события, но под сканы.
Ключевой признак производных данных — избыточность: они дублируют то, что уже есть в источнике, ради скорости чтения.
Главное свойство: производное можно выбросить
Из этого определения следует вывод, который меняет отношение к данным: производные данные одноразовы. Потеряли кеш, индекс или read-модель — не катастрофа: пересобрали из источника, и система как новенькая. «Не катастрофа» не значит «бесплатно»: переиндексация полусотни миллионов документов идёт часами и в ночное окно может не влезть, так что время пересборки надо знать заранее, а не выяснять в момент аварии. А вот потеряли саму систему записи — вот это настоящая беда, восстанавливать неоткуда.
Отсюда и остальное: бэкапить и защищать в первую очередь нужно источник правды. У каждого производного набора должна быть процедура пересборки с нуля (переиндексировать, прогреть кеш, перестроить проекцию).
Три способа получать производные данные
Как данные текут из источника в производные наборы? Есть три режима обработки, и различать их полезно:
- Online (сервис). Ждёт запрос и отвечает как можно быстрее; главное — время отклика и доступность. Так работает само приложение. Но для получения больших производных наборов online не подходит.
- Batch (пакетная обработка). Берёт большой набор данных на вход, молотит его минуты или часы, выдаёт результат; главное тут — пропускная способность, а запускается это по расписанию (ночная переиндексация, пересчёт витрины). Классика больших данных — MapReduce/Spark; в обычном бэкенде проще — фоновый обработчик на очереди.
- Stream (потоковая обработка). Нечто среднее: реагирует на события по мере их поступления, с низкой задержкой, но обрабатывает не весь набор сразу, а поток изменений. Так производные данные держат свежими почти в реальном времени.
Дисциплина: деривировать, а не дублировать записью
Главная ошибка при работе с производными данными — писать в них напрямую, в обход источника. Классический пример — двойная запись: приложение пишет и в базу, и в кеш (или в поисковый индекс) двумя отдельными операциями. Рано или поздно одна операция пройдёт, а вторая упадёт — и данные молча разъедутся.
Приложение пишет двумя отдельными операциями: в базу заказ лёг как 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 — поток событий как источник правды.