kafka-logger
kafka-logger 插件将请求和响应日志作为 JSON 对象分批推送到 Apache Kafka 集群,并支持自定义日志格式。
示例
以下示例展示了如何在不同场景下配置 kafka-logger 插件。
要跟随示例操作,请使用以下 Docker compose 文件启动一个 Kafka 集群示例:
- Docker
- Kubernetes
services:
zookeeper:
image: confluentinc/cp-zookeeper:7.8.0
container_name: zookeeper
environment:
ZOOKEEPER_CLIENT_PORT: 2181
ZOOKEEPER_TICK_TIME: 2000
networks:
- kafka_net
notkafka:
image: confluentinc/cp-kafka:7.8.0
container_name: notkafka
depends_on:
- zookeeper
environment:
KAFKA_BROKER_ID: 1
KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://notkafka:29092,PLAINTEXT_HOST://127.0.0.1:9092
KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1
KAFKA_AUTO_CREATE_TOPICS_ENABLE: "true"
ports:
- "9092:9092"
networks:
- kafka_net
networks:
kafka_net:
driver: bridge
启动容器:
docker compose up -d
为 Zookeeper 和 Kafka Deployment 创建 Kubernetes 清单文件:
apiVersion: apps/v1
kind: Deployment
metadata:
namespace: aic
name: zookeeper
spec:
replicas: 1
selector:
matchLabels:
app: zookeeper
template:
metadata:
labels:
app: zookeeper
spec:
containers:
- name: zookeeper
image: confluentinc/cp-zookeeper:7.8.0
env:
- name: ZOOKEEPER_CLIENT_PORT
value: "2181"
- name: ZOOKEEPER_TICK_TIME
value: "2000"
ports:
- containerPort: 2181
---
apiVersion: v1
kind: Service
metadata:
namespace: aic
name: zookeeper
spec:
selector:
app: zookeeper
ports:
- port: 2181
targetPort: 2181
type: ClusterIP
---
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: confluentinc/cp-kafka:7.8.0
env:
- name: KAFKA_BROKER_ID
value: "1"
- name: KAFKA_ZOOKEEPER_CONNECT
value: "zookeeper:2181"
- name: KAFKA_LISTENER_SECURITY_PROTOCOL_MAP
value: "PLAINTEXT:PLAINTEXT"
- name: KAFKA_ADVERTISED_LISTENERS
value: "PLAINTEXT://kafka-server.aic.svc:9092"
- name: KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR
value: "1"
- name: KAFKA_AUTO_CREATE_TOPICS_ENABLE
value: "true"
ports:
- containerPort: 9092
---
apiVersion: v1
kind: Service
metadata:
namespace: aic
name: kafka-server
spec:
selector:
app: kafka-server
ports:
- port: 9092
targetPort: 9092
type: ClusterIP
在配置的 Kafka 主题中等待消息:
kubectl apply -f kafka-deployment.yaml
在配置的 Kafka 主题中等待消息:
- Docker
- Kubernetes
docker exec -it notkafka kafka-console-consumer --bootstrap-server localhost:9092 --topic test2 --from-beginning
kubectl exec -n aic deploy/kafka-server -- kafka-console-consumer --bootstrap-server kafka-server.aic.svc:9092 --topic test2 --from-beginning
打开一个新的终端会话以执行以下与 APISIX 相关的步骤。
以不同的元日志格式记录日志
以下示例演示了如 何在路由上启用 kafka-logger 插件,该插件记录对路由的客户端请求并将日志推送到 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": "notkafka",
"port": 29092
}
],
"kafka_topic": "test2",
"key": "key1",
"batch_max_size": 1
}
},
"upstream": {
"nodes": {
"httpbin.org:80": 1
},
"type": "roundrobin"
}
}'
services:
- name: httpbin
routes:
- uris:
- /get
name: kafka-logger-route
plugins:
kafka-logger:
meta_format: "default"
brokers:
- host: "notkafka"
port: 29092
kafka_topic: "test2"
key: "key1"
batch_max_size: 1
upstream:
type: roundrobin
nodes:
- host: httpbin.org
port: 80
weight: 1
将配置同步到网关:
adc sync -f adc.yaml
- Gateway API
- APISIX CRD
apiVersion: v1
kind: Service
metadata:
namespace: aic
name: httpbin-external-domain
spec:
type: ExternalName
externalName: httpbin.org
---
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: "test2"
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: "test2"
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": 411.00001335144,
"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",
"x-forwarded-port": "9080"
},
"method": "GET",
"size": 83,
"uri": "/get",
"url": "http://127.0.0.1:9080/get"
},
"response": {
"headers": {
"content-length": "233",
"access-control-allow-credentials": "true",
"content-type": "application/json",
"connection": "close",
"access-control-allow-origin": "*",
"date": "Fri, 10 Nov 2023 06:02:44 GMT",
"server": "APISIX/3.16.0"
},
"status": 200,
"size": 475
},
"route_id": "kafka-logger-route",
"client_ip": "127.0.0.1",
"server": {
"hostname": "apisix",
"version": "3.16.0"
},
"apisix_latency": 18.00001335144,
"service_id": "",
"upstream_latency": 393,
"start_time": 1699596164550,
"upstream": "54.90.18.68: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
routes:
- uris:
- /get
name: kafka-logger-route
plugins:
kafka-logger:
meta_format: "origin"
brokers:
- host: "notkafka"
port: 29092
kafka_topic: "test2"
key: "key1"
batch_max_size: 1
upstream:
type: roundrobin
nodes:
- host: httpbin.org
port: 80
weight: 1
将配置同步到网关:
adc sync -f adc.yaml
- 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: "test2"
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: "test2"
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: */*
使用插件元数据记录请求和响应头
以下示例演示了如何使用插件元数据和内置变量自定义日志格式,以记录请求和响应中的特定头。
在 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": "notkafka",
"port": 29092
}
],
"kafka_topic": "test2",
"key": "key1",
"batch_max_size": 1
}
},
"upstream": {
"nodes": {
"httpbin.org:80": 1
},
"type": "roundrobin"
}
}'
services:
- name: httpbin
routes:
- uris:
- /get
name: kafka-logger-route
plugins:
kafka-logger:
meta_format: "default"
brokers:
- host: "notkafka"
port: 29092
kafka_topic: "test2"
key: "key1"
batch_max_size: 1
upstream:
type: roundrobin
nodes:
- host: httpbin.org
port: 80
weight: 1
将配置同步到网关:
adc sync -f adc.yaml
- Gateway API
- APISIX CRD
apiVersion: v1
kind: Service
metadata:
namespace: aic
name: httpbin-external-domain
spec:
type: ExternalName
externalName: httpbin.org
---
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: "test2"
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: "test2"
key: "key1"
batch_max_size: 1
应用配置:
kubectl apply -f kafka-logger-ic.yaml
❶ meta_format: 设置为 default 日志格式。请务必注意,如果你想使用插件元数据自定义日志格式,这是强制性的。如果 meta_format 设置为 origin,日志条目将保持 origin 格式。
❷ 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": {
"host": "$host",
"@timestamp": "$time_iso8601",
"client_ip": "$remote_addr",
"env": "$http_env",
"resp_content_type": "$sent_http_Content_Type"
}
}'
plugin_metadata:
- name: kafka-logger
log_format:
host: "$host"
"@timestamp": "$time_iso8601"
client_ip: "$remote_addr"
env: "$http_env"
resp_content_type: "$sent_http_Content_Type"
将配置同步到网关:
adc sync -f adc.yaml
apiVersion: apisix.apache.org/v1alpha1
kind: GatewayProxy
metadata:
namespace: aic
name: apisix-config
spec:
provider:
type: ControlPlane
controlPlane:
# ...
# 控制面连接配置
pluginMetadata:
kafka-logger:
log_format:
host: "$host"
"@timestamp": "$time_iso8601"
client_ip: "$remote_addr"
env: "$http_env"
resp_content_type: "$sent_http_Content_Type"
应用配置:
kubectl apply -f gatewayproxy.yaml
❶ 记录自定义请求头 env。
❷ 记录响应头 Content-Type。
向路由发送带有 env 头的请求:
curl -i "http://127.0.0.1:9080/get" -H "env: dev"
你应该在 Kafka 主题中看到类似于以下的日志条目:
{
"@timestamp": "2023-11-10T23:09:04+00:00",
"host": "127.0.0.1",
"client_ip": "127.0.0.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": "notkafka",
"port": 29092
}
],
"kafka_topic": "test2",
"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
routes:
- uris:
- /post
name: kafka-logger-route
plugins:
kafka-logger:
brokers:
- host: "notkafka"
port: 29092
kafka_topic: "test2"
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 sync -f adc.yaml
- Gateway API
- APISIX CRD
apiVersion: v1
kind: Service
metadata:
namespace: aic
name: httpbin-external-domain
spec:
type: ExternalName
externalName: httpbin.org
---
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: "test2"
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: "test2"
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"}'
你不应在日志中观察到请求体。
如果你除了将 include_req_body 或 include_resp_body 设置为 true 之外还自定义了 log_format,则插件将不会在日志中包含这些主体。
作为解决方法,你可以在日志格式中使用 NGINX 变量 $request_body,例如:
{
"kafka-logger": {
...,
"log_format": {"body": "$request_body"}
}
}