kafka-logger
kafka-logger 插件将请求和响应日志以 JSON 对象的形式批量发送到 Apache Kafka,并支持自定义日志格式。
示例
以下示例展示 如何将网关请求日志发送到 Kafka、自定义日志内容,以及使用 TLS 保护与 broker 的连接。
以下示例使用官方 Apache Kafka 3.9.2 镜像,并以单节点 KRaft 模式运行。
kafka-logger 支持 Kafka Produce API 版本 0–2。Kafka 4 已移除这些版本,因此在插件支持 Produce API 版本 3 或更高版本之前,请使用 Kafka 3.x broker。
- Docker
- Kubernetes
将 GATEWAY_CONTAINER 设置为正在运行的 APISIX 或 API7 网关容器:
export GATEWAY_CONTAINER=replace-with-gateway-container-name
为网关和 Kafka 创建专用网络:
docker network create gateway-kafka-net
将网关连接到该网络:
docker network connect gateway-kafka-net "$GATEWAY_CONTAINER"
创建以下 Docker Compose 文件:
services:
kafka-server:
image: apache/kafka:3.9.2
container_name: kafka-server
environment:
KAFKA_NODE_ID: 1
KAFKA_PROCESS_ROLES: broker,controller
KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka-server:9092
KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka-server:9093
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1
KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1
KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0
networks:
- kafka
networks:
kafka:
name: gateway-kafka-net
external: true
启动 broker:
docker compose up -d
容器启动后,验证 broker 能否接受请求:
docker exec kafka-server /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--list
如果命令报告连接错误,请等待几秒钟后重试。
为以 broker 和 controller 组合模式运行的单节点 Kafka broker 创建清单:
apiVersion: apps/v1
kind: Deployment
metadata:
namespace: aic
name: kafka-server
spec:
replicas: 1
selector:
matchLabels:
app: kafka-server
template:
metadata:
labels:
app: kafka-server
spec:
containers:
- name: kafka-server
image: apache/kafka:3.9.2
env:
- name: KAFKA_NODE_ID
value: "1"
- name: KAFKA_PROCESS_ROLES
value: broker,controller
- name: KAFKA_LISTENERS
value: PLAINTEXT://:9092,CONTROLLER://:9093
- name: KAFKA_ADVERTISED_LISTENERS
value: PLAINTEXT://kafka-server.aic.svc:9092
- name: KAFKA_CONTROLLER_LISTENER_NAMES
value: CONTROLLER
- name: KAFKA_LISTENER_SECURITY_PROTOCOL_MAP
value: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
- name: KAFKA_CONTROLLER_QUORUM_VOTERS
value: 1@localhost:9093
- name: KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR
value: "1"
- name: KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR
value: "1"
- name: KAFKA_TRANSACTION_STATE_LOG_MIN_ISR
value: "1"
- name: KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS
value: "0"
ports:
- name: kafka
containerPort: 9092
- name: controller
containerPort: 9093
readinessProbe:
exec:
command:
- /opt/kafka/bin/kafka-topics.sh
- --bootstrap-server
- localhost:9092
- --list
initialDelaySeconds: 5
periodSeconds: 10
timeoutSeconds: 5
---
apiVersion: v1
kind: Service
metadata:
namespace: aic
name: kafka-server
spec:
selector:
app: kafka-server
ports:
- name: kafka
port: 9092
targetPort: kafka
应用清单:
kubectl apply -f kafka-server.yaml
等待 broker 就绪:
kubectl rollout status -n aic deployment/kafka-server
创建示例使用的主题:
- Docker
- Kubernetes
docker exec kafka-server /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--create \
--if-not-exists \
--topic apisix-logs
kubectl exec -n aic deploy/kafka-server -- \
/opt/kafka/bin/kafka-topics.sh \
--bootstrap-server localhost:9092 \
--create \
--if-not-exists \
--topic apisix-logs
如需查看示例生成的记录,请在另一个终端中运行 consumer:
- Docker
- Kubernetes
docker exec -it kafka-server /opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server localhost:9092 \
--topic apisix-logs \
--from-beginning
kubectl exec -n aic -it deploy/kafka-server -- \
/opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server localhost:9092 \
--topic apisix-logs \
--from-beginning
以不同的元日志格式记录日志
以下示例将路由请求日志发送到 Kafka,并展示 default 与 origin 元日志格式的区别。
按如下方式创建启用了 kafka-logger 的路由:
- Admin API
- ADC
- Ingress Controller
curl "http://127.0.0.1:9180/apisix/admin/routes" -X PUT \
-H "X-API-KEY: ${ADMIN_API_KEY}" \
-d '{
"id": "kafka-logger-route",
"uri": "/get",
"plugins": {
"kafka-logger": {
"meta_format": "default",
"brokers": [
{
"host": "kafka-server",
"port": 9092
}
],
"kafka_topic": "apisix-logs",
"key": "key1",
"batch_max_size": 1
}
},
"upstream": {
"nodes": {
"httpbin.org:80": 1
},
"type": "roundrobin"
}
}'
services:
- name: httpbin
labels:
docs-example: kafka-logging
routes:
- uris:
- /get
name: kafka-logger-route
plugins:
kafka-logger:
meta_format: "default"
brokers:
- host: "kafka-server"
port: 9092
kafka_topic: "apisix-logs"
key: "key1"
batch_max_size: 1
upstream:
type: roundrobin
nodes:
- host: httpbin.org
port: 80
weight: 1
预览带有示例标签的服务变更:
adc diff -f adc.yaml \
--include-resource-type service \
--label-selector docs-example=kafka-logging
同步已确认的变更:
adc sync -f adc.yaml \
--include-resource-type service \
--label-selector docs-example=kafka-logging
- Gateway API
- APISIX CRD
apiVersion: v1
kind: Service
metadata:
namespace: aic
name: httpbin-external-domain
spec:
type: ExternalName
externalName: httpbin.org
ports:
- name: http
port: 80
targetPort: 80
---
apiVersion: apisix.apache.org/v1alpha1
kind: PluginConfig
metadata:
namespace: aic
name: kafka-logger-plugin-config
spec:
plugins:
- name: kafka-logger
config:
meta_format: "default"
brokers:
- host: "kafka-server.aic.svc"
port: 9092
kafka_topic: "apisix-logs"
key: "key1"
batch_max_size: 1
---
apiVersion: gateway.networking.k8s.io/v1
kind: HTTPRoute
metadata:
namespace: aic
name: kafka-logger-route
spec:
parentRefs:
- name: apisix
rules:
- matches:
- path:
type: Exact
value: /get
filters:
- type: ExtensionRef
extensionRef:
group: apisix.apache.org
kind: PluginConfig
name: kafka-logger-plugin-config
backendRefs:
- name: httpbin-external-domain
port: 80
apiVersion: apisix.apache.org/v2
kind: ApisixUpstream
metadata:
namespace: aic
name: httpbin-external-domain
spec:
ingressClassName: apisix
externalNodes:
- type: Domain
name: httpbin.org
---
apiVersion: apisix.apache.org/v2
kind: ApisixRoute
metadata:
namespace: aic
name: kafka-logger-route
spec:
ingressClassName: apisix
http:
- name: kafka-logger-route
match:
paths:
- /get
methods:
- GET
upstreams:
- name: httpbin-external-domain
plugins:
- name: kafka-logger
config:
meta_format: "default"
brokers:
- host: "kafka-server.aic.svc"
port: 9092
kafka_topic: "apisix-logs"
key: "key1"
batch_max_size: 1
应用配置:
kubectl apply -f kafka-logger-ic.yaml
❶ meta_format:设置为 default 日志格式。
❷ batch_max_size:设置为 1,以立即发送日志条目。
向路由发送请求以生成日志条目:
curl -i "http://127.0.0.1:9080/get"
你应该会看到 HTTP/1.1 200 OK 响应。
你应该会在 Kafka 主题中看到类似以下内容的日志条目。地址、时间值和版本详情会因环境而异:
{
"latency": 1030.9998989105,
"request": {
"querystring": {},
"headers": {
"host": "127.0.0.1:9080",
"user-agent": "curl/8.7.1",
"accept": "*/*",
"x-forwarded-proto": "http",
"x-forwarded-host": "127.0.0.1:9080",
"x-forwarded-port": "9080"
},
"method": "GET",
"size": 80,
"uri": "/get",
"url": "http://127.0.0.1:9080/get"
},
"response": {
"headers": {
"content-length": "311",
"access-control-allow-credentials": "true",
"content-type": "application/json",
"connection": "close",
"access-control-allow-origin": "*",
"date": "Fri, 18 Sep 2026 03:32:31 GMT",
"server": "APISIX/3.18.0"
},
"status": 200,
"size": 539
},
"route_id": "kafka-logger-route",
"client_ip": "192.168.155.1",
"server": {
"hostname": "dd2886d0b7bf",
"version": "3.18.0"
},
"apisix_latency": 408.99989891052,
"service_id": "",
"upstream_latency": 622,
"start_time": 1789702349311,
"upstream": "98.88.64.13:80"
}
将元日志格式更新为 origin:
- Admin API
- ADC
- Ingress Controller
curl "http://127.0.0.1:9180/apisix/admin/routes/kafka-logger-route" -X PATCH \
-H "X-API-KEY: ${ADMIN_API_KEY}" \
-d '{
"plugins": {
"kafka-logger": {
"meta_format": "origin"
}
}
}'
更新 adc.yaml,将 meta_format 设置为 origin:
services:
- name: httpbin
labels:
docs-example: kafka-logging
routes:
- uris:
- /get
name: kafka-logger-route
plugins:
kafka-logger:
meta_format: "origin"
brokers:
- host: "kafka-server"
port: 9092
kafka_topic: "apisix-logs"
key: "key1"
batch_max_size: 1
upstream:
type: roundrobin
nodes:
- host: httpbin.org
port: 80
weight: 1
预览带有示例标签的服务变更:
adc diff -f adc.yaml \
--include-resource-type service \
--label-selector docs-example=kafka-logging
同步已确认的变更:
adc sync -f adc.yaml \
--include-resource-type service \
--label-selector docs-example=kafka-logging
- Gateway API
- APISIX CRD
更新 kafka-logger-ic.yaml,将 meta_format 设置为 origin:
apiVersion: apisix.apache.org/v1alpha1
kind: PluginConfig
metadata:
namespace: aic
name: kafka-logger-plugin-config
spec:
plugins:
- name: kafka-logger
config:
meta_format: "origin"
brokers:
- host: "kafka-server.aic.svc"
port: 9092
kafka_topic: "apisix-logs"
key: "key1"
batch_max_size: 1
更新 kafka-logger-ic.yaml,将 meta_format 设置为 origin:
apiVersion: apisix.apache.org/v2
kind: ApisixRoute
metadata:
namespace: aic
name: kafka-logger-route
spec:
ingressClassName: apisix
http:
- name: kafka-logger-route
match:
paths:
- /get
methods:
- GET
upstreams:
- name: httpbin-external-domain
plugins:
- name: kafka-logger
config:
meta_format: "origin"
brokers:
- host: "kafka-server.aic.svc"
port: 9092
kafka_topic: "apisix-logs"
key: "key1"
batch_max_size: 1
应用更新后的配置:
kubectl apply -f kafka-logger-ic.yaml
再次向路由发送请求以生成新的日志条目:
curl -i "http://127.0.0.1:9080/get"
你应该会看到 HTTP/1.1 200 OK 响应。
你应该会在 Kafka 主题中看到类似以下内容的日志条目:
GET /get HTTP/1.1
x-forwarded-proto: http
x-forwarded-host: 127.0.0.1
user-agent: curl/8.7.1
x-forwarded-port: 9080
host: 127.0.0.1:9080
accept: */*
将日志发送到启用了 TLS 的 broker
以下 Docker 示例在 gateway-kafka-net 网络中启动一个使用 CA 签名 TLS 证书的本地 Kafka 3.9.2 broker。随后,示例会在启用证书校验的情况下连接该 broker,并使用 Produce API 版本 2,使 Kafka 能够记录消息时间戳。继续操作前,请安装 OpenSSL 和 Java keytool 命令 。
生成示例 CA、用于 kafka-tls 容器主机名的 broker 证书,以及 Kafka 所需的 Java 密钥库和信任库文件:
mkdir -p kafka-tls-certs
openssl req -x509 -newkey rsa:2048 -nodes -days 365 \
-subj "/CN=kafka-example-ca" \
-keyout kafka-tls-certs/ca.key \
-out kafka-tls-certs/ca.crt
openssl req -newkey rsa:2048 -nodes \
-subj "/CN=kafka-tls" \
-keyout kafka-tls-certs/server.key \
-out kafka-tls-certs/server.csr
printf "subjectAltName=DNS:kafka-tls\n" > kafka-tls-certs/server-ext.cnf
openssl x509 -req -days 365 \
-in kafka-tls-certs/server.csr \
-CA kafka-tls-certs/ca.crt \
-CAkey kafka-tls-certs/ca.key \
-CAcreateserial \
-extfile kafka-tls-certs/server-ext.cnf \
-out kafka-tls-certs/server.crt
openssl pkcs12 -export \
-name kafka-tls \
-in kafka-tls-certs/server.crt \
-inkey kafka-tls-certs/server.key \
-certfile kafka-tls-certs/ca.crt \
-out kafka-tls-certs/kafka.keystore.p12 \
-passout pass:changeit
keytool -importkeystore -noprompt \
-srckeystore kafka-tls-certs/kafka.keystore.p12 \
-srcstoretype PKCS12 \
-srcstorepass changeit \
-destkeystore kafka-tls-certs/kafka.keystore.jks \
-deststoretype JKS \
-deststorepass changeit \
-destkeypass changeit
keytool -importcert -noprompt \
-alias kafka-example-ca \
-file kafka-tls-certs/ca.crt \
-keystore kafka-tls-certs/kafka.truststore.jks \
-storepass changeit
printf "changeit\n" > kafka-tls-certs/kafka_keystore_creds
printf "changeit\n" > kafka-tls-certs/kafka_ssl_key_creds
创建稍后用于验证记录的客户端配置:
security.protocol=SSL
ssl.truststore.location=/etc/kafka/secrets/kafka.truststore.jks
ssl.truststore.password=changeit
ssl.endpoint.identification.algorithm=https
在网关网络中启动启用了 TLS 的 broker:
docker run -d \
--name kafka-tls \
--hostname kafka-tls \
--network gateway-kafka-net \
-v "${PWD}/kafka-tls-certs:/etc/kafka/secrets:ro" \
-e KAFKA_NODE_ID=1 \
-e KAFKA_PROCESS_ROLES=broker,controller \
-e KAFKA_LISTENER_SECURITY_PROTOCOL_MAP="SSL:SSL,CONTROLLER:PLAINTEXT" \
-e KAFKA_ADVERTISED_LISTENERS="SSL://kafka-tls:9093" \
-e KAFKA_LISTENERS="SSL://:9093,CONTROLLER://:29093" \
-e KAFKA_CONTROLLER_QUORUM_VOTERS="1@kafka-tls:29093" \
-e KAFKA_CONTROLLER_LISTENER_NAMES=CONTROLLER \
-e KAFKA_INTER_BROKER_LISTENER_NAME=SSL \
-e KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1 \
-e KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS=0 \
-e KAFKA_TRANSACTION_STATE_LOG_MIN_ISR=1 \
-e KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR=1 \
-e KAFKA_SSL_KEYSTORE_FILENAME=kafka.keystore.jks \
-e KAFKA_SSL_KEYSTORE_CREDENTIALS=kafka_keystore_creds \
-e KAFKA_SSL_KEY_CREDENTIALS=kafka_ssl_key_creds \
-e KAFKA_SSL_TRUSTSTORE_LOCATION=/etc/kafka/secrets/kafka.truststore.jks \
-e KAFKA_SSL_TRUSTSTORE_PASSWORD=changeit \
-e KAFKA_SSL_CLIENT_AUTH=none \
-e CLUSTER_ID="4L6g3nShT-eMCtK--X86sw" \
apache/kafka:3.9.2
容器启动后,验证 TLS 监听器能否接受请求:
docker exec kafka-tls /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server kafka-tls:9093 \
--command-config /etc/kafka/secrets/client.properties \
--list
如果命令报告连接错误,请等待几秒钟后重试。
创建主题:
docker exec kafka-tls /opt/kafka/bin/kafka-topics.sh \
--bootstrap-server kafka-tls:9093 \
--command-config /etc/kafka/secrets/client.properties \
--create \
--if-not-exists \
--topic apisix-logs \
--partitions 1 \
--replication-factor 1
为路由配置设置 broker 地址和主题:
export KAFKA_TLS_HOST="kafka-tls"
export KAFKA_TLS_PORT="9093"
export KAFKA_TOPIC="apisix-logs"
将生成的 CA 证书复制到 GATEWAY_CONTAINER 指定的网关容器中:
docker cp kafka-tls-certs/ca.crt \
"$GATEWAY_CONTAINER":/usr/local/apisix/conf/kafka-example-ca.crt
将该证书追加到现有系统信任包中,使网关继续信任容器内已经安装的公共 CA 证书:
docker exec "$GATEWAY_CONTAINER" sh -c '
cat /etc/ssl/certs/ca-certificates.crt \
/usr/local/apisix/conf/kafka-example-ca.crt \
> /usr/local/apisix/conf/combined-ca-bundle.pem
'
更新快速入门配置中的信任包路径,并重新加载网关:
docker exec -u 0 "$GATEWAY_CONTAINER" sh -c '
config=/usr/local/apisix/conf/config.yaml
certificate=/usr/local/apisix/conf/combined-ca-bundle.pem
output=$(mktemp /tmp/kafka-config.XXXXXX)
if grep -q "^ ssl_trusted_certificate:" "$config"; then
sed "s#ssl_trusted_certificate:.*#ssl_trusted_certificate: $certificate#" "$config"
elif grep -q "^ ssl:$" "$config"; then
sed "/^ ssl:$/a\\
ssl_trusted_certificate: $certificate" "$config"
else
sed "/^apisix:$/a\\
ssl:\\
ssl_trusted_certificate: $certificate" "$config"
fi > "$output"
cat "$output" > "$config"
rm "$output"
apisix reload
'
对于多实例部署,请将合并后的信任包和配置变更分发到每个网关实例。kafka-tls 主机名与示例 broker 证书中的 DNS 名称一致。
创建一个将每条日志立即发送到 TLS 监听器的路由:
curl "http://127.0.0.1:9180/apisix/admin/routes" -X PUT \
-H "X-API-KEY: ${ADMIN_API_KEY}" \
-d @- <<EOF
{
"id": "kafka-logger-tls-route",
"uri": "/get",
"plugins": {
"kafka-logger": {
"brokers": [
{
"host": "${KAFKA_TLS_HOST}",
"port": ${KAFKA_TLS_PORT}
}
],
"kafka_topic": "${KAFKA_TOPIC}",
"api_version": 2,
"batch_max_size": 1,
"tls": {
"verify": true
}
}
},
"upstream": {
"nodes": {
"httpbin.org:80": 1
},
"type": "roundrobin"
}
}
EOF
发送请求以生成日志条目:
curl -i "http://127.0.0.1:9080/get"
从 broker 容器中消费一条记录并输出其时间戳:
docker exec kafka-tls /opt/kafka/bin/kafka-console-consumer.sh \
--bootstrap-server kafka-tls:9093 \
--topic apisix-logs \
--consumer.config /etc/kafka/secrets/client.properties \
--from-beginning \
--max-messages 1 \
--property print.timestamp=true
消费到的记录应包含请求日志以及晚于 Unix 纪元的时间戳。如果证书校验失败,请确认网关能够读取 CA 信任包,并且 KAFKA_TLS_HOST 与 broker 证书中的某个名称一致。
该记录至少应包含以下字段和 broker 时间戳。默认日志条目还包含请求、响应、延迟和网关等字段。
{
"route_id": "kafka-logger-tls-route",
"request": {
"method": "GET",
"uri": "/get"
},
"response": {
"status": 200
}
}
使用插件元数据添加请求头和响应头
以下示例使用插件元数据和内置变量,为每个 kafka-logger 实例添加指定的请求头和响应头。
在 APISIX 中,插件元数据用于配置同一插件所有实例的通用元数据字段。当插件在多个资源上 启用,且需要统一更新这些实例的元数据字段时,此功能非常有用。
首先,按如下方式创建启用了 kafka-logger 的路由:
- Admin API
- ADC
- Ingress Controller
curl "http://127.0.0.1:9180/apisix/admin/routes" -X PUT \
-H "X-API-KEY: ${ADMIN_API_KEY}" \
-d '{
"id": "kafka-logger-route",
"uri": "/get",
"plugins": {
"kafka-logger": {
"meta_format": "default",
"brokers": [
{
"host": "kafka-server",
"port": 9092
}
],
"kafka_topic": "apisix-logs",
"key": "key1",
"batch_max_size": 1
}
},
"upstream": {
"nodes": {
"httpbin.org:80": 1
},
"type": "roundrobin"
}
}'
services:
- name: httpbin
labels:
docs-example: kafka-logging
routes:
- uris:
- /get
name: kafka-logger-route
plugins:
kafka-logger:
meta_format: "default"
brokers:
- host: "kafka-server"
port: 9092
kafka_topic: "apisix-logs"
key: "key1"
batch_max_size: 1
upstream:
type: roundrobin
nodes:
- host: httpbin.org
port: 80
weight: 1
预览带有示例标签的服务变更:
adc diff -f adc.yaml \
--include-resource-type service \
--label-selector docs-example=kafka-logging
同步已确认的变更:
adc sync -f adc.yaml \
--include-resource-type service \
--label-selector docs-example=kafka-logging
- Gateway API
- APISIX CRD
apiVersion: v1
kind: Service
metadata:
namespace: aic
name: httpbin-external-domain
spec:
type: ExternalName
externalName: httpbin.org
ports:
- name: http
port: 80
targetPort: 80
---
apiVersion: apisix.apache.org/v1alpha1
kind: PluginConfig
metadata:
namespace: aic
name: kafka-logger-plugin-config
spec:
plugins:
- name: kafka-logger
config:
meta_format: "default"
brokers:
- host: "kafka-server.aic.svc"
port: 9092
kafka_topic: "apisix-logs"
key: "key1"
batch_max_size: 1
---
apiVersion: gateway.networking.k8s.io/v1
kind: HTTPRoute
metadata:
namespace: aic
name: kafka-logger-route
spec:
parentRefs:
- name: apisix
rules:
- matches:
- path:
type: Exact
value: /get
filters:
- type: ExtensionRef
extensionRef:
group: apisix.apache.org
kind: PluginConfig
name: kafka-logger-plugin-config
backendRefs:
- name: httpbin-external-domain
port: 80
apiVersion: apisix.apache.org/v2
kind: ApisixUpstream
metadata:
namespace: aic
name: httpbin-external-domain
spec:
ingressClassName: apisix
externalNodes:
- type: Domain
name: httpbin.org
---
apiVersion: apisix.apache.org/v2
kind: ApisixRoute
metadata:
namespace: aic
name: kafka-logger-route
spec:
ingressClassName: apisix
http:
- name: kafka-logger-route
match:
paths:
- /get
methods:
- GET
upstreams:
- name: httpbin-external-domain
plugins:
- name: kafka-logger
config:
meta_format: "default"
brokers:
- host: "kafka-server.aic.svc"
port: 9092
kafka_topic: "apisix-logs"
key: "key1"
batch_max_size: 1
应用配置:
kubectl apply -f kafka-logger-ic.yaml
❶ meta_format:保留 default 格式,以包含插件元数据中的字段。当此字段设置为 origin 时,通过 log_format_extra 配置的字段会被忽略。
❷ batch_max_size:设置为 1,以立即发送日志条目。
接下来,为 kafka-logger 配置插件元数据:
- Admin API
- ADC
- Ingress Controller
curl "http://127.0.0.1:9180/apisix/admin/plugin_metadata/kafka-logger" -X PUT \
-H "X-API-KEY: ${ADMIN_API_KEY}" \
-d '{
"log_format_extra": {
"host": "$host",
"@timestamp": "$time_iso8601",
"client_ip": "$remote_addr",
"env": "$http_env",
"resp_content_type": "$sent_http_Content_Type"
}
}'
插件元数据是全局集合,无法通过标签选择器隔离。更改此条目前,请先导出完整集合:
adc dump -o adc-metadata.yaml --with-id \
--include-resource-type plugin_metadata
在保留 adc-metadata.yaml 中所有其他条目的同时,添加或更新 kafka-logger 条目:
plugin_metadata:
# 保留导出文件中的所有其他插件元数据条目。
kafka-logger:
log_format_extra:
host: "$host"
"@timestamp": "$time_iso8601"
client_ip: "$remote_addr"
env: "$http_env"
resp_content_type: "$sent_http_Content_Type"
预览完整的元数据变更,并确认其中没有非预期的更新或删除:
adc diff -f adc-metadata.yaml \
--include-resource-type plugin_metadata
同步已确认的插件元数据集合:
adc sync -f adc-metadata.yaml \
--include-resource-type plugin_metadata
在部署使用的完整 GatewayProxy 清单中的 spec.pluginMetadata 下添加以下条目:
kafka-logger:
log_format_extra:
host: "$host"
"@timestamp": "$time_iso8601"
client_ip: "$remote_addr"
env: "$http_env"
resp_content_type: "$sent_http_Content_Type"
通过部署的常规 Kubernetes 或 GitOps 工作流应用更新后的完整清单。
❶ 记录自定义请求头 env。
❷ 记录响应头 Content-Type。
向路由发送带有 env 请求头的请求:
curl -i "http://127.0.0.1:9080/get" -H "env: dev"
你应该会在 Kafka 主题中看到类似以下内容的日志条目:
{
"@timestamp": "2026-09-18T03:32:31+00:00",
"host": "127.0.0.1",
"client_ip": "192.168.155.1",
"route_id": "kafka-logger-route",
"env": "dev",
"resp_content_type": "application/json"
}
有条件地记录 请求体
以下示例仅在查询参数满足已配置的表达式时记录请求体。
按如下方式创建启用了 kafka-logger 的路由:
- Admin API
- ADC
- Ingress Controller
curl "http://127.0.0.1:9180/apisix/admin/routes" -X PUT \
-H "X-API-KEY: ${ADMIN_API_KEY}" \
-d '{
"id": "kafka-logger-route",
"uri": "/post",
"plugins": {
"kafka-logger": {
"brokers": [
{
"host": "kafka-server",
"port": 9092
}
],
"kafka_topic": "apisix-logs",
"key": "key1",
"batch_max_size": 1,
"include_req_body": true,
"include_req_body_expr": [["arg_log_body", "==", "yes"]]
}
},
"upstream": {
"nodes": {
"httpbin.org:80": 1
},
"type": "roundrobin"
}
}'
services:
- name: httpbin
labels:
docs-example: kafka-logging
routes:
- uris:
- /post
name: kafka-logger-route
plugins:
kafka-logger:
brokers:
- host: "kafka-server"
port: 9092
kafka_topic: "apisix-logs"
key: "key1"
batch_max_size: 1
include_req_body: true
include_req_body_expr:
- - "arg_log_body"
- "=="
- "yes"
upstream:
type: roundrobin
nodes:
- host: httpbin.org
port: 80
weight: 1
预览带有示例标签的服务变更:
adc diff -f adc.yaml \
--include-resource-type service \
--label-selector docs-example=kafka-logging
同步已确认的变更:
adc sync -f adc.yaml \
--include-resource-type service \
--label-selector docs-example=kafka-logging
- Gateway API
- APISIX CRD
apiVersion: v1
kind: Service
metadata:
namespace: aic
name: httpbin-external-domain
spec:
type: ExternalName
externalName: httpbin.org
ports:
- name: http
port: 80
targetPort: 80
---
apiVersion: apisix.apache.org/v1alpha1
kind: PluginConfig
metadata:
namespace: aic
name: kafka-logger-plugin-config
spec:
plugins:
- name: kafka-logger
config:
brokers:
- host: "kafka-server.aic.svc"
port: 9092
kafka_topic: "apisix-logs"
key: "key1"
batch_max_size: 1
include_req_body: true
include_req_body_expr:
- - "arg_log_body"
- "=="
- "yes"
---
apiVersion: gateway.networking.k8s.io/v1
kind: HTTPRoute
metadata:
namespace: aic
name: kafka-logger-route
spec:
parentRefs:
- name: apisix
rules:
- matches:
- path:
type: Exact
value: /post
filters:
- type: ExtensionRef
extensionRef:
group: apisix.apache.org
kind: PluginConfig
name: kafka-logger-plugin-config
backendRefs:
- name: httpbin-external-domain
port: 80
apiVersion: apisix.apache.org/v2
kind: ApisixUpstream
metadata:
namespace: aic
name: httpbin-external-domain
spec:
ingressClassName: apisix
externalNodes:
- type: Domain
name: httpbin.org
---
apiVersion: apisix.apache.org/v2
kind: ApisixRoute
metadata:
namespace: aic
name: kafka-logger-route
spec:
ingressClassName: apisix
http:
- name: kafka-logger-route
match:
paths:
- /post
methods:
- POST
upstreams:
- name: httpbin-external-domain
plugins:
- name: kafka-logger
config:
brokers:
- host: "kafka-server.aic.svc"
port: 9092
kafka_topic: "apisix-logs"
key: "key1"
batch_max_size: 1
include_req_body: true
include_req_body_expr:
- - "arg_log_body"
- "=="
- "yes"
应用配置:
kubectl apply -f kafka-logger-ic.yaml
❶ include_req_body:设置为 true 以记录请求体。
❷ include_req_body_expr:仅当 URL 查询字符串 log_body 为 yes 时记录请求体。
向路由发送带有满足条件的 URL 查询字符串的请求:
curl -i "http://127.0.0.1:9080/post?log_body=yes" -X POST -d '{"env": "dev"}'
你应该会看到请求体已记录在日志中:
{
...,
"method": "POST",
"body": "{\"env\": \"dev\"}",
"size": 179
}
}
向路由发送不带任何 URL 查询字符串的请求:
curl -i "http://127.0.0.1:9080/post" -X POST -d '{"env": "dev"}'
你应该不会在日志中看到请求体。
上面使用的 log_format_extra 字段会保留默 认日志条目,其中包括插件收集的请求体和响应体。如果改为配置 log_format,请显式包含相应变量:
{
"include_req_body": true,
"include_resp_body": true,
"log_format": {
"request_body": "$request_body",
"response_body": "$resp_body"
}
}
请求体和响应体的大小限制仍然适用。使用 log_format_extra 添加自定义字段,而无需替换默认日志条目。