Плагин kafka-logger

Плагин kafka-logger используется для отправки логов в кластеры Apache Kafka. Логи можно передавать в формате JSON.

Плагин работает как Kafka-клиент для модуля ngx_lua в Nginx.

Получение логов может занимать некоторое время. Они будут автоматически отправлены после срабатывания таймера batch processor.

Параметр Обязательный По умолчанию Допустимые значения Описание

broker_list

Да

Устаревший параметр
Использовать brokers
Список брокеров Kafka

brokers

Да

Список брокеров Kafka

brokers.host*

Да

Адрес брокера Kafka

brokers.port

Да

[0, 65535]

Порт брокера Kafka

brokers.sasl_config

Нет

Параметры SASL-аутентификации брокера Kafka

brokers.sasl_config.mechanism

Нет

PLAIN

["PLAIN",
"SCRAM-SHA-256",
"SCRAM-SHA-512"]

Механизм SASL-аутентификации

brokers.sasl_config.user

Да

Имя пользователя для SASL-аутентификации
Обязателен при использовании sasl_config

brokers.sasl_config.password

Да

Пароль для SASL-аутентификации
Обязателен при использовании sasl_config

kafka_topic

Да

Целевой топик Kafka для отправки логов

producer_type

Нет

async

["async", "sync"]

Режим отправки сообщений
Асинхронный или синхронный

required_acks

Нет

1

[1, -1]

Количество подтверждений от лидера брокера для завершения отправки сообщения
Определяет надежность доставки
Значение 0 не поддерживается

key

Нет

Ключ для распределения сообщений по партициям

timeout

Нет

3

[1,…​]

Таймаут отправки данных upstream в секундах

name

Нет

kafka logger

Уникальный идентификатор batch-процессора
Используется в метриках Prometheus как apisix_batch_process_entries

meta_format

Нет

default

["default", "origin"]

Формат сбора информации о запросе
default — JSON
origin — исходный HTTP-запрос

log_format

Нет

Формат логов в виде JSON с парами
ключ-значение. Значения поддерживают
строки и вложенные объекты
(до 5 уровней вложенности,
более глубокие уровни обрезаются).
В строках допускается использование
переменных NGINX с префиксом $

include_req_body

Нет

false

[false, true]

При значении true
в лог добавляется тело запроса

include_req_body_expr

Нет

Условие для логирования тела запроса
при включённом include_req_body.
Тело запроса записывается только
если выражение возвращает true

max_req_body_bytes

Нет

524288

>=1

Максимальный размер тела запроса для
логирования. При превышении значение
обрезается

include_resp_body

Нет

false

[false, true]

Включение тела ответа в лог

include_resp_body_expr

Нет

При значении true
в лог добавляется тело ответа

max_resp_body_bytes

Нет

524288

>=1

Максимальный размер тела ответа для
логирования. При превышении значение
обрезается

cluster_name

Нет

1

[0,…​]

Идентификатор кластера Kafka
Используется при работе с несколькими кластерами
Доступен только при producer_type=async

producer_batch_num

Нет

200

[1,…​]

Параметр batch_num библиотеки lua-resty-kafka
Количество сообщений в одном батче

producer_batch_size

Нет

1048576

[0,…​]

Параметр batch_size библиотеки lua-resty-kafka
Размер батча в байтах

producer_max_buffering

Нет

50000

[1,…​]

Параметр max_buffering библиотеки lua-resty-kafka
Максимальный размер буфера
Количество сообщений

producer_time_linger

Нет

1

[1,…​]

Параметр flush_time библиотеки lua-resty-kafka
Интервал отправки батча в секундах

meta_refresh_interval

Нет

30

[1,…​]

Параметр refresh_interval библиотеки lua-resty-kafka
Интервал обновления метаданных в секундах

Плагин поддерживает batch processor для агрегации и пакетной обработки данных. Это позволяет уменьшить частоту отправки.

Batch processor отправляет данные:

  • каждые 5 секунд

  • или при накоплении 1000 записей

Данные сначала записываются в буфер.

Когда размер буфера превышает параметры batch_max_size или buffer_duration, данные отправляются в Kafka, а буфер очищается.

Если операция успешна — возвращается true. При ошибке возвращается nil и строка "buffer overflow".

Пример meta_format:

  • default

{
  "upstream": "127.0.0.1:1980",
  "start_time": 1619414294760,
  "client_ip": "127.0.0.1",
  "service_id": "",
  "route_id": "1",
  "request": {
    "querystring": {
      "ab": "cd"
    },
    "size": 90,
    "uri": "/hello?ab=cd",
    "url": "http://localhost:1984/hello?ab=cd",
    "headers": {
      "host": "localhost",
      "content-length": "6",
      "connection": "close"
    },
    "body": "abcdef",
    "method": "GET"
  },
  "response": {
    "headers": {
      "connection": "close",
      "content-type": "text/plain; charset=utf-8",
      "date": "Mon, 26 Apr 2021 05:18:14 GMT",
      "server": "APISIX/2.5",
      "transfer-encoding": "chunked"
    },
    "size": 190,
    "status": 200
  },
  "server": {
    "hostname": "localhost",
    "version": "2.5"
  },
  "latency": 0
}
  • origin

GET /hello?ab=cd HTTP/1.1
host: localhost
content-length: 6
connection: close

abcdef

Метаданные

Формат логов может быть задан на уровне метаданных плагина.

Имя Тип Обязательный Описание

log_format

object

Нет

Формат логов в виде JSON

max_pending_entries

integer

Нет

Максимальное количество записей в буфере batch processor

Настройка метаданных является глобальной и применяется ко всем маршрутам и сервисам, где используется плагин kafka-logger.

admin_key можно получить из config.yaml:

admin_key=$(yq '.deployment.admin.admin_key[0].key' conf/config.yaml | sed 's/"//g')
curl http://127.0.0.1:9180/apisix/admin/plugin_metadata/kafka-logger -H "X-API-KEY: $admin_key" -X PUT -d '
{
    "log_format": {
        "host": "$host",
        "@timestamp": "$time_iso8601",
        "client_ip": "$remote_addr",
        "request": { "method": "$request_method", "uri": "$request_uri" },
        "response": { "status": "$status" }
    }
}'

С этой конфигурацией логи будут иметь следующий формат:

{"host":"localhost","@timestamp":"2020-09-23T19:05:05-04:00","client_ip":"127.0.0.1","request":{"method":"GET","uri":"/hello"},"response":{"status":200},"route_id":"1"}
{"host":"localhost","@timestamp":"2020-09-23T19:05:05-04:00","client_ip":"127.0.0.1","request":{"method":"GET","uri":"/hello"},"response":{"status":200},"route_id":"1"}

Включение плагина

Включение плагина выполняется в веб-интерфейсе или в Admin API.

Включение в веб-интерфейсе

Перейти в раздел Маршруты → Плагины. Плагин, добавленный в этом разделе, применяется только к выбранному маршруту.

Либо перейти в раздел Конфигурация плагинов и выполнить настройку плагина. В этом случае создаётся конфигурация с ID, которую можно подключать к разным маршрутам и сервисам для переиспользования.

log kafka 1

В строке поиска ввести kafka-logger и выбрать плагин. Нажать Add.

log kafka 2

Откроется форма конфигурации плагина. Необходимо задать параметры подключения к Syslog-серверу:

  • host — адрес брокера

  • port — порт брокера

  • kafka_topic — имя топика, в который будут отправляться логи

Необязательные параметры:

  • key — ключ сообщения

  • batch_max_size — максимальное количество сообщений в батче перед отправкой

  • name — имя логгера

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

Включение в Admin API

Ниже приведён пример включения плагина для конкретного маршрута:

curl http://127.0.0.1:9180/apisix/admin/routes/5 -H "X-API-KEY: $admin_key" -X PUT -d '
{
    "plugins": {
       "kafka-logger": {
           "brokers" : [
             {
               "host" :"127.0.0.1",
               "port" : 9092
             }
            ],
           "kafka_topic" : "test2",
           "key" : "key1",
           "batch_max_size": 1,
           "name": "kafka logger"
       }
    },
    "upstream": {
       "nodes": {
           "127.0.0.1:1980": 1
       },
       "type": "roundrobin"
    },
    "uri": "/hello"
}'

Плагин поддерживает отправку в несколько брокеров:

"brokers" : [
  {
    "host" :"127.0.0.1",
    "port" : 9092
  },
  {
    "host" :"127.0.0.1",
    "port" : 9093
  }
]

Пример использования

После включения плагина каждый обработанный запрос отправляется на Syslog-сервер.

curl -i http://127.0.0.1:9080/hello

Удаление плагина

Чтобы удалить плагин syslog, нужно удалить соответствующую JSON-конфигурацию из конфигурации плагина.

curl http://127.0.0.1:9180/apisix/admin/routes/1  -H "X-API-KEY: $admin_key" -X PUT -d '
{
    "methods": ["GET"],
    "uri": "/hello",
    "plugins": {},
    "upstream": {
        "type": "roundrobin",
        "nodes": {
            "127.0.0.1:1980": 1
        }
    }
}'