Перейти к основному содержимому

Руководство по обеспечению надежности

Механизм для ограничения входящих запросов Rate Limiter​

Для ограничения количества запросов в единицу времени GradeLy использует java-based и/или SRLS-based решение.

Java-based:

  • фильтрация происходит на уровне сервлет-контейнера, изолируя основное приложение;
  • количество проходящих Rate Limiter запросов в секунду отбрасываются в качестве метрики;
  • конфигурация задается через ConfigMap в DropApp.

SRLS-based:

  • фильтрация происходит на уровне отдельного сервиса SRLS, на моменте получения запроса граничным прокси (IGEG);
  • не требует повторного развертывания приложения;
  • конфигурация задается через ConfigMap в DropApp.

С помощью Rate Limiter можно точечно настроить:

  • защиту конкретных эндпойнтов;
  • максимальное число запросов в промежуток времени;
  • ограничение трафика по CN отправителя;
  • промежуток в единицах времени, для которого выполняется ограничение.

Механизм Rate Limiter одинаково подключается и эксплуатируется для Console (java-based и SRLS) и Worker(java-based).

Определение максимальной нагрузки​

Для определения максимальной нагрузки, которую следует задать с помощью механизма Rate Limiter, проведите тест определения максимальной производительности.

Внимание!

Нагрузка должна быть ограничена с помощью механизма Rate Limiter на уровне, не превышающем показатели теста максимальной производительности.

Сценарий тестирования​

При тестировании происходит пошаговое увеличение нагрузки с нуля до предельной (с шагом 20% от плановой). Пошаговое увеличение происходит до тех пор, пока не нарушится критерий успешности по количеству ошибок и/или времени отклика (что наступит раньше). Время работы теста на каждом шаге после стабилизации нагрузки составляет 10 минут.

Цель тестирования​

По результатам тестирования устанавливается:

  • уровень нагрузки L0 — последний шаг нагрузки, на котором не были нарушены критерии успешности;
  • уровень нагрузки Llim — предельный уровень нагрузки, при котором не был нарушен критерий по количеству ошибок;
  • уровень утилизации CPUlim — утилизация CPU на уровне нагрузки Llim.

Ожидаемый результат​

  • Определен уровень максимальной производительности. L0 удовлетворяет одному из следующих уровней нагрузки (при использовании ресурсов всеми подами не больше, чем на OSE):

    • Lmax для предыдущего релиза;
    • Запросите аналитическую оценку команды по плановой нагрузке на сервис на год вперед, если не достигнут Lmax (при суммарных ресурсах как на OSE);
    • Используйте критерий 10х от промышленной нагрузки на сервис, если не достигнута плановая нагрузка на год вперед.
  • Определен уровень нагрузки Llim.

  • Определен уровень утилизации CPUlim.

Профиль нагрузки​

Моделирование нагрузки производится с использованием средств нагрузочного тестирования путем эмуляции действий определенного количества пользователей. Каждый виртуальный пользователь (программный процесс, эмулирующий действия физического пользователя АС) циклически выполняет пользовательский сценарий.

Подача нагрузки​

  • JMeter подает нагрузку на БД Source Pangolin, выполняя PL/pgSQL-функцию NewOrder, которая включает в себя три операции INSERT и одну операцию UPDATE;
  • Первый модуль захвата и изменений GraDeLy вычитывает сгенерированные данные из БД Source Pangolin и записывает их в Kafka;
  • Модуль применения изменений вычитывает данные из Kafka и применяет их на БД Replica Pangolin.
Функция NewOrder — это бизнес-транзакция для выполнения нового заказа с помощью одной транзакции базы данных.
Представляет собой транзакцию чтения-записи среднего веса с высокой частотой выполнения. Эта транзакция является основой рабочей нагрузки.
Позволяет моделировать переменную нагрузку, характерную для оперативной активности базы данных в производственных средах.

Структура транзакции NewOrder​

CREATE OR REPLACE FUNCTION benchbase.new_order(w_id smallint, d_id bigint, c_id integer, o_entry_d timestamp without time zone, i_ids integer[], i_w_ids integer[], i_qtys integer[])
RETURNS benchbase.t_order_res
LANGUAGE plpgsql
AS $function$
DECLARE
getWarehouseTaxRate text := 'SELECT W_TAX FROM benchbase.WAREHOUSE WHERE W_ID = $1';
getDistrict text := 'SELECT D_TAX, nextval(''benchbase.D_NEXT_O_ID'') FROM benchbase.DISTRICT WHERE D_ID = $1 AND D_W_ID = $2';
getCustomer text := 'SELECT C_DISCOUNT, C_LAST, C_CREDIT FROM benchbase.CUSTOMER WHERE C_W_ID = $1 AND C_D_ID = $2 AND C_ID = $3';
createOrder text := 'INSERT INTO benchbase.OORDER (O_ID, O_D_ID, O_W_ID, O_C_ID, O_ENTRY_D, O_CARRIER_ID, O_OL_CNT, O_ALL_LOCAL) VALUES ($1, $2, $3, $4, $5, $6, $7, $8)';
createNewOrder text := 'INSERT INTO benchbase.NEW_ORDER (NO_O_ID, NO_D_ID, NO_W_ID) VALUES ($1, $2, $3)';
getItemInfo text := 'SELECT I_PRICE, I_NAME, I_DATA FROM benchbase.ITEM WHERE I_ID = $1';
getStockInfo text := 'SELECT S_QUANTITY, S_DATA, S_YTD, S_ORDER_CNT, S_REMOTE_CNT, S_DIST_%s FROM benchbase.STOCK WHERE S_I_ID = $1 AND S_W_ID = $2';
updateStock text := 'UPDATE benchbase.STOCK SET S_QUANTITY = $1, S_YTD = $2, S_ORDER_CNT = $3, S_REMOTE_CNT = $4 WHERE S_I_ID = $5 AND S_W_ID = $6';
createOrderLine text := 'INSERT INTO benchbase.ORDER_LINE (OL_O_ID, OL_D_ID, OL_W_ID, OL_NUMBER, OL_I_ID, OL_SUPPLY_W_ID, OL_DELIVERY_D, OL_QUANTITY, OL_AMOUNT, OL_DIST_INFO) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10)';

all_local boolean := TRUE;
ol_cnt integer;
items benchbase.t_order_item[];
item benchbase.t_order_item;
customer_info benchbase.t_order_customer;
c_discount float;
o_carrier_id integer;
stock_info benchbase.t_order_stock;
w_tax float;
d_tax float;
d_next_o_id integer;
brand_generic char;
ol_amount float;
total float := 0;
item_data_s benchbase.t_order_item_data;
item_data benchbase.t_order_item_data[];
misc benchbase.t_order_misc;

BEGIN
ol_cnt = array_length(i_ids, 1);
ASSERT ol_cnt != 0;
ASSERT ol_cnt = array_length(i_w_ids, 1);
ASSERT ol_cnt = array_length(i_qtys, 1);

FOR i IN 1..ol_cnt LOOP
all_local = all_local AND i_w_ids[i] = w_id;
EXECUTE getItemInfo INTO item USING i_ids[i];
ASSERT item IS NOT NULL;
items = array_append(items, item);
END LOOP;
ASSERT ol_cnt = array_length(items, 1);
EXECUTE getWarehouseTaxRate INTO w_tax USING w_id;
EXECUTE getDistrict INTO d_tax, d_next_o_id USING d_id, w_id;
EXECUTE getCustomer INTO customer_info USING w_id, d_id, c_id;
c_discount = customer_info.C_DISCOUNT;
o_carrier_id = 0;
EXECUTE createOrder USING d_next_o_id, d_id, w_id, c_id, o_entry_d, o_carrier_id, ol_cnt, cast(all_local as integer);
EXECUTE createNewOrder USING d_next_o_id, d_id, w_id;

FOR i IN 1..ol_cnt LOOP
EXECUTE FORMAT(getStockInfo, to_char(d_id, 'FM09')) INTO stock_info USING i_ids[i], i_w_ids[i];
stock_info.S_YTD = stock_info.S_YTD + i_qtys[i];
IF stock_info.S_QUANTITY >= i_qtys[i] + 10 THEN
stock_info.S_QUANTITY = stock_info.S_QUANTITY - i_qtys[i];
ELSE
stock_info.S_QUANTITY = stock_info.S_QUANTITY + 91 - i_qtys[i];
END IF;
stock_info.S_ORDER_CNT = stock_info.S_ORDER_CNT + 1;

IF i_w_ids[i] != w_id THEN
stock_info.S_REMOTE_CNT = stock_info.S_REMOTE_CNT + 1;
END IF;

EXECUTE updateStock USING stock_info.S_QUANTITY, stock_info.S_YTD, stock_info.S_ORDER_CNT, stock_info.S_REMOTE_CNT, i_ids[i], i_w_ids[i];

IF position('ORIGINAL' IN items[i].I_DATA) != 0 AND position('ORIGINAL' IN stock_info.S_DATA) != 0 THEN
brand_generic = 'B';
ELSE
brand_generic = 'G';
END IF;

ol_amount = i_qtys[i] * items[i].I_PRICE;
total = total + ol_amount;
EXECUTE createOrderLine USING d_next_o_id, d_id, w_id, i, i_ids[i], i_w_ids[i], o_entry_d, i_qtys[i], ol_amount, stock_info.S_DIST;
item_data_s = ROW(items[i].I_NAME, stock_info.S_QUANTITY, brand_generic, items[i].I_PRICE, ol_amount);
item_data = array_append(item_data, item_data_s);
END LOOP;
total = total * (1 - c_discount) * (1 + w_tax + d_tax);
misc = ROW(w_tax, d_tax, d_next_o_id, total);
RETURN ROW(customer_info, misc, item_data);
END;
$function$
;

Вызов функции NewOrder​

select benchbase.new_order((CEIL(RANDOM() * 10))::smallint, (CEIL(RANDOM() * 10))::int8, CEIL(RANDOM() * 3000)::integer, current_timestamp::timestamp, ARRAY[CEIL(RANDOM() * 100000)::int], ARRAY[CEIL(RANDOM() * 10)::int], ARRAY[1]);

Схема данных​

Схема данных

Тестовые данные​

Тестовые данные сгенерированы с помощью BenchBase с использованием бенчмарка TPC-C.

Пример конфигурации для генерации тестовых данных​

<?xml version="1.0"?>
<parameters>


<type>POSTGRES</type>
<driver>org.postgresql.Driver</driver>
<url> </url>
<username> </username>
<password> </password>
<reconnectOnConnectionFailure>true</reconnectOnConnectionFailure>


<isolation>TRANSACTION_READ_COMMITTED</isolation>
<batchsize>100</batchsize>


<scalefactor>10</scalefactor>


<terminals>1000</terminals>
<works>
<work>
<time>1200</time>
<rate>3000</rate>
<weights>45,43,4,4,4</weights>
</work>
</works>


<transactiontypes>
<transactiontype>
<name>NewOrder</name>
</transactiontype>
<transactiontype>
<name>Payment</name>
</transactiontype>
<transactiontype>
<name>OrderStatus</name>
</transactiontype>
<transactiontype>
<name>Delivery</name>
</transactiontype>
<transactiontype>
<name>StockLevel</name>
</transactiontype>
</transactiontypes>
</parameters>

Тестирование базовой конфигурации​

grdl.properties​

gdl.apply-cache-size=20
gdl.capture-cache-size=12000
gdl.kafka-consumer-queue-size=12000

Параметры модуля применения изменений Applier​

Параметры модуля применения изменений Applier

Параметры соединения БД Replica​

{
"max.pool.size": 20,
"apply.thread.count": 18,
"transaction.size": 2000,
"db.linger.ms": 3000
}

Параметры соединения Kafka​

{
"max.poll.records": "1",
"batch.size": 16384,
"max.request.size": 209715200,
"compression.type": "zstd",
"compression.zstd.level": 3
}

Метрики теста​

  • Модуль захвата изменений capture: vm-gdev-db-nosw-515.vdc11.db.dev.sbt
  • Модуль применения изменений applier: vm-gdev-db-nosw-516.vdc11.db.dev.sbt

Network recieve/transmit​

network_recieve_transmit_basic

OPS(insert/update/delete)​

  • vm-grdl-db-psql-430.vdc11.db.dev.sbt — БД Source
  • vm-grdl-db-psql-433.vdc11.db.dev.sbt — БД Replica

ops_basic

Итог​

  • Максимальная производительность репликации OPS: 28900
  • Максимальная производительность репликации MB/s: 8.2

Тестирование минимальной конфигурации​

grdl.properties​

gdl.apply-cache-size=20
gdl.capture-cache-size=12000
gdl.kafka-consumer-queue-size=500

Параметры модуля применения изменений Applier​

Параметры модуля применения изменений Applier

Параметры соединения БД Replica​

{
"max.pool.size": 5,
"apply.thread.count": 2,
"transaction.size": 500,
"db.linger.ms": 1000
}

Параметры соединения Kafka​

{
"max.poll.records": "1",
"batch.size": 16384,
"max.request.size": 209715200,
"compression.type": "zstd",
"compression.zstd.level": 3
}

Метрики теста​

  • Модуль захвата изменений capture: vm-gdev-db-nosw-513.vdc11.db.dev.sbt
  • Модуль применения изменений applier: vm-gdev-db-nosw-514.vdc11.db.dev.sbt

Network recieve/transmit​

network_recieve_transmit_min

OPS(insert/update/delete)​

  • vm-grdl-db-psql-430.vdc11.db.dev.sbt — БД Source
  • vm-grdl-db-psql-433.vdc11.db.dev.sbt — БД Replica

ops_min

Итог​

  • Максимальная производительность репликации OPS: 3900
  • Максимальная производительность репликации MB/s: 1.15

Рекомендации по настройке очередей​

Параметры конфигурации могут быть заданы:

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

Параметры соединения могут быть заданы в UI консоли GRDL.

Обозначения и типы параметров​

1) Параметры ConfigMap модулей GraDeLy (очереди/кэши модулей)​

ПараметрОписание
capture-cache-sizeРазмер внутреннего кэша потока pgoutput-сообщений в модуле capture
consumer-cache-sizeМаксимальное количество закешированных векторов изменений, готовых к отправке в Kafka
apply-cache-sizeМаксимальное количество закешированных векторов изменений, вычитанных из Kafka, но еще не примененных

2) Параметры БД-приёмника в UI GraDeLy​

ПараметрОписание
apply.thread.countПараллелизм применения
transaction.sizeГруппировка операций при применении

3) Расчетные величины (для сайзинга)​

ПараметрОписание
op_size_maxМаксимальный размер одной операции, МБ (определяется по результатам тестирования)
ops_in_tx_maxМаксимальное число операций в транзакции (определяется по результатам тестирования)

Дополнительные расчетные величины (оценки памяти)​

Ниже приведены расчетные величины, используемые для оценки потребления памяти отдельными очередями/кэшами модулей GraDeLy.
Это не параметры конфигурации, а производные оценки, вычисляемые по формулам на основе выбранных значений параметров и принятых верхних оценок нагрузки.

ПараметрОписание
mem_capture_cacheОценка памяти, занимаемой внутренним кэшем потока pgoutput-сообщений в модуле capture, в МБ
mem_capture_kafka_queueОценка памяти, занимаемой очередью подготовленных векторов изменений в модуле capture (буфер перед отправкой в Kafka), в МБ
mem_applier_kafka_queueОценка памяти, занимаемой очередью векторов изменений в модуле applier (вычитаны из Kafka, но еще не применены), в МБ
mem_apply_cacheОценка памяти, занимаемой кэшем подготовленных bulk-запросов к БД-приёмнику в модуле applier, в МБ

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

При выборе размеров очередей и кэшей ориентируемся на:

  • op_size_max
  • ops_in_tx_max
  • apply.thread.count
  • transaction.size

Параметры ConfigMap модуля capture​

capture-cache-size​

Назначение​

Размер внутреннего кэша потока pgoutput-сообщений в модуле capture.

Ограничение​

capture-cache-size >= 2 * ops_in_tx_max

где ops_in_tx_max - максимальное количество изменений в одной транзакции

Оценка памяти​

mem_capture_cache ≈ op_size_max * capture-cache-size

Пример​

  • op_size_max = 1 МБ
  • ops_in_tx_max = 500 → capture-cache-size = 500 * 2 = 1000
  • mem_capture_cache ≈ 1 МБ * 1000 = 1000 МБ (~1 ГБ)

kafka-consumer-queue-size​

Назначение​

Максимальное количество закешированных векторов изменений, готовых к отправке в Kafka.

Внимание!

Значение должно обеспечивать буферизацию как минимум одной транзакции. Рекомендуемое значение: 100 (стартовая рекомендация).

Оценка памяти​

mem_capture_kafka_queue ≈ kafka-consumer-queue-size * op_size_max

Пример​

  • op_size_max = 1 МБ
  • kafka-consumer-queue-size = 100
  • mem_capture_kafka_queue ≈ 100 МБ

Параметры очереди модуля applier​

kafka-consumer-queue-size​

Назначение​

Максимальное количество закешированных векторов изменений, вычитанных из Kafka, но еще не примененных.

Внимание!

Значение должно обеспечивать буферизацию как минимум одной транзакции. Рекомендуемое значение: 100 (стартовая рекомендация).

Оценка памяти​

mem_applier_kafka_queue ≈ kafka-consumer-queue-size * op_size_max

Пример​

  • op_size_max = 1 МБ
  • kafka-consumer-queue-size = 100
  • mem_applier_kafka_queue ≈ 100 МБ

apply-cache-size​

Назначение​

Максимальное количество подготовленных bulk-запросов к БД-реплике, ожидающих выполнения.

Рекомендуемое значение​

По одному на поток применения:

apply-cache-size = apply.thread.count

Оценка памяти​

mem_apply_cache ≈ apply-cache-size * op_size_max * transaction.size

где transaction.size — размер пачки/группы операций на применение; при этом не меньше одной транзакции источника.

Пример​

  • apply-cache-size = 10 (если apply.thread.count = 10)
  • op_size_max = 10 МБ
  • transaction.size = 20
  • mem_apply_cache ≈ 10 * 10 МБ * 20 = 2000 МБ (~2 ГБ)

Практические рекомендации​

  1. Соберите исходные оценки (op_size_max, ops_in_tx_max) на реальных данных.
  2. Посчитайте память по формулам для каждого кэша/очереди и суммарно.
  3. Заложите overhead (метаданные векторов изменений, структуры данных, GC): обычно добавляют запас ×2–×10 к расчетному объёму.
  4. Если упираетесь в память — снижайте в первую очередь:
    1. kafka-consumer-queue-size (с учетом “не меньше 1 транзакции”),
    2. затем аккуратно transaction.size,
    3. и только потом — apply.thread.count / apply-cache-size (это уже влияет на пропускную способность).

Пример расчета​

Исходные данные​

  • op_size_max = 2 МБ
  • ops_in_tx_max = 300
  • kafka-consumer-queue-size = 100
  • apply.thread.count = 8
  • transaction.size = 25

Принимаем​

  • capture-cache-size = 2 * ops_in_tx_max = 600
  • apply-cache-size = apply.thread.count = 8

Сapture — итоговая формула и расчет​

  1. mem_capture_total ≈ op_size_max * (capture-cache-size + kafka-consumer-queue-size)
  2. mem_capture_total ≈ 2 * (600 + 100) = 1400 МБ

Итого для capture:

  • ~1400 МБ (~1.4 ГБ)

Суммарно с запасом ×3:

  • 4200 МБ

Applier — итоговая формула и расчет​

  1. mem_applier_total ≈ op_size_max * (kafka-consumer-queue-size + apply-cache-size * transaction.size)
  2. mem_applier_total ≈ 2 * (100 + 8 * 25) = 2 * 300 = 600 МБ

Итого для applier:

  • ~600 МБ

Суммарно с запасом ×3:

  • 1800 МБ

Настройка механизма для ограничения входящих запросов Rate Limiter​

Для Java-based:

Задайте следующие параметры в ConfigMap консоли или воркера:

ПараметрТипЗначение по умолчаниюОписание
grdl-console.rate-limiter-enabledbooleantrueВключает и выключает Rate Limiter для консоли
grdl-console.rate-limiter-metric-prefixstringgrdl_consoleПриставка для имени метрики, которая показывает, сколько запросов прошли через Rate Limiter на консоли.
Например, при значении по умолчанию полное название метрики будет: grdl_console.rate_limiter.requests
grdl-console.rate-limiter-bucketsjson[{"uri":"/general", "maxTokens":1000, "refillIntervalSec":1}]}Конфигурация Rate Limiter для каждого эндпойнта консоли
grdl-module.rate-limiter-enabledbooleanfalseВключает и выключает Rate Limiter для воркера
grdl-module.rate-limiter-metric-prefixstringgrdl_moduleПриставка для имени метрики, которая показывает, сколько запросов прошли через Rate Limiter на воркере.
Например, при значении по-умолчанию полное название метрики будет: grdl_module.rate_limiter.requests
grdl- module.rate-limiter-bucketsjson[{"uri":"/general", "maxTokens":1000, "refillIntervalSec":1}]}Конфигурация Rate Limiter для каждого эндпойнта воркера

Json-конфигурация позволяет точечно определить ограничение по количеству вызовов для каждого эндпойнта.

Эта конфигурация представляет массив объектов, в каждом из которых указывается:

  • uri — эндпойнт, который необходимо защитить. Поддерживаются regexp, например, /module/.*/state. Есть ключевое значение /general, с помощью него все эндпойнты, которые не заданы явно, объединяются в одну группу с общим лимитом;
  • maxTokens — максимальное количество запросов, которое пользователь разрешает к эндпойнту;
  • refillIntervalSec — интервал в секундах, за который доступное количество запросов восстанавливается.
примечание

Если maxTokens=1000, а refillIntervalSec=1, пользователь не ждет, пока пройдет секунда, — при каждом запросе увеличиваем доступное количество запросов на величину, пропорциональную прошедшему времени. Например, если с момента прошлого запроса прошло пол секунды, к доступному лимиту добавляется 500 запросов.

Важно, чтобы этот json представлял из себя одну строку без пробелов.

Для SRLS-based:

Задайте следующие параметры в ConfigMap консоли или воркера:

ПараметрТипЗначение по умолчаниюОписание
grdl.srls.enabledbooleantrueВключает и выключает Rate Limiter для консоли
grdl.k8s.common.srls.base.headerstringiv-userНазвание заголовка которая идентифицирует уникального пользователя
grdl.k8s.common.srls.console.headerstringx-console-srlsНазвание заголовка которая идентифицирует запросы для API console
grdl.k8s.common.srls.worker.headerstringx-srlsНазвание заголовка которая идентифицирует запросы от модулей
grdl.k8s.istio.ingress.srls.shortnamestringgrdlДескриптор для определения правил SRLS Rate Limits
grdl.k8s.globalRateLimit.spec.envoyVersionstring1.25Версия Envoy для SRLS
grdl.k8s.common.srls.serverstring``Адрес сервиса для подключения к SRLS
grdl.k8s.common.srls.portstring8081Порт для подключения к SRLS
grdl.k8s.common.srls.namespacestring``Namespace k8s SRLS
grdl.srls.endpointsjson[{"name":"grdl","shortname":"grdl","overall_limit":10000,"by_header":{"header":"iv-user,x-console-srls,x-srls","value":500,"anon_value": 0,"invokers":{{"header_value":"2","name":"console","value":500},{"header_value":"1","name":"worker","value":500}}}}]Конфигурация Rate Limiter by header для endpoint консоли
grdl-console.envoy-cn-enablebooleantrueВключает и выключает защиту от DDOS для консоли
grdl-console.envoy-cn-liststring^.*solution%.sbt$Список доступных доменов, которые не блокируются, через запятую

Json-конфигурация позволяет гибко настравить Rate Limits для эндпойнтов.

Эта конфигурация представляет массив объектов, которые соотвествуют спецификации SRLS.

В примере, указанном в дистрибутиве - в SRLS ставится дескриптор grdl и любой трафик которые не будет иметь ниодного из заголовков iv-user,x-console-srlsx-srls будет заблокирован, в обратном случее применится общий rate limit - 500 запрос в сек (значение по умолчанию), если один из заголовков будет иметь значение "1" или "2" то квота для него будет пересмотрена.

Защита от DDoS организована через EnvoyFilter lua скриптом, поэтому при формировании списка grdl-console.envoy-cn-list нужно учитывать синтаксис lua (например изоляцию спецсимволов через %), поддерживает регулярные выражения, в значении по умолчанию wildcard для определенных доменов 1 и 2-го уровня.

Bypass некритичных сервисов​

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

  • One-Time Password (OTP) / OTT;

  • Мониторинг (MONA) и журналирование (LOGA).

    Внимание!

    Если мониторинг и журналирование не доступны, информация о событиях записывается в локальное хранилище. Размер локального хранилища зависит от настроек k8s.

    После переполнения памяти локального хранилища, консоль перестанет функционировать, но репликация продолжится за счет принципа непрерывности репликации, обеспечиваемого архитектурой GraDeLy.

Непрерывность репликации​

Непрерывность репликации GraDeLy при недоступности UI/Консоли/репозитарной БД обеспечивается за счет архитектуры: воркеры работают изолированно друг от друга и от консоли. После создания задачи репликации они выполняют ее до получения следующей команды от консоли.

Гарантированная доставка данных​

GraDeLy обеспечивает сохранность, целостность транзакций и защиту данных, в том числе гарантированную доставку, за счет принципа непрерывности репликации и отказоустойчивости Kafka.

Убедиться в полноте транзакционной целостности можно за счет:

Параметра flush_lsn на базе-источнике​

flush_lsn, или сброс lsn, внутренней метки (Log Sequence Number), обозначающей позицию в журнале транзакций (WAL), до которой данные уже безопасно переданы и подтверждены, при чтении слота репликации гарантированно происходит, только после получения ответа от Kafka об успешной доставке сообщения в топик.

Параметров acks=all и min.insync.replicas в Kafka​

Параметр acks определяет, сколько брокеров должны подтвердить получение сообщения до того, как оно будет считаться успешно записанным.

  • acks=0: сообщение считается доставленным сразу после отправки, подтверждение от брокера не требуется.
  • acks=1: сообщение подтверждается, когда его принял только лидер партиции.
  • acks=all: сообщение подтверждается только после записи во все актуальные in-sync реплики (ISR).

В версии Kafka >= v3.0 значение acks по умолчанию all.

Параметр min.insync.replicas определяет минимальное количество in-sync реплик, которые должны подтвердить получение сообщения, чтобы оно считалось доставленным.

Например, при min.insync.replicas=2 сообщение считается доставленным, если минимум 2 ISR подтвердили получение данных. Если меньше ISR — продюсеры получат ошибку NotEnoughReplicasException.

Чтобы гарантировать устойчивую доставку сообщений:

  • Используйте acks=all.
  • Установите min.insync.replicas ≥ 2 на брокере/топике.

В этом случае сообщение будет записано минимум на двух репликах, и если одна из них выйдет из строя — данные не будут потеряны.

Установка min.insync.replicas на уровне брокера​

Инструкция для администраторов Kafka

  1. Установка параметра min.insync.replicas:

    kafka-configs.sh --bootstrap-server localhost:9092 --alter \
    --entity-type brokers --entity-default \
    --add-config min.insync.replicas=2
  2. Проверьте конфигурации брокера:

    kafka-configs.sh --bootstrap-server localhost:9092 --describe \
    --entity-type brokers --entity-default

Или через файл конфигурации:

Откройте config/server.properties и добавьте:

min.insync.replicas=2
Внимание!

Требуется перезапуск брокера для применения параметра.

Если Kafka не отправляет подтверждение записи модулю Capture​

  1. Настройте обработчик ошибок следующим образом:

    DML_ERROR=ABORT
  2. При остановке репликации из-за ошибки "Duplicate key value violates unique constraint" или иных ошибок, связанных с нарушением ограничений:

    • откройте топик error в Kafka;
    • найдите grdl_id сообщения об ошибке;
    • запустите модуль Applier с позиции: в качестве типа позиции выберите "gradely", в поле Позиция укажите grdl_id из сообщения об ошибке +1.
Внимание!

При других типах поведения обработчика ошибок на DML_ERROR exactly once доставка не работает. Происходит переход на at least once.

Записи GraDeLyID на базе-приемнике​

Идентификатор последней примененной транзакции GraDeLyID записывается в Kafka в столбцы key и value, а также в столбец trx_id технической таблицы $GRADELY_APPLY_POSITION$, которая находится на базе приемника.

Умная балансировка​

GraDeLy обеспечивает балансировку нагрузки между узлами за счет проверок работоспособности (healthchecks) прикладной функциональности.

Console предоставляет следующие heatlhchecks:

  • Readiness probe: uri - Settings/isReady;
  • Liveness probe: uri - Settings/isActive.

Для Worker внешние healthchecks не требуются, так как ими управляет Console.

В качестве балансировщика нагрузки может использоваться HAProxy или его аналог. Балансировка настраивается по инструкции используемого балансировщика.

Поэтапное обновление GraDeLy​

Балансировщик нагрузки и java-based Rate Limiter позволяют контролировать трафик на контуре с новой версией по нагрузке.

Принятие решения об откате или дальнейшей установке происходит по итогам анализа метрик и фона ошибок:

  • утилизация CPU и RAM подов приложения не превышает 80%;
  • приложение прошло проверку базовой работоспособности.

Проверка базовой работоспособности приложения считается успешной, если:

  • UI консоль доступна;

  • через UI можно редактировать конфигурации, соединения и графы;

  • тестирование конфигураций и соединений проходит успешно.

    подсказка

    Для экономии ресурсов в средах опытной эксплуатации не создается резервный граф для проверки воркера. Считается, что его работоспособность подтверждается на тестовых стендах. Console — малонагруженный компонент, поэтому достаточно базовой проверки работоспособности.

Промежуток времени между установкой на Stand-By и Master определяет заказчик продукта в зависимости от того, сколько времени ему потребуется на проверку базовой работоспособности приложения.

Внимание!

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

В GraDeLy при разработке используется инкрементальный подход для изменения структуры базы данных, где каждое изменение обратно совместимо. Технически это обеспечивается добавлением новых changeset в Liquibase-скрипты.

Благодаря этому дополнительных действий с базой данных, кроме непосредственного исполнения Liquibase-скриптов, при обновлении не требуется.

Базовые метрики Console и Worker​

Console передает четыре ключевые метрики:

МетрикаТехническое названиеЕдиница измерения
Количество worker, готовых взять задачу (статус READY)grdl_console_active_worker.countcount
Количество worker, выполняющих задачу (статус ACTIVE)grdl_console_ready_worker.countcount
Количество worker, к которым нет доступа (статус DETACHED)grdl_console_detached_worker.countcount
Количество произошедших ошибокgrdl_console_error.countcount

Worker передает метрики, которые относятся к процессу репликации:

МетрикаТехническое названиеЕдиница измерения
Количество транзакций в секундуgrdl_module.tpscount/second
Текущий размер очереди из транзакций на источникеgrdl_module.captureCacheSizecount
Размер очереди из транзакций на приемнике в момент записи в базуgrdl_module.applyCacheSizecount
Размер очереди транзакций на приемнике в момент чтение из Kafkagrdl_module.kafkaConsumerQueueSizecount

Инфраструктурные метрики​

GraDeLy рекомендует настроить получение следующих инфраструктурных метрик для подов с console и worker:

МетрикаЕдиница измерения
pod CPU%
CPU контейнера в поде%
pod RAM%
RAM контейнера в поде%

Порядок поэтапного обновления GraDeLy​

подсказка

Workers обратно совместимы с Console.

Обновление Console не влияет на работающую репликацию.

Обновлять workers можно поэтапно. При обновлении workers не происходит потери данных, но может возникнуть незначительная задержка в репликации, поэтому рекомендуется проводить обновление workers в техническое окно.

Предусловия:

  1. Обновление происходит при полной работоспособности Master кластера.
  2. Работает балансировщик нагрузки между кластерами.

Шаги:

  1. Запустите Liquibase-скрипты для БД с конфигурацией Console.
  2. Обновите компонент (UI-Console/Worker) на Stand-By кластере.
  3. Проведите базовую проверку работоспособности компонента на Stand-By кластере.
  4. Повторите шаги 2 и 3 для Master кластера.

Откат изменений​

При обновлении с ошибкой на любом из пунктов основного сценария проводим процедуру отката до предыдущей версии: выполните основной сценарий с последней работоспособной версией.

Поведение GraDeLy при недоступности Kafka или БД​

Вне зависимости от причины недоступности БД источника, БД приемника или Kafka, поведение воркеров GraDeLy определяется только установленным обработчиком ошибок.

подсказка

GraDeLy умеет реагировать на сбои Kafka благодаря обработчику ошибок. Для этого после сохранения схемы графа на все его соединения по умолчанию добавляется обработчик ошибок со следующими параметрами:

  • Тип: CONNECTION_ERROR;
  • Действие: RETRY;
  • Попыток реконнекта: 180;
  • Тайм-аут реконнекта (сек): 10.

Режимы обработчика ошибок​

abort​

  • Первая недоступность БД/Kafka приводит к:
    • остановке модулей на графе репликации;
    • записи ошибки в процессы и логи.
  • Репликация не продолжается автоматически.

connection_retry​

  • При недоступности сервисов выполняются попытки переподключения:
    • каждый промежуток, заданный в Тайм-аут реконекта;
    • не более, чем число попыток, указанных в Попытках реконекта.
  • Если подключение успешно, процесс восстанавливается.
  • Если все попытки исчерпаны, сценарий переходит в режим abort.

Сценарии работы​

1. Недоступность БД источника (режим connection_retry)​

  • Capture фиксирует последний commitLsn, полученный из pgoutput.
  • Capture пытается переподключаться по схеме: Тайм-аут реконекта × Попытк реконекта.
  • При успешном восстановлении соединения:
    • Capture переподключается к БД-источнику.
    • Продолжает получать изменения из pgoutput с сохраненного commitLsn.

2. Недоступность БД приемника (режим connection_retry)​

  • Applier выполняет переподключения по схеме: Тайм-аут реконекта × Попыток реконекта.
  • При успешном восстановлении соединения:
    • Applier продолжает репликацию с следующего offset из "$GRADELY_APPLY_POSITION$".

3. Недоступность Kafka (режим connection_retry)​

  • Оба воркера (Capture и Applier) переподключаются по схеме connection_retry.
  • Capture фиксирует последний commitLsn, полученный из pgoutput.
  • На время простоя Kafka:
    • Capture прекращает чтение из pgoutput.

После восстановления Kafka:

  • Capture подключается к БД источнику и продолжает получать изменения из pgoutput с сохраненного commitLsn.
  • Apply продолжает репликацию с следующего offset из "$GRADELY_APPLY_POSITION$".

Итог​

  • abort = остановка при первой же ошибке.
  • connection_retry = устойчивый режим: воркеры ждут восстановления сервисов и продолжают с последней сохраненной позиции (commitLsn или GRADELY_APPLY_POSITION).

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

Запас производительности​

Тест надежности​

Сценарий тестирования​

Тест проводится на уровне нагрузки Lstab = 70% от Lmax. Длительность стабильной нагрузки не менее 24 часов. В ходе теста фиксируются все отклонения от нормального поведения системы, в том числе деградация производительности, утечки.

Ожидаемый результат​

Не были нарушены критерии успешности. В ходе теста показатели интенсивности операций и времени отклика были стабильны. Отсутствует утечка ресурсов. В ходе теста наблюдалась хотя бы одна отработка FullGC (если этого не произошло, необходим дополнительный тест на фиксацию утечки памяти). Получен AWR-отчет за 10-минутный интервал теста. На основании AWR-отчета получено подтверждение от архитекторов БД, что тестируемое приложение не оказывает негативного влияния на работу БД.

Повышение производительности​

Для повышения производительности в GraDeLy предусмотрена возможность смены порядка транзакций.

GraDeLy оптимизирует порядок транзакций за счет функции reorder. Благодаря ей можно менять порядок команд INSERT, UPDATE, DELETE и объединять их в prepared statements. Смена порядка транзакций повышает скорость записи на 30-50%.

Для использования функции reorder должны быть установлены отложенные ограничения (deferred constraints) (инструкция по настройке приведена в «Руководстве администратора», разделе «Сценарии администрирования», подразделе «Настройка многопоточной репликации»).

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

примечание

По умолчанию "reorder": false, что значит: транзакции идут в том порядке, в котором они были сделаны.

  1. Нажмите модуль на графе.

    Соединения, Редактирование модуля

  2. Нажмите Редактировать в открывшемся окне Модуль.

    Соединения, Редактирование модуля

  3. Включите переупорядочивание операций в транзакции для оптимизации последовательности транзакций. Смена порядка транзакций повышает скорость записи на 30-50%.

    Соединения, Редактирование модуля

  4. Нажмите Сохранить.

Тест подтверждения максимальной нагрузки​

Сценарий тестирования​

Тест проводится на ступени нагрузки, предшествующей L0 (или на уровне нагрузки 90% от L0). Длительность стабильной нагрузки не менее 1 часа. Если в процессе тестирования система оказалась недогружена или перегружена, то значение нагрузки корректируется и второй тест проводится повторно. В случае увеличения нагрузки новый уровень может быть рассчитан на основе данных об утилизации ресурсов.

Ожидаемый результат​

В ходе теста зафиксирован максимальный достигнутый уровень нагрузки Lmax.

Работа GraDeLy с PostgreSQL в кластерном режиме​

Console поддерживает работу с PostgreSQL в кластерном режиме по умолчанию. Для работы воркера в кластерном режиме нужно в настройках соединения с БД изменить строку подключения по следующему шаблону (смотрите раздел Настройка соединения в Руководстве пользователя интерфейса консоли управления):

jdbc:postgresql://host1:port1,host2:port2c/database?preferQueryMode=simple&targetServerType=primary
Внимание!

Если не задать "jdbc:postgresql://host1:port1,host2:port2c/database?targetServerType=primary", при старте графа воркеры упадут с ошибками.

Обязательные параметры для работы с PostgreSQL в кластерном режиме:

  • host1:port1 - адрес узла кластера;

  • host2:port2 - адрес следующего узла кластера;

  • preferQueryMode - режим запросов к БД, ставится simple;

    Внимание!
    • Возможен только simple;
    • extended, extended_cache_everything, extended_for_prepared невозможны.

    Параметр preferQueryMode опциональный: его нужно прописать для работы в кластерном режиме. Этот режим запросов к БД немного снижает производительность, поэтому используйте его, только если хотите реплицировать из кластера.

  • targetServerType - роль участника кластера, ставится primary.

Для бесперебойной репликации после сохранения схемы графа на все его соединения по умолчанию добавляется обработчик ошибок со следующими параметрами:

  • Тип: CONNECTION_ERROR;
  • Действие: RETRY;
  • Попыток реконнекта: 180;
  • Тайм-аут реконнекта (сек): 10.

Резервирование и отказоустойчивость​

Отказоустойчивость GraDeLy обеспечивается благодаря работе в более чем одном кластере для консоли администратора (UI) и консоли управления (console), и модуля захвата и применения (worker) только в рамках одного кластера через масштабирование Deployment.

Схема, отображающая отказоустойчивость и поэтапное обновление консоли:

примечание

Схема для UI аналогична за исключением подключения к БД.

Соединения, Редактирование модуля

Console​

UI и Console работают в режиме Active-Active: один экземпляр разворачивается в DropApp кластер Master, а другой — в DropApp кластер Stand-by.

Все запросы к UI и Console обрабатываются балансировщиком, который на основе анализа healthcheck отправляет запрос на доступный кластер.

Все экземпляры Console работают с кластером БД Patroni, который воспринимается как единая БД, где хранятся данные для настроек правил репликации.

Worker​

Workers работают в режиме Active-Passive: одни экземпляры развертываются в DropApp кластер Master, а другие — в DropApp кластер Stand-by.

Все запросы к worker обрабатываются консолью, выполняющей роль балансировщика, в том же кластере, таким образом:

  • при доступности Master работают консоли Master и Stand-by, и они отправляет запросы на Master workers;
  • при недоступности Master работает консоль Stand-by и отправляет запросы на Master Stand-by.

Все экземпляры worker работают с кластером Kafka и имеют конфигурационный параметр, указывающий на кластер, в котором они работают.

Алгоритм работы при сбое Worker​

При выходе из строя worker происходит попытка самовосстановления через обработчик ошибок из Master-кластера.

Если worker не может восстановиться, все доступные на момент консоли из Master и Stand-by регистрируют его сбой.

Если в Master остались свободные workers, одному из них console назначает задачу вышедшего из строя worker.

Если весь Master недоступен или в нем не осталось свободных workers, доступная console назначает worker из Stand-by кластера задачу вышедшего из строя worker.

Восстановление Master​

Если worker из Stand-by кластера включился в репликацию из-за недоступности worker из Master, он продолжает работу даже после восстановления работоспособности worker из Master.

Обратное переключение производится вручную, после того как администраторы убедятся в работоспособности кластера.

Резервирование и отказоустойчивость СУБД​

Резервирование и отказоустойчивость СУБД обеспечивается за счет поддержки работы Console с Patroni-кластером PostgreSQL.

Полнота мониторинга на оперативном дашборде​

Полнота мониторинга на оперативном дашборде обеспечивается за счет передачи бизнес-метрик с Console и Worker, описанных в разделе «Поэтапное обновление GraDeLy».

Наличие включенной ротации сертификатов​

Ротация сертификатов контролируется со стороны администраторов инфраструктуры. Со стороны GraDeLy для работы с секретами используется vault-агент.

Обеспечение минимального набора реплик pod​

GraDeLy автоматически не определяет, сколько реплик понадобится для его работы. Число реплик настраивается при развертывании приложения:

  • для Worker — 1 под на 1 модуль графа;
  • для Console — минимальное количество подов определяется на основе НТ тестов.

Использование liveness probe и readiness probe​

В Console GraDeLy есть стандартные пробы:

  • Readiness probe: uri - Settings/isReady;
  • Liveness probe: uri - Settings/isActive.

Запрет доступа одной АС к БД другой АС​

На сетевом уровне подключение к БД происходит по fqdn с сертификатами из HashiCorp Vault. На логическом уровне по логину и паролю из HashiCorp Vault. Запрет доступа осуществляется за счет контроля доступа команд к секретам друг друга со стороны администраторов.

Требование по георезервированию​

Георезервирование полностью настраивается администраторами инфраструктуры. Со стороны GraDeLy специфических ограничений или требований нет, так как GraDeLy работает со всеми базами по одному принципу.

Обработка ошибок CONSTRAINT VIOLATION​

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

  • Нарушение первичного ключа (PK) — например, вставка уже существующего ключа.
  • Нарушение внешнего ключа (FK) — Попытка вставить ссылку на несуществующий ключ (INSERT, UPDATE) или попытка удалить из родительской таблицы ключ, на который есть ссылка в дочерней таблице.
  • Невыполнение проверок — недопустимое или отсутствующее значение, которое не проходит проверку типа CHECK (включая NOT NULL).

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

Ошибочные транзакции сохраняются в отдельное хранилище — теневую таблицу $CONFLICT_RESOLUTION_TX$.

Для создания таблицы пропишите и выполните в редакторе СУБД SQL-скрипт:

CREATE TABLE <schema_name>."$CONFLICT_RESOLUTION_TX$" (
source_id int8 NOT NULL,
gradely_id int8 NOT NULL,
change_vector_seq int4 NOT NULL,
transaction_id int8 NOT NULL,
table_schema varchar(128) NOT NULL,
table_name varchar(128) NOT NULL,
opcode bpchar(1) NOT NULL,
error_code bpchar(1) NOT NULL,
error_message text NULL,
vector_data jsonb NOT NULL,
error_handle_action bpchar(1) NOT NULL,
process_id uuid NOT NULL,
created_at timestamp NOT NULL,
CONSTRAINT "$CONFLICT_RESOLUTION_TX$_pkey" PRIMARY KEY (source_id, gradely_id, change_vector_seq)
);

Ошибочные транзакции сохраняются в теневой таблице, если в поле options для соединения базы данных параметр save.skipped.tx.to.db: true. По умолчанию save.skipped.tx.to.db: true.

По умолчанию таблица находится в схеме grdl.

Структура технической таблицы:

КолонкаТипОписание
source_idbigintID базы-источника
gradely_idbigintвнутренний ID транзакции GraDeLy
change_vector_seqintпорядковый номер вектора изменений в транзакции
transaction_idbigintID транзакции
table_schemavarchar(128)Схема таблицей-приемником
table_namevarchar(128)Название таблицы-приемника
opcodechar(1)Операция (I, U, D)
error_codechar(1)Тип ошибки (F — нарушение внешнего ключа, P — нарушение первичного ключа, C — нарушение проверки значений)
error_messagevarchar2(256)Человекочитаемое описание ошибки
vector_datajsonJSON с содержимым вектора
error_handle_actionchar(1)Тип реакции на ошибку ( A — ABORT, C — CONTINUE)

Если указаны кластер и топик, GraDeLy может также сохранять ошибочные векторы в Kafka.

Ошибочные транзакции сохраняются в Kafka topic, если в поле options для соединения базы данных параметр save.skipped.tx.to.kafka: true. По умолчанию save.skipped.tx.to.kafka: false.

Структура сообщения об ошибке в Kafka topic:

openapi: 3.1.0
info:
title: Error message for REP020
description: Структура сообщения об ошибке
version: 0.0.1
paths:
/:
description: Фиктивный path
components:
schemas:
LogMessage:
description: JSON, который записывается в поле сообщения
type: object
properties:
gradely_id:
description: Идентификатор транзакции, в которой произошла ошибка
type: integer
format: int64
minimum: 1
maximum: 9223372036854775807
example: 83487
source_id:
description: Источник, из которого пришло изменение
type: integer
format: int64
minimum: -9223372036854775808
maximum: 9223372036854775807
example: 87238507
opcode:
description: Код операции
type: string
enum: [ I, U, D ]
exception:
description: Ошибка, полученная драйвером jdbc
type: object
properties:
sql_state:
description: SQLException.getSQLState()
type: string
maxLength: 5
example: 0A000
vendor_code:
description: SQLException.getErrorCode()
type: integer
minimum: -9223372036854775808
maximum: 9223372036854775807
example: 600
additionalProperties: false
vector:
description: Вектор изменений
type: object
additionalProperties:
description: Описание полей
type: string
maxLength: 128
example: The event description
required: [ gradely_id, source_id, opcode, exception, vector ]
additionalProperties: false

Коэффициент объёма WAL к объёму данных на Kafka​

Цель​

Получить практический коэффициент, который связывает объем WAL, проходящий через логический слот PG, с приростом занятого диска Kafka. Коэффициент используется для расчета требуемого диска Kafka под retention.

примечание

Коэффициент рассчитывается для воркеров потоков репликации.

Определения​

  • WAL_slot_diff — объем WAL в байтах, соответствующий прогрессу логического слота, измеряется как разница confirmed_flush_lsn перед началом теста и после его окончания.

    select confirmed_flush_lsn lsn1 from pg_replication_slots; -- Перед началом теста
    select confirmed_flush_lsn lsn2 from pg_replication_slots; -- После окончания теста

    select lsn2::pg_lsn - lsn1::pg_lsn;
  • kafka_disk_diff — прирост физического объема файлов Kafka на диске для data-топика.

    du -sb /KAFKADATA/demo-* -- Перед началом теста
    du -sb /KAFKADATA/demo-* -- После окончания теста
  • W / K — сколько WAL приходится на 1 байт Kafka на диске.

  • K / W — сколько диска Kafka получается на 1 байт WAL.

Методика измерения​

Значения параметров по умолчанию, влияющие на измерение коэффициентов​

PG:

wal_compression = off

Kafka:

batch.size = 16384
compression.type = zstd
compression.level = 3

Набор тестов и их смысл​

Тесты разбиты на 3 группы:

  1. Best-case — данные максимально сжимаемы.

    • DDL:

      CREATE TABLE compression_test.best_case (
      id bigint PRIMARY KEY GENERATED ALWAYS AS IDENTITY,
      tenant_id smallint NOT NULL,
      payload text NOT NULL
      );
      CREATE OR REPLACE PROCEDURE compression_test.best_case_proc(
      p_target_mb int DEFAULT 1024,
      p_payload_len int DEFAULT 960,
      p_rows_per_tx int DEFAULT 1
      )
      LANGUAGE plpgsql
      AS $$
      DECLARE
      target_bytes bigint := p_target_mb::bigint * 1024 * 1024;
      i bigint := 0;
      j int;
      BEGIN
      LOOP
      EXIT WHEN pg_table_size('compression_test.best_case'::regclass) >= target_bytes;

      FOR j IN 1..p_rows_per_tx LOOP
      i := i + 1;

      INSERT INTO compression_test.best_case(tenant_id, payload)
      VALUES (((i - 1) % 10 + 1)::smallint, repeat('A', p_payload_len));
      END LOOP;

      COMMIT;
      END LOOP;
      END;
      $$;
    • Best-1: мелкие транзакции (1 row/tx), короткий payload:

      CALL compression_test.best_case_proc(
      p_target_mb => 1024,
      p_payload_len => 960,
      p_rows_per_tx => 1
      );
    • Best-2: крупные транзакции (1024 rows/tx), короткий payload

      CALL compression_test.best_case_proc(
      p_target_mb => 1024,
      p_payload_len => 960,
      p_rows_per_tx => 1024
      );
    • Best-3: мелкие транзакции (1 row/tx), очень длинный payload(TOAST)

      CALL compression_test.best_case_proc(
      p_target_mb => 1024,
      p_payload_len => 960000,
      p_rows_per_tx => 1
      );
  2. Worst-case — данные максимально несжимаемы.

    • DDL:

      CREATE TABLE compression_test.worst_case (
      id bigint PRIMARY KEY GENERATED ALWAYS AS IDENTITY,
      tenant_id smallint NOT NULL,
      payload text NOT NULL
      );



      CREATE OR REPLACE PROCEDURE compression_test.worst_case_proc(
      p_target_mb int DEFAULT 1024,
      p_payload_len int DEFAULT 960,
      p_rows_per_tx int DEFAULT 1,
      p_table regclass DEFAULT 'compression_test.worst_case'::regclass
      )
      LANGUAGE plpgsql
      AS $$
      DECLARE
      target_bytes bigint := p_target_mb::bigint * 1024 * 1024;
      i bigint := 0;
      j int;
      tenant smallint;
      BEGIN
      IF p_rows_per_tx < 1 THEN
      RAISE EXCEPTION 'p_rows_per_tx must be >= 1';
      END IF;

      LOOP
      EXIT WHEN pg_table_size(p_table) >= target_bytes;

      FOR j IN 1..p_rows_per_tx LOOP
      i := i + 1;

      tenant := ((i * 7919) % 30000 + 1)::smallint;

      INSERT INTO compression_test.worst_case(tenant_id, payload)
      VALUES (tenant, substr(encode(compression_test.gen_random_bytes((p_payload_len + 1) / 2), 'hex'), 1, p_payload_len));
      END LOOP;

      COMMIT;
      END LOOP;
      END;
      $$;



      CREATE OR REPLACE PROCEDURE compression_test.worst_case_toast_proc(
      p_target_mb int DEFAULT 1024,
      p_payload_len int DEFAULT 960,
      p_rows_per_tx int DEFAULT 1,
      p_table regclass DEFAULT 'compression_test.worst_case'::regclass
      )
      LANGUAGE plpgsql
      AS $$
      DECLARE
      target_bytes bigint := p_target_mb::bigint * 1024 * 1024;
      i bigint := 0;
      j int;
      tenant smallint;
      BEGIN
      IF p_rows_per_tx < 1 THEN
      RAISE EXCEPTION 'p_rows_per_tx must be >= 1';
      END IF;

      LOOP
      EXIT WHEN pg_table_size(p_table) >= target_bytes;

      FOR j IN 1..p_rows_per_tx LOOP
      i := i + 1;

      tenant := ((i * 7919) % 30000 + 1)::smallint;

      INSERT INTO compression_test.worst_case(tenant_id, payload)
      VALUES (
      tenant,
      (
      SELECT substr(
      string_agg(
      encode(compression_test.gen_random_bytes(chunk_bytes), 'hex'),
      '' ORDER BY idx
      ),
      1,
      p_payload_len
      )
      FROM (
      SELECT
      gs AS idx,
      LEAST(1024, need_bytes - gs * 1024) AS chunk_bytes
      FROM (
      SELECT ((p_payload_len + 1) / 2)::int AS need_bytes
      ) nb,
      generate_series(0, (nb.need_bytes - 1) / 1024) gs
      ) s
      )
      );
      END LOOP;

      COMMIT;
      END LOOP;
      END;
      $$;
    • Worst-1: мелкие транзакции (1 row/tx), payload ≈960:

      CALL compression_test.worst_case_proc(
      p_target_mb => 1024,
      p_payload_len => 960,
      p_rows_per_tx => 1
      );
    • Worst-2: крупные транзакции (1024 rows/tx), payload ≈960:

      CALL compression_test.worst_case_proc(
      p_target_mb => 1024,
      p_payload_len => 960,
      p_rows_per_tx => 1024
      );
    • Worst-3: мелкие транзакции (1 row/tx), очень длинный payload (TOAST):

      CALL compression_test.worst_case_toast_proc(
      p_target_mb => 1024,
      p_payload_len => 960000,
      p_rows_per_tx => 1
      );
  3. pgbench — приближение к "типичному OLTP".

    • Тестовые данные (не реплицируются):

      pgbench -i -s 100 -I dtGvp \
      -d "host= port= dbname= user="
    • Запуск теста:

      pgbench -T 600 -c 16 -j 4 -M prepared -P 10 \
      -d "host= port= dbname= user="

Цель набора тестов — получить диапазон C:

  • C_best (нижняя граница);
  • C_worst (верхняя граница);
  • C_typical (реалистичная опора).

Результаты тестирования​

СценарийКоличество сообщенийWAL, MiBKafka demo-0 delta, MiBK / WW / K
Best-1 (A×960, 1 row/tx)1,049,4521179.93259.090.21964.55
Best-2 (A×960, 1024 rows/tx)1,0241131.957.030.0062161.01
Best-3 (A×960000, 1 row/tx)85,9871120.0024.690.022045.35
Worst-1 (random×960, 1 row/tx)1,048,2811184.00836.300.70631.42
Worst-2 (random×960, 1024 rows/tx)1,0241122.86539.400.48042.08
Worst-3 (random×960000, 1 row/tx)1,0781075.74512.980.47692.10
pgbench (TPC-B-like)550,5391661.06190.290.11468.73

Итог и рекомендации по использованию коэффициента​

  1. Рекомендуемый подход — по выше описанной методике измерьте коэффициент на своем стенде/прод-профиле.

  2. Если нет возможности провести тестирование на своем стенде:

    • Практически безопасная стандартная рекомендация для расчета диска Kafka под retention, опираясь на worst-cases — (K/W=0.48) (W/K=2.08).
    • Самостоятельно оцените сжимаемость и размер транзакций и, опираясь на таблицу тестирования, подберите коэффициент.
  3. Для системных топиков (topic-error, topic-conflict-tx) — заложите размер KAFKADATA 100 МБ на каждый топик.

Соотношение размера выгружаемых данных к объему данных в Kafka​

Цель​

Получить cоотношение размера выгружаемых данных к объему данных в Kafka.

Результаты тестирования​

По результатам тестирования ориентировочное соотношение размера выгружаемых данных к объему, занимаемому ими в каталоге /KAFKADATA, составляет:

pg_table_size/KAFKADATA ≈ 2.1

Указанное соотношение получено для случая, когда replication factor = 1. Это означает, что при выгрузке INIT объем данных, занимаемых в Kafka, может отличаться от исходного объема таблиц в PostgreSQL и должен учитываться при расчете требований к дисковому пространству.

Для replication factor = N суммарный объем данных, занимаемый в Kafka-кластере, можно оценить по формуле:

KAFKADATA_total ≈ (pg_table_size / 2.1) × N

При этом эквивалентное соотношение примет вид:

pg_table_size / KAFKADATA_total ≈ 2.1 / N

где: KAFKADATA_total — суммарный объем данных по всем брокерам Kafka;N — значение replication factor.

Итог и рекомендации по использованию коэффициента​

Объем данных в Kafka при выгрузке INIT может отличаться от исходного объема таблиц в PostgreSQL и должен учитываться при расчете требований к дисковому пространству.