a6-plugin-kafka-logger
概览
kafka-logger 插件会将请求/响应日志推送到 Apache Kafka
topic。它支持多个 broker、SASL 认证、异步/同步
生产、自定义日志格式和批处理,以实现高效投递。
适用场景
- 将访问日志流式写入 Kafka,供下游处理
- 为实时 API 分析流水线提供数据
- 与基于 Kafka 的日志基础设施集成
- 适用于需要 SASL 认证的 Kafka 集群
插件配置参考
核心参数
| 字段 | 类型 | 是否必填 | 默认值 | 说明 |
|---|---|---|---|---|
brokers | array | 是 | — | Kafka broker 列表 |
brokers[].host | string | 是 | — | Broker 主机名或 IP |
brokers[].port | integer | 是 | — | Broker 端口(1-65535) |
kafka_topic | string | 是 | — | 目标 Kafka topic |
key | string | 否 | — | 用于路由的分区 key |
timeout | integer | 否 | 3 | 连接超时时间,单位秒 |
SASL 认证
| 字段 | 类型 | 是否必填 | 默认值 | 说明 |
|---|---|---|---|---|
brokers[].sasl_config | object | 否 | — | 每个 broker 的 SASL 配置 |
brokers[].sasl_config.mechanism | string | 否 | "PLAIN" | PLAIN、SCRAM-SHA-256 或 SCRAM-SHA-512 |
brokers[].sasl_config.user | string | 是* | — | SASL 用户名(设置 sasl_config 时必填) |
brokers[].sasl_config.password | string | 是* | — | SASL 密码(设置 sasl_config 时必填) |
生产者配置
| 字段 | 类型 | 默认 | 描述 |
|---|---|---|---|
producer_type | 字符串 | "async" | async(批量发送,默认)或 sync(立即发送) |
required_acks | 整数 | 1 | 1(主副本确认)或 -1(所有副本确认) |
producer_batch_num | 整数 | 200 | 每批发送的 Kafka 消息数量 |
producer_batch_size | 整数 | 1048576 | 批次大小(字节,默认 1 MB) |
producer_max_buffering | 整数 | 50000 | 最大缓冲消息数量 |
producer_time_linger | 整数 | 1 | 刷新间隔(秒) |
meta_refresh_interval | 整数 | 30 | Kafka 元数据刷新间隔(秒) |
cluster_name | 整数 | 1 | 集群标识符(用于多集群) |
日志格式选项
| 字段 | 类型 | 默认 | 描述 |
|---|---|---|---|
meta_format | 字符串 | "default" | default(JSON)或 origin(原始 HTTP) |
log_format | 对象 | — | 使用 $variable 语法的自定义日志格式 |
include_req_body | 布尔 | false | 包含请求正文 |
include_req_body_expr | 数组 | — | 有条件的请求主体记录 |
include_resp_body | 布尔 | false | 包含响应体 |
include_resp_body_expr | 数组 | — | 按条件记录响应体 |
max_req_body_bytes | 整数 | 524288 | 可记录的最大请求体大小(512 KB) |
max_resp_body_bytes | 整数 | 524288 | 可记录的最大响应体大小(512 KB) |
批处理参数
| 字段 | 类型 | 默认 | 描述 |
|---|---|---|---|
batch_max_size | 整数 | 1000 | 每个批次的最大条目数 |
inactive_timeout | 整数 | 5 | 刷新未满批次前的等待时间(秒) |
buffer_duration | 整数 | 60 | 最早条目的最长保留时间(秒) |
max_retry_count | 整数 | 0 | 发送失败后的重试次数 |
retry_delay | 整数 | 1 | 重试之间的秒数 |
分步指南:将日志发送到 Kafka
1. 创建启用 kafka-logger 的路由
a6 route create -f - <<'EOF'
{
"id": "kafka-logged-api",
"uri": "/api/*",
"plugins": {
"kafka-logger": {
"brokers": [
{"host": "kafka-1", "port": 9092},
{"host": "kafka-2", "port": 9092}
],
"kafka_topic": "apisix-logs",
"batch_max_size": 100
}
},
"upstream": {
"type": "roundrobin",
"nodes": {
"backend:8080": 1
}
}
}
EOF
2. 在 Kafka 中验证消息
kafka-console-consumer --bootstrap-server kafka-1:9092 --topic apisix-logs --from-beginning
常见模式
使用 SASL 认证的 Kafka 集群
{
"plugins": {
"kafka-logger": {
"brokers": [
{
"host": "kafka.example.com",
"port": 9092,
"sasl_config": {
"mechanism": "SCRAM-SHA-256",
"user": "apisix",
"password": "secret"
}
}
],
"kafka_topic": "api-logs",
"required_acks": -1
}
}
}
自定义日志格式
{
"plugins": {
"kafka-logger": {
"brokers": [{"host": "kafka", "port": 9092}],
"kafka_topic": "api-logs",
"log_format": {
"@timestamp": "$time_iso8601",
"client_ip": "$remote_addr",
"method": "$request_method",
"uri": "$request_uri",
"status": "$status",
"latency": "$request_time",
"upstream": "$upstream_addr"
}
}
}
}
按路由 ID 分区
{
"plugins": {
"kafka-logger": {
"brokers": [{"host": "kafka", "port": 9092}],
"kafka_topic": "api-logs",
"key": "$route_id"
}
}
}
高吞吐量调优
{
"plugins": {
"kafka-logger": {
"brokers": [
{"host": "kafka-1", "port": 9092},
{"host": "kafka-2", "port": 9092},
{"host": "kafka-3", "port": 9092}
],
"kafka_topic": "api-logs",
"producer_type": "async",
"producer_batch_num": 500,
"producer_batch_size": 2097152,
"producer_max_buffering": 100000,
"producer_time_linger": 2,
"batch_max_size": 5000,
"inactive_timeout": 10,
"required_acks": 1
}
}
}
原始HTTP日志格式
{
"plugins": {
"kafka-logger": {
"brokers": [{"host": "kafka", "port": 9092}],
"kafka_topic": "raw-logs",
"meta_format": "origin"
}
}
}
生成原始HTTP请求文本而不是JSON。
配置同步示例
version: "1"
routes:
- id: kafka-logged-api
uri: /api/*
plugins:
kafka-logger:
brokers:
- host: kafka-1
port: 9092
- host: kafka-2
port: 9092
kafka_topic: apisix-logs
producer_type: async
required_acks: 1
batch_max_size: 200
inactive_timeout: 5
upstream_id: my-upstream
故障排查
| 现象 | 原因 | 修复方式 |
|---|---|---|
| Kafka 中没有消息 | Broker 不可达 | 验证 Broker 主机和端口,并检查网关节点的防火墙规则 |
| SASL 身份验证失败 | 凭证或机制错误 | 验证用户名和密码,并确保机制与 Kafka 配置匹配 |
| 消息延迟 | 批次或超时设置过大 | 降低 inactive_timeout 和 producer_time_linger |
| 消息丢失 | 缓冲区溢出 | 增大 producer_max_buffering,并增加 broker |
| 未找到 Topic | Topic 不存在且已禁用自动创建 | 手动创建 Topic,或启用 auto.create.topics.enable |
| 延迟较高 | required_acks: -1 且副本响应缓慢 | 使用 required_acks: 1 降低延迟,但持久性会相应降低 |
本文根据 api7/a6 仓库中的 a6-plugin-kafka-logger/SKILL.md 生成。可在 AI Agent Skills 页面浏览全部 Skill。