Хранилище окон агрегации#
Варианты распределенного хранилища#
Kafka, compacted topic.
PostgreSQL.
Apache Cassandra.
Apache Ignite.
Reactive Stream Adapter (EVTA).
Особенности реализации хранилища для Kafka compacted topic#
Хранилище разделено на локальный кеш в памяти и в компактном топике Kafka. Каждый экземпляр обработчика подключается к одному компактному топику Kafka и вычитывает сообщения из всех его партиций.
При инициализации хранилища происходит наполнение локального кеша из содержимого топика. Далее локальный кеш синхронизируется с топиком в фоновой задаче с определенным интервалом.
В ключ сообщения Kafka записывается ключ агрегации, имя шага агрегации и время старта окна агрегации.
В тело сообщения Kafka записываются накопленные сообщения, время закрытия окна по тайм-ауту, состояние триггера и признак удаленной записи.
Сообщения хранятся в формате Protobuf, партиция при добавлении записи выбирается на основе ключа агрегации и имени шага агрегации.
Рекомендации по созданию топика Kafka#
При создании топика рекомендуется указать следующие настройки:
Название |
Описание |
Значение |
|---|---|---|
Name |
Имя топика |
- |
Partition Count |
Количество партиций |
10 |
Replica Count |
Фактор репликации |
2 |
cleanup.policy |
Стратегия очистки данных |
compact |
retention.ms |
Время, через которое очищается топик (мс) |
1800000 |
Особенности реализации хранилища для PostgreSQL#
Пример создания используемой таблицы:
create table EVPC_STORAGE
(
ID serial
constraint storages_pk
primary key,
AGGREGATION_KEY text,
TIMESTAMP bigint,
STEP text,
JOB text,
ENTRY bytea,
CLOSED_WINDOW boolean default false
);
Дополнительных настроек не требуется.
Особенности реализации хранилища для Apache Cassandra#
Пример создания используемой таблицы для хранения событий внутри окон агрегации:
CREATE TABLE <keyspace>.EVPC_STORAGE (
ADDITION_TIME timestamp,
AGGREGATION_KEY text,
STEP text,
JOB text,
ENTRY blob,
PRIMARY KEY ((AGGREGATION_KEY, STEP, JOB), ADDITION_TIME)
) WITH CLUSTERING ORDER BY (ADDITION_TIME ASC);
Пример создания используемой таблицы для хранения времени открытия окон агрегации:
CREATE TABLE <keyspace>.EVPC_WINDOW_STORAGE (
AGGREGATION_KEY text PRIMARY KEY,
TIMESTAMP bigint
);
Пример конфигурации универсального обработчика с хранилищем окон агрегации в Apache Cassandra:
{
kafka: {
type: "kafka"
"bootstrap.servers":"localhost:9092"
consumer: {
"group.id": "event-process-flow-group"
"client.id": "event-process-flow-client"
}
producer: {
"client.id": "event-process-flow-producer"
}
}
source: {
name: "source"
type: "source"
topic: "input"
config: ${kafka}
destination: ${aggregationStep}
}
destination: {
name: "output"
type: "destination"
topic: "output"
config: ${kafka}
}
aggregationStep: {
name: "simple-aggregation"
type: "aggregation"
format: {
input: "json"
output: "json"
}
key {
type: "dsl"
source: "file"
path: "testDsl/key.tr"
}
trigger: {
type: "dsl"
source: "file"
path: "testDsl/trigger.tr"
timeout: 3600000000
}
dsl: {
source: "file"
path: "testDsl/simple.tr"
}
storage: {
type: "timeseries"
database: {
keyspace: "example" // обязательное поле
tableName: EVPC_STORAGE,
windowTableName: EVPC_WINDOW_STORAGE,
contactPoints: "{ IP_ADDRESS }"
username: "<example>",
password: "<password>",
localDataCenter: "datacenter1",
ssl.enabled: "false"
driverConfigPath: /path/to/custom/cassandra/driver/config
}
}
destination: ${destination}
}
flow: {
name: "EVPC"
source: [${source}]
}
}
Настройка driverConfigPath позволяет указать путь к файлу application.conf, где переопределяются настройки Datastax Java Driver, используемого для подключения к Apache Cassandra.
Пример заполнения application.conf с переопределенными настройками драйвера:
datastax-java-driver {
advanced.protocol.version = V4
profiles {
slow {
basic.request.timeout = 10 seconds
}
}
}
Расширенные параметры настроек приведены на странице DataStax Java Driver.
Особенности реализации хранилища для Apache Ignite#
Для использования key-value хранилища окон агрегации Apache Ignite необходимо добавить jar-файл ignite-core-<ignite version>.jar, содержащий JDBC драйвер для подключения к Ignite.
Apache Ignite позволяет создавать кеш и эквивалентную ему таблицу с помощью SQL-запроса.
Пример создания используемой таблицы:
CREATE TABLE IF NOT EXISTS EVPC_STORAGE
(
AGGREGATION_KEY VARCHAR NOT NULL,
TIMESTAMP BIGINT NOT NULL,
STEP VARCHAR NOT NULL,
JOB VARCHAR NOT NULL,
ENTRY BINARY NOT NULL,
CLOSED_WINDOW BOOLEAN DEFAULT FALSE,
PRIMARY KEY (AGGREGATION_KEY, TIMESTAMP, STEP, JOB)
) WITH "CACHE_NAME=EVPC_STORAGE, KEY_TYPE=String";
При создании таблицы EVPC_STORAGE для хранения окон агрегации также создается эквивалентный ей кеш, где ключом является строка, содержащая поля AGGREGATION_KEY, TIMESTAMP, STEP, JOB, а значением является поле ENTRY в бинарном формате.
Поле CLOSED_WINDOW — текущее состояние окна:
false — окно агрегации еще открыто;
true — окно агрегации закрыто.
Дополнительных настроек не требуется.
Особенности реализации хранилища агрегации с использованием Reactive Stream Adapter (EVTA)#
Используется интерфейс GRPC для доступа до топика Kafka.
Подробнее о настройках хранилища агрегации с использованием EVTA описано в документе «Руководство администратора» в разделе «Конфигурация универсального обработчика».