Агрегации, столбцы, окна и ANN: практикум
Оглавление · MapReduce · Векторы
JavaScript-части работают с маленьким воспроизводимым набором. Цель — сравнить результаты, потом стоимость и гарантии. Микробенчмарк трёх строк не выбирает production-СУБД. SQL выполняется в отдельной учебной базе; pgvector и ClickHouse нужны только для соответствующих частей, основной проект их не получает.
Один набор в JavaScript и SQL
import assert from "node:assert/strict";
const events = [
{ id: 1, topic: "react", value: 10 },
{ id: 2, topic: "react", value: 30 },
{ id: 3, topic: "data", value: 5 },
];
function aggregate(rows) {
const groups = new Map();
for (const row of rows) {
const pair = groups.get(row.topic) ?? { sum: 0, count: 0 };
pair.sum += row.value;
pair.count++;
groups.set(row.topic, pair);
}
return [...groups]
.sort(([a], [b]) => a.localeCompare(b))
.map(([topic, pair]) => ({ topic, ...pair, average: pair.sum / pair.count }));
}
assert.deepEqual(aggregate(events), [
{ topic: "data", sum: 5, count: 1, average: 5 },
{ topic: "react", sum: 40, count: 2, average: 20 },
]);
console.log(aggregate(events));
CREATE SCHEMA lesson_analytics;
CREATE TABLE lesson_analytics.event (
id integer PRIMARY KEY, topic text NOT NULL, value integer NOT NULL
);
INSERT INTO lesson_analytics.event VALUES (1,'react',10),(2,'react',30),(3,'data',5);
SELECT topic, sum(value), count(*), avg(value)
FROM lesson_analytics.event GROUP BY topic ORDER BY topic;
Получится data: 5/1/5, react: 40/2/20. Сравните с локальной моделью MapReduce:
mapper выдаёт (topic, (value, 1)), combiner/reducer суммирует пары, деление
выполняется в конце. SQL NULL и JS null не имеют общей семантики арифметики:
в этом контракте value обязателен. Для float задайте допустимую погрешность.
Столбцовое чтение в ClickHouse
В отдельном учебном ClickHouse выполните:
CREATE TABLE lesson_event (id UInt64, topic String, value Int32)
ENGINE = MergeTree ORDER BY (topic, id);
INSERT INTO lesson_event VALUES (1,'react',10),(2,'react',30),(3,'data',5);
SELECT topic, sum(value), count(), avg(value) FROM lesson_event GROUP BY topic ORDER BY topic;
Первый итог тот же. Но ORDER BY этого engine не является UNIQUE constraint: повтор INSERT может удвоить результат. Сравните с PostgreSQL PRIMARY KEY, который отвергнет повтор ID. Дедупликация ingestion, UPDATE и фоновые merges имеют контракты конкретного engine, не следуют из слова «столбцовый».
Создайте большой учебный набор с лишними текстовыми колонками, запросите только
topic/value, измерьте rows/bytes read и CPU на одинаковом фильтре. Отдельно
измерьте запись одной строки, большой batch и актуальность аналитической копии.
Измерение не обязано показать превосходство одного продукта во всех задачах.
Столбцовая физическая организация не равна wide-column модели Cassandra.
Окна по времени события
Для окна [0,10) событие с eventTime 8 может прийти после события с eventTime 12.
Processing time — время получения, event time — значение из события. Watermark
явно задаёт границу принятия поздних данных, а не доказывает отсутствие опозданий.
import assert from "node:assert/strict";
let watermark = -Infinity;
const counts = new Map();
const seen = new Set();
const late = [];
function accept(event) {
if (seen.has(event.id)) return;
seen.add(event.id);
const start = Math.floor(event.time / 10) * 10;
if (start + 10 <= watermark) {
late.push(event.id);
return;
}
counts.set(start, (counts.get(start) ?? 0) + 1);
}
accept({ id: 1, time: 2 });
accept({ id: 2, time: 12 });
accept({ id: 3, time: 8 });
accept({ id: 3, time: 8 });
watermark = 10;
accept({ id: 4, time: 9 });
assert.deepEqual(
[...counts],
[
[0, 2],
[10, 1],
],
);
assert.deepEqual(late, [4]);
console.log("Windows:", [...counts], "late:", late);
Здесь policy — отложить позднюю запись, а не пересчитать закрытое окно. Другая policy могла бы выпустить исправление результата. Для реального потока нужны checkpoint, восстановление seen/counts, ограничение памяти, согласованная фиксация offset и результата. Повтор после рестарта без сохранённого seen удвоит счётчик. Изоляция локального алгоритма не моделирует кластер Spark/Flink.
Исходник схемы
flowchart LR Arrival[Порядок поступления] --> Dedup[Дедупликация ID] Dedup --> Window[Окно event time] Watermark[Watermark и policy] --> Window Window --> Result[Итог или исправление]
Точный поиск и ANN pgvector: проектирование опыта
Три строки ниже — подготовленный пример запросов. Сравнение качества ANN на большом наборе — задание на проектирование опыта: генератор набора и 100 векторов запросов читатель готовит сам. Зафиксируйте seed, размерность, нормализацию, метрику, k и версии; сохраните входы, чтобы baseline и ANN читали одинаковые данные. Среда готова после проверки расширения, исходного набора и exact-запроса; на основном PostgreSQL CMS расширение не устанавливайте.
Нужен PostgreSQL 17 с установленным расширением pgvector; для воспроизводимости
фиксируйте версию расширения и SELECT extversion FROM pg_extension.
Команды ниже используют HNSW и cosine; базовая точная геометрия уже разобрана
в главе о векторах.
CREATE EXTENSION IF NOT EXISTS vector;
CREATE SCHEMA lesson_vector;
CREATE TABLE lesson_vector.item (id integer PRIMARY KEY, embedding vector(3) NOT NULL);
INSERT INTO lesson_vector.item VALUES (1,'[1,0,0]'),(2,'[0.9,0.1,0]'),(3,'[0,1,0]');
SELECT id, embedding <=> '[1,0,0]'::vector AS distance
FROM lesson_vector.item ORDER BY embedding <=> '[1,0,0]'::vector, id LIMIT 2;
CREATE INDEX item_hnsw ON lesson_vector.item USING hnsw (embedding vector_cosine_ops);
ANALYZE lesson_vector.item;
EXPLAIN (ANALYZE, BUFFERS) SELECT id FROM lesson_vector.item
ORDER BY embedding <=> '[1,0,0]'::vector LIMIT 2;
Точный top-2 — IDs 1, 2. На трёх строках planner может выбрать Seq Scan, поэтому появление индекса не доказывает ANN-выполнение. Создайте большой детерминированный набор и 100 фиксированных query vectors; для baseline в отдельной транзакции отключите index scans, получите точные top-k и план. Для ANN верните настройки и убедитесь по плану, что используется HNSW. Tie-breaking baseline задайте по ID; ANN membership при равных расстояниях может отличаться без потери релевантности — определите policy оценки.
Для каждого query вычислите recall@k = |exact ∩ approximate| / k, сравните
распределение recall, p50/p95 latency, время построения и размер индекса.
Меняйте hnsw.ef_search в пределах session/transaction. При фильтре обычный ANN
может вернуть меньше k результатов: фильтрация и candidate budget взаимодействуют.
Проверьте редкий фильтр, свежую вставку, удаление и несовпадение размерности.
Cosine для нулевого вектора не имеет обычного геометрического смысла;
валидируйте такие входы и смотрите правила индексирования pgvector.
Приёмка: одинаковые JS/SQL/MapReduce итоги, policy поздних событий, фактические планы exact/ANN и таблица recall/latency. SQL базовой агрегации можно проверить embedded PostgreSQL; ClickHouse, pgvector и кластерный runtime требуют своих отдельных запусков. Не заявляйте результаты ANN по одному наличию CREATE INDEX.
Источники: MergeTree, pgvector, watermarks Flink.