Хранилище окон агрегации#

Варианты распределенного хранилища#

  1. Kafka, compacted topic.

  2. PostgreSQL.

  3. Apache Cassandra.

  4. Apache Ignite.

  5. 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 описано в документе «Руководство администратора» в разделе «Конфигурация универсального обработчика».