Точки назначения событий#

Описание точки назначения событий#

Назначение описывается аналогично источника, поле type имеет значение destination:

  1. name — название источника событий;

  2. topic — имя топика для записи сообщений по умолчанию, опциональный, по умолчанию «output»;

  3. config — конфигурация транспорта назначения;

  4. deadLetterTopic - имя топика/очереди, куда будут отправляться сообщения с ошибочным destination.

Поле config содержит специфическую конфигурацию транспорта:

Apache Kafka: аналогично источнику

Простая запись в файл (запись выполняется во время чекпоинтов и при завершении работы обработчика, рекомендуется для использования с конечными источниками событий (в данный момент только «file»)):

  1. type — тип транспорта, значение: «file» или «file-stream», аналогично источнику;

  2. path — полный путь до директории назначения;

  3. partPrefix — префикс имен файлов;

  4. partSuffix — суффикс имен файлов.

Запись в файл с настраиваемой rollingPolicy (рекомендуется для использования со всеми бесконечными источниками событий (все, кроме «file»)):

  1. type — тип транспорта, значение: «file-stream»;

  2. path — полный путь до файла или директории с файлами;

  3. partPrefix — префикс имен файлов;

  4. partSuffix — суффикс имен файлов;

  5. rollingPolicy — политика записи событий в файл.

Поле rollingPolicy содержит конфигурацию политики записи событий в файл:

  1. type — тип, значение «onCheckpoint» — события сохраняются в файл во время чекпоинтов (новый чекпоинт — новый файл) или «default» — политика настраивается с помощью параметров ниже;

  2. rolloverInterval — интервал в мс, по истечению которого текущий файл будет сохранен;

  3. maxPartSize — максимальный размер файла в байтах;

  4. inactivityInterval — максимальный интервал не активности файла в мс, по истечению которого файл будет сохранен.

Конечные имена файлов будут иметь вид:

<partPrefix>-part-<subtaskIndex>-<partFileIndex>-<partSuffix>

Точка назначения событий#

destination: {
   name: "kafka"
   type: "destination"
   topic: "output"
   config: {
      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"
      }
   }
}
destination: {
   name: "file"
   type: "destination"
   topic: "output"
   config: {
      type: "file"
      path: "/full/path/to/directory/"
      partPrefix: "output"
      partSuffix: ".txt"
        }
    }
destination: {
   name: "file-stream"
   type: "destination"
   topic: "output"
   config: {
      type: "file-stream"
      path: "/full/path/to/directory/"
      partPrefix: "output"
      partSuffix: ".txt"
      rollingPolicy: {
         type: "default"
         rolloverInterval: "60000"
         maxPartSize: "1048576"
         inactivityInterval: "60000"
        }
      }
    }

Описание точки назначения событий с подключением к Active MQ Artemis#

Адрес и имя очереди, в которую необходимо писать сообщения, указываются в поле topic конфигурации источника событий через ::. Например, необходимо писать в очередь queue1 по адресу test-input, тогда в поле topic необходимо указать test-input::queue1. В случае, если имя очереди совпадает с адресом, то в данном поле можно указать только адрес. Например, необходимо писать в очередь test с адресом test, тогда в поле topic необходимо указать test.

  1. type - тип транспорта, значение: «artemis_mq».

  2. connectors - список брокеров для подключения, в формате host:port через запятую.

  3. connectorConfigs - SSL-конфигурация брокера для подключения (в случае использования HashiCorp Vault здесь также указываются настройки подключения к HashiCorp Vault):

    • sslEnabled включение/отключение SSL, по умолчанию false (в случае, если указаны настройки HashiCorp Vault, данную настройку можно не указывать, будет использоваться значение true);

    • keyStorePath путь до хранилища приватного ключа;

    • keyStorePassword пароль от хранилища приватного ключа;

    • trustStorePath путь до хранилища доверенных сертификатов;

    • trustStorePassword пароль от хранилища доверенных сертификатов;

    • verifyHost включение/отключение проверки имени хоста клиента и CN сертификата клиента, по умолчанию true.

  4. factory параметры к org.apache.activemq.artemis.jms.client.ActiveMQConnectionFactory:

    • blockOnqueue синхронная отправка сообщений в очередь, по умолчанию false;

    • loadBalancingPolicyClassName имя балансировщика нагрузки, по умолчанию org.apache.activemq.artemis.api.core.client.loadbalance.RoundRobinConnectionLoadBalancingPolicy (первым подключением выбирается случайное из списка). В дистрибутиве присутствует альтернативный балансировщик ru.sbt.flink.streaming.connectors.jms.artemismq.utils.RoundRobinLoadBalancingPolicy, который всегда начинает с первого подключения;

    • reconnectAttempts количество повторных попыток подключения, по умолчанию 0;

    • recoveryInterval интервал между попытками восстановления соединения, по умолчанию 2000 мс;

    • recoveryIntervalMultiplier множитель, на который будет умножаться каждый следующий интервал между подключениями, по умолчанию 1;

    • confirmationWindowSize размер буфера команд, переданных от клиента серверу, по умолчанию -1, не вести буфер;

    • callTimeout время при отправке через кластерное соединение, по умолчанию 30000 мс.

  5. disableMessageTimestamp - отключение временных меток сообщений, по умолчанию false.

  6. deliveryMode - гарантия доставки сообщений. Поле может быть persistent или non_persistent. По умолчанию non_persistent. Опция persistent - сообщение записывается в стабильное хранилище как часть операции отправки сообщения клиентом. Опцияnon_persistent - сообщение может быть утеряно, но оно не должно быть доставлено дважды.

  7. messagePriority - устанавливает приоритет сообщений в диапазоне от 0 до 9, по умолчанию 4.

  8. messageTimeToLive - устанавливает период времени с момента отправки (в мс), в течение которого созданное сообщение должно храниться в очереди, по умолчанию 0 (время хранения не ограничено).

Точка назначения событий#

  destination: {
   name: "Active MQ Artemis sink"
   type: "destination"
   topic: "flink-flow-test-output"
   config: {
      type: "artemis_mq"
      secret: "secret"
      connectors: "host:port"//, localhost:61616",
      connectorConfigs: {
         sslEnabled: "true"
         keyStorePath: "keystore-path"
         keyStorePassword: "__PLACEHOLDER_"
         trustStorePath: "truststore-path"
         trustStorePassword: "__PLACEHOLDER_"
         verifyHost: "true"
      }
      factory: {
         reconnectAttempts: 1
         recoveryInterval: 3000
         recoveryIntervalMultiplier: 1010
         blockOnQueue: "true"
         confirmationWindowSize: "1"
         loadBalancingPolicyClassName: ru.sbt.flink.streaming.connectors.jms.artemismq.utils.RoundRobinLoadBalancingPolicy
         callTimeout: 3000
      }
      disableMassageTimeStamp: false
      deliveryMode: non_persistent
      messagePriority: 5
      messageTimeToLive: 3600000
   }
}