Apache APISIX kafka-logger 插件实战:Kafka 访问日志采集的配置、格式定制与源码解析
2026/9/15 1:05:44 网站建设 项目流程

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")

插件的priority403version0.1,其生命周期钩子分别负责不同职责(见 apisix/plugins/kafka-logger.lua):

  • access阶段:当开启include_req_body时读取请求体,支持用include_req_body_expr表达式做条件过滤;
  • body_filter阶段:调用log_util.collect_body按需收集响应体(含解压逻辑);
  • log阶段:根据meta_format选择日志条目格式,先尝试交给批处理器,若批处理器不存在则创建新的处理器并异步发送。

配置属性全解析

下表完整列出插件的全部属性(与官方文档一致,并补充了源码中的取值范围与默认值说明):

名称类型必填默认值合法值描述
broker_listobject已废弃,请改用brokers。Kafka 节点列表,键为主机名、值为端口
brokersarrayKafka 节点(broker)列表
brokers.hoststringKafka broker 主机,例如192.168.1.1
brokers.portinteger[0, 65535]Kafka broker 端口
brokers.sasl_configobjectKafka broker 的 SASL 认证配置
brokers.sasl_config.mechanismstring"PLAIN"["PLAIN"]SASL 认证机制
brokers.sasl_config.userstringSASL 用户名;存在 sasl_config 时必填
brokers.sasl_config.passwordstringSASL 密码;存在 sasl_config 时必填
kafka_topicstring日志推送的目标 topic
producer_typestringasync["async", "sync"]producer 的消息发送模式
required_acksinteger1[1, -1]leader 需要收到多少确认才认为请求完成,控制消息的持久性,语义与 Kafka 的acks一致;required_acks不能为 0
keystring用于消息分区分配的 key
timeoutinteger3[1,...]上游发送数据的超时时间(秒)
namestring"kafka logger"批处理器的唯一标识;使用 Prometheus 监控 APISIX 指标时,该名称会导出在apisix_batch_process_entries
meta_formatenum"default"["default", "origin"]请求信息的收集格式;default为 JSON 格式,origin为原始 HTTP 请求格式
log_formatobject以 JSON 键值对声明的日志格式,值仅支持字符串,可用$前缀引用 APISIX 或 Nginx 变量
include_req_bodybooleanfalse[false, true]true时在日志中包含请求体;若请求体过大无法保存在内存中,受 Nginx 限制将无法记录
include_req_body_exprarrayinclude_req_bodytrue时生效的过滤条件,仅当表达式求值为true时才记录请求体,语法参考 lua-resty-expr
max_req_body_bytesinteger524288>=1该大小以内的请求体会被推送到 Kafka,超出配置值时会在推送前被截断
include_resp_bodybooleanfalse[false, true]true时在日志中包含响应体
include_resp_body_exprarrayinclude_resp_bodytrue时生效的过滤条件,仅当表达式求值为true时才记录响应体
max_resp_body_bytesinteger524288>=1该大小以内的响应体会被推送到 Kafka,超出配置值时会在推送前被截断
cluster_nameinteger1[0,...]集群名称;存在两个及以上 Kafka 集群时使用,仅在producer_typeasync时生效
producer_batch_numinteger可选200[1,...]lua-resty-kafka 的batch_num参数,合并消息后批量发送,单位为消息条数
producer_batch_sizeinteger可选1048576[0,...]lua-resty-kafka 的batch_size参数,单位为字节
producer_max_bufferinginteger可选50000[1,...]lua-resty-kafka 的max_buffering参数,表示最大缓冲大小,单位为消息条数
producer_time_lingerinteger可选1[1,...]lua-resty-kafka 的flush_time参数,单位为秒
meta_refresh_intervalinteger可选30[1,...]lua-resty-kafka 的refresh_interval参数,指定元数据自动刷新时间,单位为秒

对照源码中的 schema 定义 可以看到,插件对属性做了严格校验:

  • brokers要求minItems = 1uniqueItems = true,每个节点必须包含hostportport取值范围为 1~65535);
  • sasl_config存在时userpassword为必填,mechanism目前仅支持PLAIN
  • broker_listbrokers二选一,且必须与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_bytesmax_resp_body_bytes必须为整数,若传入字符串会在 schema 校验阶段直接被拒绝,参见 t/plugin/kafka-logger-large-body.t 中的校验测试。

批处理机制:攒批发送,避免频繁提交

kafka-logger支持使用批处理器(batch processor)聚合日志并在一个批次内统一处理,从而避免频繁向 Kafka 提交数据。默认情况下,批处理器每5秒提交一次数据,或当队列中的数据达到1000条时提交。具体配置方式见 Batch Processor。

批处理器的配置项同样写在插件配置中,相关参数如下:

名称类型必填默认值描述
namestringlogger's name批处理器唯一标识,默认为调用它的 logger 插件名,如 http-logger 的 name 为 "http logger"
batch_max_sizeinteger1000每批发送的最大日志条数,达到上限后自动推送
inactive_timeoutinteger5刷新缓冲区的最长时间(秒),到期后无论条数是否达到上限都会推送
buffer_durationinteger60批次中最旧条目的最大存活时间(秒),超过后该批必须被处理
max_retry_countinteger0出错时从处理管道移除条目前的最大重试次数
retry_delayinteger1执行失败后延迟重试的秒数

需要注意的数据流约定(官方文档以 IMPORTANT 标注):

数据首先写入缓冲区。当缓冲区超过batch_max_sizebuffer_duration时,数据被发送到 Kafka 服务器并清空缓冲区。若处理成功返回true,失败则返回nil并附带 "buffer overflow" 错误字符串。

从源码实现看,batch-processor-manager.lua 为每个插件维护一组缓冲区:add_entry找到已存在的处理器时直接push;找不到时由add_entry_to_new_processor用插件的namebatch_max_sizemax_retry_countretry_delaybuffer_durationinactive_timeout等配置创建一个新的批处理器实例。批处理器还会每 30 分钟清理一次空缓冲区对象,避免内存泄漏。官方建议将inactive_timeout配置得比buffer_duration小,以取得最佳使用效果。

meta_format:两种日志格式详解

meta_format决定日志条目的收集格式,取值为defaultorigin

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)、upstreamservice_idroute_idconsumerclient_ipstart_time(毫秒时间戳)以及latencyupstream_latencyapisix_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: origininclude_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_idroute_id都会被自动附加。

通过 Plugin Metadata 全局配置

除了在 Route/Service 的插件配置中声明log_format,你还可以通过插件元数据(Metadata)进行全局配置。元数据配置对所有使用kafka-logger插件的 Route 和 Service 生效。

名称类型必填描述
log_formatobject以 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 1TEST 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_type
  • required_acks=required_acks
  • batch_num=producer_batch_num
  • batch_size=producer_batch_size
  • max_buffering=producer_max_buffering
  • flush_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一旦存在,userpassword即为必填项(源码 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 的源码可以看到发送层的两个关键行为:

  1. 批量编码差异:当batch_max_size为 1 时,单条日志以 JSON 对象({})形式发送;大于 1 时以 JSON 数组([{}])形式发送。这直接决定了 Kafka 下游消费者需要解析的数据结构。
  2. 失败处理:发送失败时返回错误信息,包含 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询