Точки назначения событий#
Описание точки назначения событий#
Назначение описывается аналогично источника, поле type имеет значение destination:
name— название источника событий;topic— имя топика для записи сообщений по умолчанию, опциональный, по умолчанию «output»;config— конфигурация транспорта назначения;deadLetterTopic- имя топика/очереди, куда будут отправляться сообщения с ошибочным destination.
Поле config содержит специфическую конфигурацию транспорта:
Apache Kafka: аналогично источнику
Простая запись в файл (запись выполняется во время чекпоинтов и при завершении работы обработчика, рекомендуется для использования с конечными источниками событий (в данный момент только «file»)):
type— тип транспорта, значение: «file» или «file-stream», аналогично источнику;path— полный путь до директории назначения;partPrefix— префикс имен файлов;partSuffix— суффикс имен файлов.
Запись в файл с настраиваемой rollingPolicy (рекомендуется для использования со всеми бесконечными источниками событий (все, кроме «file»)):
type— тип транспорта, значение: «file-stream»;path— полный путь до файла или директории с файлами;partPrefix— префикс имен файлов;partSuffix— суффикс имен файлов;rollingPolicy— политика записи событий в файл.
Поле rollingPolicy содержит конфигурацию политики записи событий в файл:
type— тип, значение «onCheckpoint» — события сохраняются в файл во время чекпоинтов (новый чекпоинт — новый файл) или «default» — политика настраивается с помощью параметров ниже;rolloverInterval— интервал в мс, по истечению которого текущий файл будет сохранен;maxPartSize— максимальный размер файла в байтах;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.
type- тип транспорта, значение: «artemis_mq».connectors- список брокеров для подключения, в формате host:port через запятую.connectorConfigs- SSL-конфигурация брокера для подключения (в случае использования HashiCorp Vault здесь также указываются настройки подключения к HashiCorp Vault):sslEnabledвключение/отключение SSL, по умолчаниюfalse(в случае, если указаны настройки HashiCorp Vault, данную настройку можно не указывать, будет использоваться значениеtrue);keyStorePathпуть до хранилища приватного ключа;keyStorePasswordпароль от хранилища приватного ключа;trustStorePathпуть до хранилища доверенных сертификатов;trustStorePasswordпароль от хранилища доверенных сертификатов;verifyHostвключение/отключение проверки имени хоста клиента и CN сертификата клиента, по умолчаниюtrue.
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 мс.
disableMessageTimestamp- отключение временных меток сообщений, по умолчаниюfalse.deliveryMode- гарантия доставки сообщений. Поле может бытьpersistentилиnon_persistent. По умолчаниюnon_persistent. Опцияpersistent- сообщение записывается в стабильное хранилище как часть операции отправки сообщения клиентом. Опцияnon_persistent- сообщение может быть утеряно, но оно не должно быть доставлено дважды.messagePriority- устанавливает приоритет сообщений в диапазоне от 0 до 9, по умолчанию4.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
}
}