Системный журнал EVTP#

Настройка системного журнала#

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

  • $kafka_logs_dir/*-standalonesession-*.log — содержит события Job Manager. Подключение и отключение Task Manager, события смены лидера в кластере, события в работе обработчиков.

  • $kafka_logs_dir/*-taskexecutor-*.log — содержит события Task Manager. События подключения к Job Manager-лидеру, события в работе обработчиков, журналы самих обработчиков.

Логирование в централизованную систему журналирования Platform V Monitor#

Для EVTP может быть настроена отправка системных журналов в централизованную систему журналирования Platform V Monitor.

Предварительно необходимо узнать названия топиков централизованной системы журналирования Platform V Monitor, куда будут отправляться журналы EVTP.

Перед установкой (переустановкой) новой версии EVTP необходимо в конфигурационном файле vars.yml заполнить блок:

logback_kafka_appender: # Настройка отправки логов в kafka
  enable: false # Включение механизма отправки логов в Kafka через logback
  topic_FlinkLogger: <topic_name> # Топик для отправки логов флинка
  topic_FlinkLogger_never_block: true # Блокировка работы приложения при недоступности Kafka (по умолчанию - блокируется, при true недоставленные сообщения в kafka отбрасываются)
  topic_FlinkLogger_discarding_threshold: 20 # Процент свободного места в очереди отправки сообщений при достижении которого будут удаляться сообщения уровня TRACE, DEBUG, INFO
  topic_FlinkLogger_queue_size: 512 # Размер очереди для отправки в Кафку
  topic_UniversalJobLogger: <topic_name> # Топик для отправки логов универсального обработчика
  topic_UniversalJobLogger_never_block: true # Блокировка работы приложения при недоступности Kafka (по умолчанию - блокируется, при true недоставленные сообщения в kafka отбрасываются)
  topic_UniversalJobLogger_discarding_threshold: 20 # Процент свободного места в очереди отправки сообщений при достижении которого будут удаляться сообщения уровня TRACE, DEBUG, INFO
  topic_UniversalJobLogger_queue_size: 512 # Размер очереди для отправки в Кафку
  retries: 3 # Количество переиницилизации продьюсера, если параметр не задан, по умолчанию значение 3
  interval: 1000 # Интервал между переиницилизациями, по умолчанию 1000 мс. Задается в мс.
  multiplier: 1 # Множитель интервала переинициализации, по умолчанию значение 1.
  max_pool_size: 8 # Опциональный параметр для установки максимального количества одновременно работающих продюсеров (размер пула), значение по умолчанию 8
  pool_size: 1 # Опциональный параметр, задающий количество одновременно работающих продюсеров (размер пула), не должен быть больше producerMaxPoolSize, значение по умолчанию 1
  pool_load_balancer: round-robin # Опциональный параметр для указания алгоритма балансировки нагрузки для параллельных продюсеров ('random' или 'round-robin'), значение по умолчанию 'round-robin'
  producer_configs: # Настройка Kafka продюсера
    - "bootstrap.servers=host1:port1,host2:port2" # Bootstrap подключения к Apache Kafka вида host:port,host2:port2
    - "security.protocol=SSL" # Тип протокола подключения. Значение по умолчанию: SSL
#      ssl.keystore.location: ssl/logback.jks # Путь до хранилища сертификатов, для отправки логов в kafka
#      ssl.keystore.password: _PLACEHOLDER_ # Пароль от keystore
#      ssl.key.password: _PLACEHOLDER_ # Пароль для key.password
#      ssl.truststore.location: ssl/logback.jks # Путь до truststore хранилища
#      ssl.truststore.password: _PLACEHOLDER_ # Пароль от truststore
#      ssl.endpoint.identification.algorithm: "" # Обязательный параметр. Значение "" не изменяемое. Отключение проверки хостнейма в сертификате, обязательно для стендов Kafka
#      ssl.engine.factory.class: ru.sbt.ss.kafka.VaultSslEngineFactory # Класс используемый для настроек vault подключения
#      ssl.vault.properties.file: conf/vault.properties # Путь до файла с настройками vault
#      ## переопределяем необходимые параметры
#      ssl.vault.pki.common.name: <common_name> # Common name (CN) для генерации сертификатов
#      ssl.keystore.location: ${vault:ssl.keystore.location} # Путь до keystore хранилища сертификата, где будет размещен сертификат сгенерированный Vault
#      ssl.truststore.location: ${vault:ssl.truststore.location} # Путь до truststore хранилища

Настройки отправляемых сообщений указываются в файле /conf/logback.xml. Отправляемые сообщения имеют формат JSON согласно шаблону, описываемому в файле logback.xml.

Для настройки отправки сообщений необходимо заполнить блок appender FlinkLogger файла logback.xml.

В данном блоке следует обратить внимание на заполнение следующих параметров:

Параметр

Описание

Значение по умолчанию

message

Указать название поля для тела сообщения

message

timestamp

Указать название поля для временной метки

timestamp

level

Уровень логирования сообщения

level

needAddThreadToMessage

Необходимость добавления имени потока в поле с телом сообщения

true

name=host

Имя хоста расположения EVTP

dns-имя или ip

name=system

Тип сервиса (job-manager или task-manager)

job-manager

topic

Топик записи сообщений логов

flink_log

client.id

Идентификатор продьюсера Kafka, по которому можно однозначно идентифицировать продьюсера Apache Kafka

пусто

bootstrap.servers

Хосты подключения к системе журналирования

host:port,host:port

security.protocol

Тип протокола подключения к системе журналирования

SSL

ssl.endpoint.identification.algorithm

Отключение проверки имени хоста в сертификате

пусто

appender-ref

Альтернативное место для сообщений журналов

По умолчанию не указывается

Пример файла FlinkLogger.xml с заполненными настройками и подробным описанием блока appender(s).

При возникновении ошибки валидации сообщений, она записывается в логи в формате: время возникновения ошибки, обработчик, описание ошибки. В описании ошибки указан атрибут, в котором возникла ошибка, и тип ошибки (ошибка типов, длины, и т.д.). Данные логи могут быть отправлены в централизованную систему журналирования, например Platform V Monitor.

Пример ошибок, попадающих в логи:

2023-04-03 09:30:32.576 [transformation-transformation -> route-branch-transformation (1/1)#0] ERROR r.s.c.f.e.process.flow.function.DeclarativeMapperMapFunction  - transformation error in step transformation: FATAL: password authentication failed for user "em"

Доступные уровни логирования системного журнала#

Имя логирования

Описание

Примечание

TRACE

Отражает менее приоритетные события для отладки

Уровень TRACE не рекомендуется включать в Промышленной среде

DEBUG

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

Уровень DEBUG не рекомендуется включать в Промышленной среде

INFO

Уровень логирования по умолчанию. События отражают информационные события системного журнала

WARN

Уровень отражает предупреждения и некритические ошибки обработки запросов/состояния сервиса

ERROR

Уровень отражает критические ошибки обработки запросов/состояния сервиса

Часто встречающиеся события в файлах логов#

Часто встречающиеся события в файлах логов типа ERROR#

  1. Не подошел пароль от хранилища (security.ssl.internal.keystore-password).

2022-06-28 13:29:45.584 [main] ERROR org.apache.flink.runtime.entrypoint.ClusterEntrypoint  - Could not start cluster entrypoint StandaloneSessionClusterEntrypoint.
org.apache.flink.runtime.entrypoint.ClusterEntrypointException: Failed to initialize the cluster entrypoint StandaloneSessionClusterEntrypoint.
        at org.apache.flink.runtime.entrypoint.ClusterEntrypoint.startCluster(ClusterEntrypoint.java:189)
        at org.apache.flink.runtime.entrypoint.ClusterEntrypoint.runClusterEntrypoint(ClusterEntrypoint.java:537)
        at org.apache.flink.runtime.entrypoint.StandaloneSessionClusterEntrypoint.main(StandaloneSessionClusterEntrypoint.java:74)
Caused by: java.io.IOException: Failed to initialize SSL for the blob server
        at org.apache.flink.runtime.blob.BlobServer.<init>(BlobServer.java:183)
        at org.apache.flink.runtime.entrypoint.ClusterEntrypoint.initializeServices(ClusterEntrypoint.java:276)
        at org.apache.flink.runtime.entrypoint.ClusterEntrypoint.runCluster(ClusterEntrypoint.java:209)
        at org.apache.flink.runtime.entrypoint.ClusterEntrypoint.lambda$startCluster$0(ClusterEntrypoint.java:171)
        at java.base/java.security.AccessController.doPrivileged(Native Method)
        at java.base/javax.security.auth.Subject.doAs(Subject.java:423)
        at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1836)
        at org.apache.flink.runtime.security.HadoopSecurityContext.runSecured(HadoopSecurityContext.java:41)
        at org.apache.flink.runtime.entrypoint.ClusterEntrypoint.startCluster(ClusterEntrypoint.java:170)
        ... 2 common frames omitted
Caused by: java.io.IOException: Keystore was tampered with, or password was incorrect
        at java.base/sun.security.provider.JavaKeyStore.engineLoad(JavaKeyStore.java:795)
        at java.base/sun.security.util.KeyStoreDelegator.engineLoad(KeyStoreDelegator.java:243)
        at java.base/java.security.KeyStore.load(KeyStore.java:1479)
        at org.apache.flink.runtime.net.SSLUtils.getKeyManagerFactory(SSLUtils.java:277)
        at org.apache.flink.runtime.net.SSLUtils.createInternalNettySSLContext(SSLUtils.java:333)
        at org.apache.flink.runtime.net.SSLUtils.createInternalSSLContext(SSLUtils.java:299)
        at org.apache.flink.runtime.net.SSLUtils.createSSLServerSocketFactory(SSLUtils.java:103)
        at org.apache.flink.runtime.blob.BlobServer.<init>(BlobServer.java:180)
        ... 10 common frames omitted
Caused by: java.security.UnrecoverableKeyException: Password verification failed
        at java.base/sun.security.provider.JavaKeyStore.engineLoad(JavaKeyStore.java:793)
        ... 17 common frames omitted
  1. Не найдено хранилище (security.ssl.internal.keystore).

2022-06-28 13:30:56.224 [main] ERROR org.apache.flink.runtime.entrypoint.ClusterEntrypoint  - Could not start cluster entrypoint StandaloneSessionClusterEntrypoint.
org.apache.flink.runtime.entrypoint.ClusterEntrypointException: Failed to initialize the cluster entrypoint StandaloneSessionClusterEntrypoint.
        at org.apache.flink.runtime.entrypoint.ClusterEntrypoint.startCluster(ClusterEntrypoint.java:189)
        at org.apache.flink.runtime.entrypoint.ClusterEntrypoint.runClusterEntrypoint(ClusterEntrypoint.java:537)
        at org.apache.flink.runtime.entrypoint.StandaloneSessionClusterEntrypoint.main(StandaloneSessionClusterEntrypoint.java:74)
Caused by: java.io.IOException: Failed to initialize SSL for the blob server
        at org.apache.flink.runtime.blob.BlobServer.<init>(BlobServer.java:183)
        at org.apache.flink.runtime.entrypoint.ClusterEntrypoint.initializeServices(ClusterEntrypoint.java:276)
        at org.apache.flink.runtime.entrypoint.ClusterEntrypoint.runCluster(ClusterEntrypoint.java:209)
        at org.apache.flink.runtime.entrypoint.ClusterEntrypoint.lambda$startCluster$0(ClusterEntrypoint.java:171)
        at java.base/java.security.AccessController.doPrivileged(Native Method)
        at java.base/javax.security.auth.Subject.doAs(Subject.java:423)
        at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1836)
        at org.apache.flink.runtime.security.HadoopSecurityContext.runSecured(HadoopSecurityContext.java:41)
        at org.apache.flink.runtime.entrypoint.ClusterEntrypoint.startCluster(ClusterEntrypoint.java:170)
        ... 2 common frames omitted
Caused by: java.nio.file.NoSuchFileException: /opt/Apache/flink/ssl/flink.jks
        at java.base/sun.nio.fs.UnixException.translateToIOException(UnixException.java:92)
        at java.base/sun.nio.fs.UnixException.rethrowAsIOException(UnixException.java:111)
        at java.base/sun.nio.fs.UnixException.rethrowAsIOException(UnixException.java:116)
        at java.base/sun.nio.fs.UnixFileSystemProvider.newByteChannel(UnixFileSystemProvider.java:219)
        at java.base/java.nio.file.Files.newByteChannel(Files.java:371)
        at java.base/java.nio.file.Files.newByteChannel(Files.java:422)
        at java.base/java.nio.file.spi.FileSystemProvider.newInputStream(FileSystemProvider.java:420)
        at java.base/java.nio.file.Files.newInputStream(Files.java:156)
        at org.apache.flink.runtime.net.SSLUtils.getKeyManagerFactory(SSLUtils.java:276)
        at org.apache.flink.runtime.net.SSLUtils.createInternalNettySSLContext(SSLUtils.java:333)
        at org.apache.flink.runtime.net.SSLUtils.createInternalSSLContext(SSLUtils.java:299)
        at org.apache.flink.runtime.net.SSLUtils.createSSLServerSocketFactory(SSLUtils.java:103)
        at org.apache.flink.runtime.blob.BlobServer.<init>(BlobServer.java:180)
        ... 10 common frames omitted
  1. Не смогли расшифровать пароль ключом из параметра security.encoding.key.

2022-06-28 13:31:55.644 [main] ERROR ru.sbt.ss.password.BaseEncryptor  - Can't decrypt by the set parameters
  1. Не смогли подключиться к Zookeeper.

2022-06-28 13:34:04.173 [main-SendThread(tkleq-snaps0003.vm.esrt.cloud.ru:2181)] INFO  o.a.flink.shaded.zookeeper.org.apache.zookeeper.ClientCnxn  - Opening socket connection to server
2022-06-28 13:34:04.173 [main-EventThread] ERROR o.a.flink.shaded.curator.org.apache.curator.ConnectionState  - Authentication failed
  1. Не смогли поднять REST интерфейс (порт rest.bind-port занят).

2022-06-28 13:37:42.014 [main] ERROR org.apache.flink.runtime.entrypoint.ClusterEntrypoint  - Could not start cluster entrypoint StandaloneSessionClusterEntrypoint.
org.apache.flink.runtime.entrypoint.ClusterEntrypointException: Failed to initialize the cluster entrypoint StandaloneSessionClusterEntrypoint.
        at org.apache.flink.runtime.entrypoint.ClusterEntrypoint.startCluster(ClusterEntrypoint.java:189)
        at org.apache.flink.runtime.entrypoint.ClusterEntrypoint.runClusterEntrypoint(ClusterEntrypoint.java:537)
        at org.apache.flink.runtime.entrypoint.StandaloneSessionClusterEntrypoint.main(StandaloneSessionClusterEntrypoint.java:74)
Caused by: org.apache.flink.util.FlinkException: Could not create the DispatcherResourceManagerComponent.
        at org.apache.flink.runtime.entrypoint.component.DefaultDispatcherResourceManagerComponentFactory.create(DefaultDispatcherResourceManagerComponentFactory.java:261)
        at org.apache.flink.runtime.entrypoint.ClusterEntrypoint.runCluster(ClusterEntrypoint.java:223)
        at org.apache.flink.runtime.entrypoint.ClusterEntrypoint.lambda$startCluster$0(ClusterEntrypoint.java:171)
        at java.base/java.security.AccessController.doPrivileged(Native Method)
        at java.base/javax.security.auth.Subject.doAs(Subject.java:423)
        at org.apache.hadoop.security.UserGroupInformation.doAs(UserGroupInformation.java:1836)
        at org.apache.flink.runtime.security.HadoopSecurityContext.runSecured(HadoopSecurityContext.java:41)
        at org.apache.flink.runtime.entrypoint.ClusterEntrypoint.startCluster(ClusterEntrypoint.java:170)
        ... 2 common frames omitted
Caused by: java.net.BindException: Could not start rest endpoint on any port in port range 8081
        at org.apache.flink.runtime.rest.RestServerEndpoint.start(RestServerEndpoint.java:219)
        at org.apache.flink.runtime.entrypoint.component.DefaultDispatcherResourceManagerComponentFactory.create(DefaultDispatcherResourceManagerComponentFactory.java:165)
        ... 9 common frames omitted

Возможно использование уровней логирования trace и debug, включение не рекомендуется в промышленных инсталляциях.