Apache APISIX kafka-logger 插件实战:Kafka 访问日志采集的配置、格式定制与源码解析
【免费下载链接】apisixThe Cloud-Native API Gateway项目地址: https://gitcode.com/GitHub_Trending/ap/apisix
Apache APISIX 的kafka-logger插件用于将网关处理过的请求日志以 JSON 对象的形式推送到 Apache Kafka 集群,是日志采集与流式分析链路中的标准组件。本文以官方文档为主线,结合当前仓库的插件源码、日志工具模块与测试用例,完整讲解其全部配置属性、批处理机制、日志格式定制方法以及底层实现原理,帮助你从零开始在 APISIX 上接入 Kafka 日志管道。
插件概述与工作原理
kafka-logger插件的定位是 APISIX 的 Kafka 客户端驱动,它基于 ngx_lua 的 Nginx 模块工作:在每个请求的生命周期内收集访问信息,在日志阶段(log阶段)将日志数据写入缓冲区,再由批处理器或直接发送到 Kafka 集群。由于日志发送是异步完成的,因此收到日志数据可能需要一些时间——数据会在批处理器中的定时器到期后自动发送。
从源码看,插件本身只是一个轻量封装,真正的日志收集、格式组装、请求体/响应体处理都复用了 apisix/utils/log-util.lua 工具模块,而底层的 Kafka 生产逻辑则委托给resty.kafka.producer(即 lua-resty-kafka 库):
local producer = require ("resty.kafka.producer") local bp_manager_mod = require("apisix.utils.batch-processor-manager") local log_util = require("apisix.utils.log-util")插件的priority为403、version为0.1,其生命周期钩子分别负责不同职责(见 apisix/plugins/kafka-logger.lua):
access阶段:当开启include_req_body时读取请求体,支持用include_req_body_expr表达式做条件过滤;body_filter阶段:调用log_util.collect_body按需收集响应体(含解压逻辑);log阶段:根据meta_format选择日志条目格式,先尝试交给批处理器,若批处理器不存在则创建新的处理器并异步发送。
配置属性全解析
下表完整列出插件的全部属性(与官方文档一致,并补充了源码中的取值范围与默认值说明):
| 名称 | 类型 | 必填 | 默认值 | 合法值 | 描述 |
|---|---|---|---|---|---|
| broker_list | object | 是 | 已废弃,请改用brokers。Kafka 节点列表,键为主机名、值为端口 | ||
| brokers | array | 是 | Kafka 节点(broker)列表 | ||
| brokers.host | string | 是 | Kafka broker 主机,例如192.168.1.1 | ||
| brokers.port | integer | 是 | [0, 65535] | Kafka broker 端口 | |
| brokers.sasl_config | object | 否 | Kafka broker 的 SASL 认证配置 | ||
| brokers.sasl_config.mechanism | string | 否 | "PLAIN" | ["PLAIN"] | SASL 认证机制 |
| brokers.sasl_config.user | string | 是 | SASL 用户名;存在 sasl_config 时必填 | ||
| brokers.sasl_config.password | string | 是 | SASL 密码;存在 sasl_config 时必填 | ||
| kafka_topic | string | 是 | 日志推送的目标 topic | ||
| producer_type | string | 否 | async | ["async", "sync"] | producer 的消息发送模式 |
| required_acks | integer | 否 | 1 | [1, -1] | leader 需要收到多少确认才认为请求完成,控制消息的持久性,语义与 Kafka 的acks一致;required_acks不能为 0 |
| key | string | 否 | 用于消息分区分配的 key | ||
| timeout | integer | 否 | 3 | [1,...] | 上游发送数据的超时时间(秒) |
| name | string | 否 | "kafka logger" | 批处理器的唯一标识;使用 Prometheus 监控 APISIX 指标时,该名称会导出在apisix_batch_process_entries中 | |
| meta_format | enum | 否 | "default" | ["default", "origin"] | 请求信息的收集格式;default为 JSON 格式,origin为原始 HTTP 请求格式 |
| log_format | object | 否 | 以 JSON 键值对声明的日志格式,值仅支持字符串,可用$前缀引用 APISIX 或 Nginx 变量 | ||
| include_req_body | boolean | 否 | false | [false, true] | 为true时在日志中包含请求体;若请求体过大无法保存在内存中,受 Nginx 限制将无法记录 |
| include_req_body_expr | array | 否 | 在include_req_body为true时生效的过滤条件,仅当表达式求值为true时才记录请求体,语法参考 lua-resty-expr | ||
| max_req_body_bytes | integer | 否 | 524288 | >=1 | 该大小以内的请求体会被推送到 Kafka,超出配置值时会在推送前被截断 |
| include_resp_body | boolean | 否 | false | [false, true] | 为true时在日志中包含响应体 |
| include_resp_body_expr | array | 否 | 在include_resp_body为true时生效的过滤条件,仅当表达式求值为true时才记录响应体 | ||
| max_resp_body_bytes | integer | 否 | 524288 | >=1 | 该大小以内的响应体会被推送到 Kafka,超出配置值时会在推送前被截断 |
| cluster_name | integer | 否 | 1 | [0,...] | 集群名称;存在两个及以上 Kafka 集群时使用,仅在producer_type为async时生效 |
| producer_batch_num | integer | 可选 | 200 | [1,...] | lua-resty-kafka 的batch_num参数,合并消息后批量发送,单位为消息条数 |
| producer_batch_size | integer | 可选 | 1048576 | [0,...] | lua-resty-kafka 的batch_size参数,单位为字节 |
| producer_max_buffering | integer | 可选 | 50000 | [1,...] | lua-resty-kafka 的max_buffering参数,表示最大缓冲大小,单位为消息条数 |
| producer_time_linger | integer | 可选 | 1 | [1,...] | lua-resty-kafka 的flush_time参数,单位为秒 |
| meta_refresh_interval | integer | 可选 | 30 | [1,...] | lua-resty-kafka 的refresh_interval参数,指定元数据自动刷新时间,单位为秒 |
对照源码中的 schema 定义 可以看到,插件对属性做了严格校验:
brokers要求minItems = 1且uniqueItems = true,每个节点必须包含host和port(port取值范围为 1~65535);sasl_config存在时user与password为必填,mechanism目前仅支持PLAIN;broker_list与brokers二选一,且必须与kafka_topic同时出现(oneOf校验),测试用例TEST 2: missing broker list即验证了缺失 broker 列表时会报错value should match only one schema, but matches none(见 t/plugin/kafka-logger.t)。
另外注意,max_req_body_bytes与max_resp_body_bytes必须为整数,若传入字符串会在 schema 校验阶段直接被拒绝,参见 t/plugin/kafka-logger-large-body.t 中的校验测试。
批处理机制:攒批发送,避免频繁提交
kafka-logger支持使用批处理器(batch processor)聚合日志并在一个批次内统一处理,从而避免频繁向 Kafka 提交数据。默认情况下,批处理器每5秒提交一次数据,或当队列中的数据达到1000条时提交。具体配置方式见 Batch Processor。
批处理器的配置项同样写在插件配置中,相关参数如下:
| 名称 | 类型 | 必填 | 默认值 | 描述 |
|---|---|---|---|---|
| name | string | 否 | logger's name | 批处理器唯一标识,默认为调用它的 logger 插件名,如 http-logger 的 name 为 "http logger" |
| batch_max_size | integer | 否 | 1000 | 每批发送的最大日志条数,达到上限后自动推送 |
| inactive_timeout | integer | 否 | 5 | 刷新缓冲区的最长时间(秒),到期后无论条数是否达到上限都会推送 |
| buffer_duration | integer | 否 | 60 | 批次中最旧条目的最大存活时间(秒),超过后该批必须被处理 |
| max_retry_count | integer | 否 | 0 | 出错时从处理管道移除条目前的最大重试次数 |
| retry_delay | integer | 否 | 1 | 执行失败后延迟重试的秒数 |
需要注意的数据流约定(官方文档以 IMPORTANT 标注):
数据首先写入缓冲区。当缓冲区超过
batch_max_size或buffer_duration时,数据被发送到 Kafka 服务器并清空缓冲区。若处理成功返回true,失败则返回nil并附带 "buffer overflow" 错误字符串。
从源码实现看,batch-processor-manager.lua 为每个插件维护一组缓冲区:add_entry找到已存在的处理器时直接push;找不到时由add_entry_to_new_processor用插件的name、batch_max_size、max_retry_count、retry_delay、buffer_duration、inactive_timeout等配置创建一个新的批处理器实例。批处理器还会每 30 分钟清理一次空缓冲区对象,避免内存泄漏。官方建议将inactive_timeout配置得比buffer_duration小,以取得最佳使用效果。
meta_format:两种日志格式详解
meta_format决定日志条目的收集格式,取值为default或origin。
default 格式(JSON 对象)
default格式将请求信息收集为结构化 JSON。官方文档示例:
{ "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 }这一结构的组装逻辑在 log-util.lua 的 get_full_log 中:包含request(url、uri、method、headers、querystring、size)、response(status、headers、size)、server(hostname、version)、upstream、service_id、route_id、consumer、client_ip、start_time(毫秒时间戳)以及latency、upstream_latency、apisix_latency三个时延指标(单位毫秒)。当开启include_req_body时还会附加request.body。
origin 格式(原始 HTTP 请求)
origin格式直接保留原始 HTTP 请求报文。官方文档示例:
GET /hello?ab=cd HTTP/1.1 host: localhost content-length: 6 connection: close abcdef源码中对应的实现是 get_req_original:将请求行ctx.var.request、每个请求头以name: value的形式拼接,空行后按需追加请求体。测试用例 t/plugin/kafka-logger.t 验证了开启meta_format: origin且include_req_body: true时,发送到 Kafka 的数据正是上述原始报文格式。
自定义日志格式:插件级与全局 Metadata
log_format允许以 JSON 键值对声明日志格式,值仅支持字符串,且可以通过$前缀引用 APISIX 变量或 Nginx 变量。例如:
{ "log_format": { "host": "$host", "@timestamp": "$time_iso8601", "client_ip": "$remote_addr" } }日志输出将变为:
{"host":"localhost","@timestamp":"2020-09-23T19:05:05-04:00","client_ip":"127.0.0.1","route_id":"1"} {"host":"localhost","@timestamp":"2020-09-23T19:05:05-04:00","client_ip":"127.0.0.1","route_id":"1"}底层实现在 get_custom_format_log:以$开头的值会被解析为变量引用,运行时从ctx.var取值;不以$开头的值则作为常量字符串直接写入日志;同时无论是否配置了自定义格式,service_id与route_id都会被自动附加。
通过 Plugin Metadata 全局配置
除了在 Route/Service 的插件配置中声明log_format,你还可以通过插件元数据(Metadata)进行全局配置。元数据配置对所有使用kafka-logger插件的 Route 和 Service 生效。
| 名称 | 类型 | 必填 | 描述 |
|---|---|---|---|
| log_format | object | 否 | 以 JSON 键值对声明的日志格式,值仅支持字符串,可用$前缀引用 APISIX 或 Nginx 变量 |
通过 Admin API 配置元数据的示例如下。首先从config.yaml中取出admin_key并保存到环境变量:
admin_key=$(yq '.deployment.admin.admin_key[0].key' conf/config.yaml | sed 's/"//g')然后调用 Admin API:
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" } }'关于log_format的优先级,从 get_log_entry 的源码逻辑可以推断:当插件配置中的conf.log_format存在时优先使用插件级格式;否则回退到元数据中的log_format;两者都没有时才使用默认的完整 JSON 日志结构。测试文件 t/plugin/kafka-logger-log-format.t 分别验证了元数据格式(TEST 1~TEST 3)与插件内格式(TEST 4)两条路径。
启用插件
以下示例展示如何在指定 Route 上启用kafka-logger:
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" }'其中batch_max_size: 1表示每条日志立即发送,便于即时验证;生产环境建议调大该值以获得批量发送的吞吐收益。
多 Broker 配置
插件支持同时向多个 broker 推送消息,在插件配置中列出多个节点即可:
"brokers" : [ { "host" :"127.0.0.1", "port" : 9092 }, { "host" :"127.0.0.1", "port" : 9093 } ],从源码看,log 阶段的 producer 复用逻辑 值得注意:插件通过core.lrucache.plugin_ctx按请求上下文缓存 producer 实例,避免每次请求都新建连接导致 Kafka 消息在分区上分布不均。producer 的底层参数映射关系如下(单位换算均在源码中完成):
request_timeout=timeout× 1000(毫秒)producer_type=producer_typerequired_acks=required_acksbatch_num=producer_batch_numbatch_size=producer_batch_sizemax_buffering=producer_max_bufferingflush_time=producer_time_linger× 1000(毫秒)refresh_interval=meta_refresh_interval× 1000(毫秒)
SASL 认证配置
若 Kafka broker 开启了 SASL 认证,可在brokers节点内配置sasl_config:
"brokers": [ { "host": "192.168.1.1", "port": 9092, "sasl_config": { "mechanism": "PLAIN", "user": "your-username", "password": "your-password" } } ]sasl_config一旦存在,user与password即为必填项(源码 schema 中通过required = {"user", "password"}强制校验)。
请求体与响应体记录
kafka-logger支持将请求体和响应体一并写入日志,便于排查问题或做报文审计:
include_req_body:置为true后,日志包含请求体。请求体由access阶段的ngx.req.read_body()读取;若请求体过大无法保存在内存中,受 Nginx 限制将无法记录。include_req_body_expr可附加条件表达式,仅当表达式为真时才读取请求体;max_req_body_bytes控制记录上限(默认 512 KiB,超出部分截断)。include_resp_body:置为true后,日志包含响应体。响应体在body_filter阶段收集,include_resp_body_expr同样支持条件过滤,max_resp_body_bytes控制上限。若响应带有 Content-Encoding 压缩(如 gzip),collect_body 会通过content_decode模块自动解码后再写入日志。
验证日志是否到达 Kafka
配置完成后,向 APISIX 发起请求即可触发日志发送:
curl -i http://127.0.0.1:9080/hello请求处理完毕后,日志会按批处理器策略被推送到 Kafka 的test2topic。测试用例中的--- error_log断言(如send data to kafka: ...)与--- wait等待机制表明:由于发送是异步的,从请求结束到日志落盘 Kafka 之间存在毫秒级到秒级的延迟窗口,验证时需要有相应的等待时间。
删除插件
要移除kafka-logger插件,将对应 Route 的插件配置清空即可。APISIX 会自动热加载,无需重启:
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 } } }'发送细节与容错机制
从 send_kafka_data 的源码可以看到发送层的两个关键行为:
- 批量编码差异:当
batch_max_size为 1 时,单条日志以 JSON 对象({})形式发送;大于 1 时以 JSON 数组([{}])形式发送。这直接决定了 Kafka 下游消费者需要解析的数据结构。 - 失败处理:发送失败时返回错误信息,包含 topic 与 broker 列表上下文,供批处理器按
max_retry_count/retry_delay策略重试;测试用例 t/plugin/kafka-logger.t 在 broker 不可达(producer_type: sync)时验证了failed to send data to Kafka topic的错误日志输出。
此外,插件对include_req_body_expr/include_resp_body_expr的表达式语法会在 schema 校验阶段(check_log_schema)提前验证,非法表达式会在配置提交时即被拒绝,而不是等到请求阶段才暴露问题(见 log-util.lua 的 check_log_schema)。
小结
kafka-logger是 APISIX 日志插件体系中面向 Kafka 的标准组件:通过brokers+kafka_topic完成最基本的接入,借助批处理器控制发送频率与吞吐,使用meta_format在结构化 JSON 与原始 HTTP 报文之间切换,再通过插件级或全局 Metadata 的log_format定制日志字段,配合include_req_body/include_resp_body实现报文级的完整审计。结合本文给出的源码路径(插件主文件、日志工具模块、批处理管理器)与测试用例(kafka-logger.t、kafka-logger-log-format.t、kafka-logger-large-body.t),你可以据此完成从配置、验证到排障的完整闭环。
【免费下载链接】apisixThe Cloud-Native API Gateway项目地址: https://gitcode.com/GitHub_Trending/ap/apisix
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考