Publish request and response logs to an Apache Kafka topic. For more information, see Kafka topics.
Kong also provides a Kafka plugin for request transformations. See Kafka Upstream.
Publish request and response logs to an Apache Kafka topic. For more information, see Kafka topics.
Kong also provides a Kafka plugin for request transformations. See Kafka Upstream.
Note: If the
max_batch_sizeargument > 1, a request is logged as an array of JSON objects.
Every request is logged separately in a JSON object, separated by a new line \n.
Expand this block to see a sample log object
{ "response": { "size": 9982, "headers": { "access-control-allow-origin": "*", "content-length": "9593", "date": "Thu, 19 Sep 2024 22:10:39 GMT", "content-type": "text/html; charset=utf-8", "via": "1.1 kong/3.8.0.0-enterprise-edition", "connection": "close", "server": "gunicorn/19.9.0", "access-control-allow-credentials": "true", "x-kong-upstream-latency": "171", "x-kong-proxy-latency": "1", "x-kong-request-id": "2f6946328ffc4946b8c9120704a4a155" }, "status": 200 }, "route": { "updated_at": 1726782477, "tags": [], "response_buffering": true, "path_handling": "v0", "protocols": [ "http", "https" ], "service": { "id": "fb4eecf8-dec2-40ef-b779-16de7e2384c7" }, "https_redirect_status_code": 426, "regex_priority": 0, "name": "example_route", "id": "0f1a4101-3327-4274-b1e4-484a4ab0c030", "strip_path": true, "preserve_host": false, "created_at": 1726782477, "request_buffering": true, "ws_id": "f381e34e-5c25-4e65-b91b-3c0a86cfc393", "paths": [ "/example-route" ] }, "workspace": "f381e34e-5c25-4e65-b91b-3c0a86cfc393", "workspace_name": "default", "tries": [ { "balancer_start": 1726783839539, "balancer_start_ns": 1.7267838395395e+18, "ip": "34.237.204.224", "balancer_latency": 0, "port": 80, "balancer_latency_ns": 27904 } ], "client_ip": "192.168.65.1", "request": { "id": "2f6946328ffc4946b8c9120704a4a155", "headers": { "accept": "*/*", "user-agent": "HTTPie/3.2.3", "host": "localhost:8000", "connection": "keep-alive", "accept-encoding": "gzip, deflate" }, "uri": "/example-route", "size": 139, "method": "GET", "querystring": {}, "url": "http://localhost:8000/example-route" }, "upstream_uri": "/", "started_at": 1726783839538, "source": "upstream", "upstream_status": "200", "latencies": { "kong": 1, "proxy": 171, "request": 173, "receive": 1 }, "service": { "write_timeout": 60000, "read_timeout": 60000, "updated_at": 1726782459, "host": "httpbin.konghq.com", "name": "example_service", "id": "fb4eecf8-dec2-40ef-b779-16de7e2384c7", "port": 80, "enabled": true, "created_at": 1726782459, "protocol": "http", "ws_id": "f381e34e-5c25-4e65-b91b-3c0a86cfc393", "connect_timeout": 60000, "retries": 5 } }
This plugin uses the lua-resty-kafka client.
When encoding request bodies, several things happen:
application/x-www-form-urlencoded, multipart/form-data,
or application/json, this plugin passes the raw request body in the body attribute, and tries
to return a parsed version of those arguments in body_args.
If this parsing fails, the plugin returns an error message and the message isn’t sent.content-type is not text/plain, text/html, application/xml, text/xml, or application/soap+xml,
then the body will be base64-encoded to ensure that the message can be sent as JSON. In that case,
the message has an extra attribute called body_base64 set to true.The custom_fields_by_lua configuration allows for the dynamic modification of
log fields using Lua code. For example, here is a snippet of an example configuration that
removes the route field from the logs:
curl -i -X POST http://localhost:8001/plugins \
--data config.name=kafka-log \
--data config.custom_fields_by_lua.route="return nil"Similarly, new fields can be added:
curl -i -X POST http://localhost:8001/plugins \
--data config.name=kafka-log \
--data config.custom_fields_by_lua.header="return kong.request.get_header('h1')"Array indices should be enclosed within square brackets. For example:
curl -i -X POST http://localhost:8001/plugins \
--header 'Accept: application/json' \
--header 'Content-Type: application/json' \
--data '{
"name": "kafka-log",
"config": {
"custom_fields_by_lua": {
"foo[1].bar[2].woo": "return 456"
}
}
}'Array indices only support positive integers.
Dot characters (.) in the field key create nested fields. You can use a backslash \ to escape a dot if you want to keep it in the field name.
For example, if you configure a field with both a regular dot and an escaped dot:
curl -i -X POST http://localhost:8001/plugins/ \
...
--data config.name=kafka-log \
--data config.custom_fields_by_lua.[my_entry.log\.field]="return foo"The field will look like this in the log:
"my_entry": {
"log.field": "foo"
}All logging plugins use the same table for logging.
If you set custom_fields_by_lua in one plugin, all logging plugins that execute after that plugin will also use the same configuration.
For example, if you configure fields via custom_fields_by_lua in File Log, those same fields will appear in Syslog, since Kafka Log executes first.
If you want all logging plugins to use the same configuration, we recommend using the Pre-function plugin to call kong.log.set_serialize_value so that the function is applied predictably and is easier to manage.
If you don’t want all logging plugins to use the same configuration, you need to manually disable the relevant fields in each plugin.
For example, if you configure a field in File Log that you don’t want appearing in Syslog, set that field to return nil in the File Log plugin:
curl -i -X POST http://localhost:8001/plugins/ \
...
--data config.name=kafka-log \
--data config.custom_fields_by_lua.my_file_log_field="return nil"See the plugin execution order reference for more details on plugin ordering.
Lua code runs in a restricted sandbox environment, whose behavior is governed
by the untrusted_lua configuration properties.
Sandboxing imposes several limitations on how custom Lua code can be executed, for heightened security. The Lua (or LuaJIT) language itself is not limited — only the available environment and the set of usable modules are restricted.
The limitations can be adjusted with the untrusted_lua=off|strict|lax|sandbox|on setting.
See the sandboxing reference for more information.
Further, as code runs in the context of the log phase, only PDK methods that can run in said phase can be used.
The Kafka Log plugin supports integration with Confluent Schema Registry for AVRO and JSON schemas.
Schema registries provide a centralized repository for managing and validating schemas for data formats like AVRO and JSON. Integrating with a schema registry allows the plugin to validate and serialize/deserialize messages in a standardized format.
Using a schema registry with Kong Gateway provides several benefits:
To learn more about Kong’s supported schema registry, see:
When a producer plugin is configured with a schema registry, the following workflow occurs:
sequenceDiagram
autonumber
participant Client
participant Kong as Kafka Log plugin
participant Registry as Schema Registry
participant Kafka
activate Client
activate Kong
Client->>Kong: Send request
deactivate Client
activate Registry
Kong->>Registry: Fetch schema from registry
Registry-->>Kong: Return schema
deactivate Registry
Kong->>Kong: Validate message against schema
Kong->>Kong: Serialize using schema
activate Kafka
Kong->>Kafka: Forward to Kafka
deactivate Kong
deactivate Kafka
If validation fails, the request is rejected with an error message.
To configure Schema Registry with the Kafka Log plugin, use the config.schema_registry parameter in your plugin configuration.
For sample configuration values, see:
By default, an Avro schema requires every union-typed value to be wrapped in a single-key object that names the union branch.
For example, a field declared as ["null", "string"] must be sent as {"string": "hello"} or {"null": null}.
This forces HTTP clients to understand Avro’s wire encoding to call your API.
Set payload_encoding: simple_json on the value_schema or the key_schema to let the Kafka Log plugin accept plain JSON instead, and resolve union branches against the schema itself:
|
Value |
Resolution |
|---|---|
null, or a nullable field is omitted
|
Encoded as the union’s null branch.
|
| A value matching exactly one non-null branch | Encoded as that branch. |
A value matching more than one non-null branch (for example a JSON number against ["int", "long"] or ["float", "double"])
|
Encoded using the first matching branch, in the order the branches are declared in the schema.
An integer that doesn’t fit in a 32-bit signed range is always encoded as long, even if int is declared first.
|
| A value that doesn’t match any branch of the union | The request is rejected with an error that includes the JSON path and the branches that were considered. |
This resolution applies at every level of the payload, including fields inside nested records, arrays, and maps. Logical types (for example timestamps, decimals, or UUIDs) are passed through unchanged once their union is resolved.
If a record field is omitted from the request body, the plugin falls back to the field’s schema default, if one exists.
Otherwise, the field must be nullable, or the plugin rejects the request as missing a required field.
Because an Avro-tagged value like {"string": "hello"} already matches a single branch by name, simple_json accepts it as-is.
This lets you migrate clients from avro_json to simple_json one at a time, instead of all at once.
For a sample configuration, see Simple JSON encoding for Avro schemas.
The Kafka Log plugin can compress message batches from a producer before sending them to the Kafka broker, using config.compression_type.
Compression reduces network bandwidth between Kong Gateway and the broker, broker disk usage, and cross-broker replication cost.
This applies only to the Kong Gateway-to-broker traffic flow.
It’s independent of any HTTP-level Content-Encoding between clients and Kong Gateway, and works the same way in both sync and async producer modes.
|
Codec |
Ratio |
Speed |
Notes |
|---|---|---|---|
none (default)
|
N/A | N/A | No compression. Preserves current behavior on upgrade. |
gzip
|
Highest | Slowest | Best when bandwidth or storage is the binding constraint and producer CPU is cheap. |
snappy
|
Moderate | Fast | A common default for throughput-sensitive Kafka workloads. |
lz4
|
Moderate | Fastest | Recommended codec for most workloads: similar ratio to Snappy, typically faster. |
Note:
zstdisn’t available as a producer-side codec, because it requires Kafka Produce API v7+ negotiation. The consume side can already decompresszstdbatches written by other producers.
Compression happens at the producer batch level. Kafka Log compresses the entire outgoing record batch as a single unit before it’s sent. Compression efficiency improves with batch size, so it’s most effective in async mode, where the plugin already accumulates messages before flushing to the broker. In sync mode, batches are typically smaller, but compression is still available and can be worthwhile for larger payloads.
If config.compression_type is misconfigured or compression fails at runtime, Kafka Log logs a warning and sends the batch uncompressed rather than dropping it.
For an example, see Compress log messages before sending to Kafka.
The Kafka Log plugin supports the following SASL authentication mechanisms for broker connections via authentication.mechanism:
|
Mechanism |
Description |
Example |
|---|---|---|
PLAIN
|
Authenticates using a username and password.
Set authentication.strategy to sasl and provide authentication.user and authentication.password.
|
Plain authentication |
SCRAM-SHA-256
|
Authenticates using a username and password with SCRAM-SHA-256 hashing.
Set authentication.strategy to sasl and provide authentication.user and authentication.password.
|
SCRAM-SHA-256 authentication |
SCRAM-SHA-512
|
Authenticates using a username and password with SCRAM-SHA-512 hashing.
Set authentication.strategy to sasl and provide authentication.user and authentication.password.
|
SCRAM-SHA-512 authentication |
OAUTHBEARER v3.15+
|
Authenticates using short-lived OAuth 2.0 access tokens fetched automatically by Kong Gateway.
Kong Gateway uses the client_credentials grant to retrieve tokens from the configured authentication.oauthbearer.token_endpoint_url, caches them until expiry, and presents them in the SASL/OAUTHBEARER handshake.
Requires the authentication.oauthbearer block.
|
SASL/OAUTHBEARER authentication |