Skip to main content
Print

Интеграция с Кафкой

Содержимое конфигурации внешней Kafka

Этот документ описывает содержимое поля integration для Kafka-интеграции: параметры клиента, правила отправки сообщений из Appliner во внешнюю Kafka и правила чтения сообщений из внешней Kafka обратно в Appliner.

Минимальная структура

{
  "bootstrap_servers": "platform-kafka-kafka-bootstrap.kafka.svc.cluster.local:9092",
  "security_protocol": "SASL_PLAINTEXT",
  "sasl_mechanism": "SCRAM-SHA-512",
  "sasl_jaas_config": "org.apache.kafka.common.security.scram.ScramLoginModule required username="developer" password="developer";",
  "producer": {
    "client_id": "appliner",
    "topics": [
      {
        "internal_key": "topic-a",
        "internal_name": "Топик 3",
        "topic_name": "external-topic-a"
      }
    ]
  },
  "consumer": {
    "group_id": "appliner-integration-${app.project-name}",
    "client_id": "appliner",
    "topics": [
      {
        "internal_key": "topic-a",
        "internal_name": "Топик 3",
        "topic_pattern": "external-topic-a"
      }
    ]
  }
}

Подключение к Kafka

Поле Обязательность  Описание
bootstrap_servers

Да 

Адреса брокеров внешней Kafka в формате Kafka bootstrap.servers, например kafka-1:9092,kafka-2:9092.
security_protocol

Нет 

Протокол безопасности Kafka-клиента, например PLAINTEXT, SASL_PLAINTEXT, SASL_SSL. Если поле не задано, используется поведение Kafka-клиента по умолчанию.
sasl_mechanism

Нет 

SASL-механизм, например PLAIN, SCRAM-SHA-256, SCRAM-SHA-512.
sasl_jaas_config

Нет 

 JAAS-конфигурация для SASL-аутентификации. Содержит секреты, поэтому храните файл как секрет окружения.

Пример для SCRAM:

{
  "security_protocol": "SASL_PLAINTEXT",
  "sasl_mechanism": "SCRAM-SHA-512",
  "sasl_jaas_config": "org.apache.kafka.common.security.scram.ScramLoginModule required username="user" password="password";"
}

Таймауты

Секция timeouts необязательна. Если её нет, применяются значения по умолчанию.

Поле  Значение по умолчанию  Описание
connect_timeout_ms  
 

20000

Таймаут установки соединения и клиентских запросов producer/consumer.
group_poll_timeout_ms 
 

1000 

Таймаут poll для проверки consumer group.
admin_request_timeout_ms

 5000 

Таймаут AdminClient-запросов, например проверки доступности кластера.

Поддерживаются оба варианта именования: connect_timeout_ms и connect-timeout-ms.

{
  "timeouts": {
    "connect_timeout_ms": 20000,
    "group_poll_timeout_ms": 1000,
    "admin_request_timeout_ms": 5000
  }
}

Отправка сообщений во внешнюю Kafka

Секция producer настраивает producer внешней Kafka. Она не обязательна: если интеграция не должна отправлять сообщения во внешнюю Kafka, секцию можно не указывать.

Поле Обязательность  Описание
producer  Нет  Настройки отправки сообщений во внешнюю Kafka.
producer.client_id Да, если задан producer  client.id Kafka producer.
producer.topics Нет  Список логических топиков, доступных для отправки из Appliner.
producer.topics[].internal_key  Да для маршрутизации по ключу  Внутренний ключ топика в Appliner. Используется, когда исходящее сообщение содержит topic_key.
producer.topics[].internal_name  Нет  Человекочитаемое название топика для интерфейса и логов.
producer.topics[].topic_name  Да  Реальное имя топика во внешней Kafka.

Как выбирается топик для отправки:
Если в исходящем сообщении задан topic_name, сервис отправляет сообщение в этот топик напрямую.
Иначе, если задан topic_key, сервис ищет запись в producer.topics с таким internal_key и берёт её topic_name.
Иначе используется первый topic_name из producer.topics.
Если топик определить нельзя, в Appliner отправляется ошибка.

Чтение сообщений из внешней Kafka

Секция consumer настраивает подписки на внешние топики. Она не обязательна: если интеграция не должна читать сообщения из внешней Kafka, секцию можно не указывать.

Поле  Обязательность  Описание
consumer Нет  Настройки чтения сообщений из внешней Kafka.
consumer.group_id  Да, если задан consumer  Kafka consumer group. Можно использовать плейсхолдеры Spring, например ${app.project-name}.
consumer.client_id  Да, если задан consumer  Базовый client.id consumer. Для каждой подписки сервис добавляет суффикс -<internal_key>.
consumer.topics  Нет  Список подписок на внешние топики. Если список пустой или отсутствует, внешние listener’ы не запускаются.
consumer.topics[].topic_pattern Да  Регулярное выражение для подписки на внешние Kafka-топики. Для точного топика укажите его имя без wildcard-символов.
consumer.topics[].internal_key Да  Внутренний ключ подписки.
consumer.topics[].internal_name  Нет  Человекочитаемое название подписки.
consumer.topics[].key_filter Нет  Список допустимых Kafka record key. Если список пустой, фильтр по key отключён.
consumer.topics[].content_filter  Нет  Список JSONPath-выражений для фильтрации record value. Если список пустой, фильтр по содержимому отключён.

Consumer читает key и value как строки и использует auto.offset.reset=earliest.
Пример с фильтрами:

{
  "consumer": {
    "group_id": "appliner-integration-${app.project-name}",
    "client_id": "appliner",
    "topics": [
      {
        "topic_pattern": "external-topic-.*",
        "internal_key": "orders",
        "internal_name": "Заказы",
        "key_filter": ["business-a", "business-b"],
        "content_filter": ["$.eventType", "$.order[?(@.status == 'PAID')]"]
      }
    ]
  }
}

Логика фильтрации:

  • пустой key_filter не ограничивает сообщения по Kafka key;
  • пустой content_filter не ограничивает сообщения по body;
  • если заданы оба фильтра, сообщение должно пройти оба;
  • key_filter сравнивает Kafka key с элементами списка точным совпадением;
  • content_filter считает сообщение подходящим, если хотя бы одно JSONPath-выражение возвращает непустой результат;
  • если body не является валидным JSON или JSONPath не может быть вычислен, сообщение не проходит content_filter.

Замечания

  • sasl_jaas_config содержит логин и пароль. Не храните реальные значения в репозитории.
  • Дублирующиеся producer.topics[].internal_key игнорируются: используется первая найденная запись.
  • Закомментированное поле correlation-id-path из старых примеров текущей конфигурацией не используется.
  • Для topic_pattern используется Java regular expression, поэтому специальные символы regex нужно экранировать, если они должны восприниматься буквально.

 

Оставить комментарий

Оглавление