База знаний Appliner
Интеграция с Кафкой
Содержимое конфигурации внешней 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 |
|
Таймаут установки соединения и клиентских запросов producer/consumer. |
group_poll_timeout_ms |
|
Таймаут poll для проверки consumer group. |
admin_request_timeout_ms |
|
Таймаут 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 нужно экранировать, если они должны восприниматься буквально.