> For the complete documentation index, see [llms.txt](https://docs.fluentbit.io/manual/llms.txt). Markdown versions of documentation pages are available by appending `.md` to page URLs; this page is available as [Markdown](https://docs.fluentbit.io/manual/data-pipeline/outputs/kafka.md).

# Kafka Producer

The *Kafka Producer* output plugin lets you ingest your records into an [Apache Kafka](https://kafka.apache.org/) service. This plugin uses the official [librdkafka C library](https://github.com/confluentinc/librdkafka).

In Fluent Bit 4.0.4 and later, the Kafka input plugin supports authentication with AWS MSK IAM, enabling integration with Amazon MSK (Managed Streaming for Apache Kafka) clusters that require IAM-based access.

## Configuration parameters

This plugin supports the following parameters:

| Key                               | Description                                                                                                                                                                                                                                                                                                                    | Default      |
| --------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ | ------------ |
| `aws_msk_iam`                     | Enable AWS MSK IAM authentication. Requires Fluent Bit 4.0.4 or later.                                                                                                                                                                                                                                                         | `false`      |
| `aws_msk_iam_cluster_arn`         | Full ARN of the MSK cluster used for region extraction. Required when `aws_msk_iam` is enabled.                                                                                                                                                                                                                                | *none*       |
| `brokers`                         | Single or multiple list of Kafka brokers. For example, `192.168.1.3:9092`, `192.168.1.4:9092`.                                                                                                                                                                                                                                 | *none*       |
| `client_id`                       | Client ID to use when connecting to Kafka.                                                                                                                                                                                                                                                                                     | *none*       |
| `dynamic_topic`                   | Adds unknown topics (found in `topic_key`) to `topics`. Only a default topic needs to be configured in `topics`.                                                                                                                                                                                                               | `false`      |
| `format`                          | Specify data format. Available formats: `avro` (requires Avro encoder build option), `gelf`, `json`, `msgpack`, `otlp_json` (supports logs, metrics, and traces events), `otlp_proto` (supports logs, metrics, and traces events), `protobuf` (requires Protobuf encoder build option, and Fluent Bit v5.1.3 or later), `raw`. | `json`       |
| `gelf_full_message_key`           | Key to use as the long message for GELF format output.                                                                                                                                                                                                                                                                         | *none*       |
| `gelf_host_key`                   | Key to use as the host for GELF format output.                                                                                                                                                                                                                                                                                 | *none*       |
| `gelf_level_key`                  | Key to use as the log level for GELF format output.                                                                                                                                                                                                                                                                            | *none*       |
| `gelf_short_message_key`          | Key to use as the short message for GELF format output.                                                                                                                                                                                                                                                                        | *none*       |
| `gelf_timestamp_key`              | Key to use as the timestamp for GELF format output.                                                                                                                                                                                                                                                                            | *none*       |
| `group_id`                        | Consumer group ID.                                                                                                                                                                                                                                                                                                             | *none*       |
| `message_key`                     | Optional key to store the message.                                                                                                                                                                                                                                                                                             | *none*       |
| `message_key_field`               | If set, the value of `message_key_field` in the record will indicate the message key. If not set or not found in the record, `message_key` is used if set.                                                                                                                                                                     | *none*       |
| `otlp_logs_partition_by_resource` | When using `otlp_json` or `otlp_proto` format for logs, send each OTLP resource's logs as a separate Kafka message.                                                                                                                                                                                                            | `false`      |
| `protobuf_message`                | Fully qualified name of the Protobuf message used to encode records. Optional when the registered root schema declares exactly one message, and required otherwise. Requires the Protobuf encoder build option. See [Protobuf support](#protobuf-support). Supported in v5.1.3 or later.                                       | *none*       |
| `queue_full_retries`              | Number of local retries to enqueue data when the `rdkafka` queue is full. The interval between retries is 1 second. Set to `0` for unlimited retries.                                                                                                                                                                          | `10`         |
| `raw_log_key`                     | When using the `raw` format, the value of `raw_log_key` in the record is sent to Kafka as the payload.                                                                                                                                                                                                                         | *none*       |
| `rdkafka.{property}`              | `{property}` can be any [librdkafka property](https://github.com/confluentinc/librdkafka/blob/master/CONFIGURATION.md).                                                                                                                                                                                                        | *none*       |
| `schema_id`                       | Avro or Protobuf schema ID. Requires the Avro or Protobuf encoder build option.                                                                                                                                                                                                                                                | *none*       |
| `schema_registry_bearer_token`    | Bearer token used to authenticate to the Schema Registry. Also accepted as `schema.registry.bearer.token`.                                                                                                                                                                                                                     | *none*       |
| `schema_registry_framing`         | Wire format used to frame messages encoded with a registry-resolved schema. Only `cp1`, the Confluent wire format, is supported.                                                                                                                                                                                               | `cp1`        |
| `schema_registry_http_passwd`     | Password for Schema Registry HTTP basic authentication. Also accepted as `schema.registry.http.password`.                                                                                                                                                                                                                      | *none*       |
| `schema_registry_http_user`       | User for Schema Registry HTTP basic authentication. Also accepted as `schema.registry.http.user`.                                                                                                                                                                                                                              | *none*       |
| `schema_registry_subject`         | Schema Registry subject to resolve the Avro or Protobuf schema from. Also accepted as `schema.registry.subject`. See [Resolve schemas from a registry](#resolve-schemas-from-a-registry).                                                                                                                                      | *none*       |
| `schema_registry_url`             | Base URL of a Confluent Schema Registry, or a comma-separated list of URLs. Also accepted as `schema.registry.url`.                                                                                                                                                                                                            | *none*       |
| `schema_registry_version`         | Version of the subject to resolve. Also accepted as `schema.registry.version`.                                                                                                                                                                                                                                                 | `latest`     |
| `schema_str`                      | Inline Avro schema string, used instead of resolving a schema from a registry. Requires the Avro encoder build option. Not supported with `format` set to `protobuf`, which always resolves its schema from a registry.                                                                                                        | *none*       |
| `timestamp_format`                | Specify the timestamp format. Allowed values: `double`, `iso8601` (seconds precision), `iso8601_ns` (nanoseconds precision).                                                                                                                                                                                                   | `double`     |
| `timestamp_key`                   | Key to store the record timestamp.                                                                                                                                                                                                                                                                                             | `@timestamp` |
| `topic_key`                       | If multiple `topics` exist, the value of `topic_key` in the record indicates the topic to use. If the value isn't present in `topics`, the first topic in the list is used.                                                                                                                                                    | *none*       |
| `topics`                          | Single topic or comma-separated list of topics that Fluent Bit will use to send messages to Kafka. If multiple topics are set, the `topic_key` field in the record selects the topic.                                                                                                                                          | `fluent-bit` |
| `workers`                         | The number of [workers](/manual/administration/multithreading.md#outputs) to perform flush operations for this output.                                                                                                                                                                                                         | `0`          |

Setting `rdkafka.log.connection.close` to `false` and `rdkafka.request.required.acks` to `1` are examples of recommended settings of `librdfkafka` properties.

## Get started

To insert records into Apache Kafka, you can run the plugin from the command line or through the configuration file.

### Command line

The Kafka plugin can read parameters through the `-p` argument (property):

```shell
fluent-bit -i cpu -o kafka -p brokers=192.168.1.3:9092 -p topics=test
```

### Configuration file

In your main configuration file append the following:

{% tabs %}
{% tab title="fluent-bit.yaml" %}

```yaml
pipeline:
  inputs:
    - name: cpu

  outputs:
    - name: kafka
      match: '*'
      brokers: 192.168.1.3:9092
      topics: test
```

{% endtab %}

{% tab title="fluent-bit.conf" %}

```
[INPUT]
  Name  cpu

[OUTPUT]
  Name        kafka
  Match       *
  Brokers     192.168.1.3:9092
  Topics      test
```

{% endtab %}
{% endtabs %}

### Avro support

Fluent Bit can encode records as Avro messages for the `out_kafka` plugin, but this support isn't included in any official release packages or container images. To use it, you must build Fluent Bit from source.

Avro support is optional and must be activated at build time by using a build definition with `cmake`: `-DFLB_AVRO_ENCODER=On` such as in the following example which activates:

* `out_kafka` with Avro encoding
* Fluent Bit Prometheus
* Metrics using an embedded HTTP endpoint
* Debugging support
* Builds the test suites

```shell
cmake -DFLB_DEV=On -DFLB_OUT_KAFKA=On -DFLB_TLS=On -DFLB_TESTS_RUNTIME=On -DFLB_TESTS_INTERNAL=On -DCMAKE_BUILD_TYPE=Debug -DFLB_HTTP_SERVER=true -DFLB_AVRO_ENCODER=On ../
```

#### Kafka configuration file with Avro encoding

In this example, the Fluent Bit configuration tails Kubernetes logs, updates the log lines with Kubernetes metadata using the Kubernetes filter. It then sends the updated log lines to a Kafka broker encoded with a specific Avro schema.

{% tabs %}
{% tab title="AVRO enabled: fluent-bit.yaml" %}

```yaml
pipeline:
  inputs:
    - name: tail
      tag: kube.*
      alias: some-alias
      path: /logdir/*.log
      db: /dbdir/some.db
      skip_long_lines: on
      refresh_interval: 10
      parser: some-parser

  filters:
    - name: kubernetes
      match: 'kube.*'
      kube_url: https://some_kube_api:443
      kube_ca_file: /certs/ca.crt
      kube_token_file: /tokens/token
      kube_tag_prefix: kube.var.log.containers.
      merge_log: on
      merge_log_key: log_processed

  outputs:
    - name: kafka
      match: '*'
      brokers: 192.168.1.3:9092
      topics: test
      # AVRO support must be enabled for schema support
      schema_str:  '{"name":"avro_logging","type":"record","fields":[{"name":"timestamp","type":"string"},{"name":"stream","type":"string"},{"name":"log","type":"string"},{"name":"kubernetes","type":{"name":"krec","type":"record","fields":[{"name":"pod_name","type":"string"},{"name":"namespace_name","type":"string"},{"name":"pod_id","type":"string"},{"name":"labels","type":{"type":"map","values":"string"}},{"name":"annotations","type":{"type":"map","values":"string"}},{"name":"host","type":"string"},{"name":"container_name","type":"string"},{"name":"docker_id","type":"string"},{"name":"container_hash","type":"string"},{"name":"container_image","type":"string"}]}},{"name":"cluster_name","type":"string"},{"name":"fabric","type":"string"}]}'
      schema_id: some_schema_id
      rdkafka.client.id: some_client_id
      rdkafka.debug: all
      rdkafka.enable.ssl.certificate.verification: true
      rdkafka.ssl.certificate.location: /certs/some.cert
      rdkafka.ssl.key.location: /certs/some.key
      rdkafka.ssl.ca.location: /certs/some-bundle.crt
      rdkafka.security.protocol: ssl
      rdkafka.request.required.acks: 1
      rdkafka.log.connection.close: false
      format: avro
      rdkafka.log_level: 7
      rdkafka.metadata.broker.list: 192.168.1.3:9092
```

{% endtab %}

{% tab title="AVRO enabled: fluent-bit.conf" %}

```
[INPUT]
  Name              tail
  Tag               kube.*
  Alias             some-alias
  Path              /logdir/*.log
  DB                /dbdir/some.db
  Skip_Long_Lines   On
  Refresh_Interval  10
  Parser            some-parser

[FILTER]
  Name                kubernetes
  Match               kube.*
  Kube_URL            https://some_kube_api:443
  Kube_CA_File        /certs/ca.crt
  Kube_Token_File     /tokens/token
  Kube_Tag_Prefix     kube.var.log.containers.
  Merge_Log           On
  Merge_Log_Key       log_processed

[OUTPUT]
  Name        kafka
  Match       *
  Brokers     192.168.1.3:9092
  Topics      test
  # AVRO support must be enabled for schema support
  Schema_Str  {"name":"avro_logging","type":"record","fields":[{"name":"timestamp","type":"string"},{"name":"stream","type":"string"},{"name":"log","type":"string"},{"name":"kubernetes","type":{"name":"krec","type":"record","fields":[{"name":"pod_name","type":"string"},{"name":"namespace_name","type":"string"},{"name":"pod_id","type":"string"},{"name":"labels","type":{"type":"map","values":"string"}},{"name":"annotations","type":{"type":"map","values":"string"}},{"name":"host","type":"string"},{"name":"container_name","type":"string"},{"name":"docker_id","type":"string"},{"name":"container_hash","type":"string"},{"name":"container_image","type":"string"}]}},{"name":"cluster_name","type":"string"},{"name":"fabric","type":"string"}]}
  Schema_Id some_schema_id
  rdkafka.client.id some_client_id
  rdkafka.debug All
  rdkafka.enable.ssl.certificate.verification true

  rdkafka.ssl.certificate.location /certs/some.cert
  rdkafka.ssl.key.location /certs/some.key
  rdkafka.ssl.ca.location /certs/some-bundle.crt
  rdkafka.security.protocol ssl
  rdkafka.request.required.acks 1
  rdkafka.log.connection.close false

  Format avro
  rdkafka.log_level 7
  rdkafka.metadata.broker.list 192.168.1.3:9092
```

{% endtab %}
{% endtabs %}

### Protobuf support

Protobuf support is available in Fluent Bit version 5.1.3 and greater. Fluent Bit can encode records as Protobuf messages for the `out_kafka` plugin, but this support isn't included in any official release packages or container images. To use it, you must build Fluent Bit from source.

Protobuf support is optional and must be activated at build time with `-DFLB_PROTOBUF_ENCODER=On`. The build requires Protobuf 3.12 or greater, including the `libprotoc` development libraries. In a build without it, setting `format` to `protobuf` fails at startup with `format protobuf requires FLB_PROTOBUF_ENCODER=On`.

Unlike Avro, Protobuf has no inline schema option. The schema always comes from a Confluent Schema Registry, so `schema_registry_url` is required. Without it, the plugin fails at startup with `format protobuf requires schema_registry_url`. See [Resolve schemas from a registry](#resolve-schemas-from-a-registry).

The schema resolved from the registry becomes the root `.proto` file. If that schema declares `references`, Fluent Bit fetches each referenced subject from the registry and registers it as an imported file, resolving references recursively. A schema graph is limited to 64 files and 4 MiB in total.

Use `protobuf_message` to select which message encodes your records:

* If the root schema declares exactly one message, you can leave `protobuf_message` unset.
* Otherwise, set it to the fully qualified message name, such as `com.example.logs.LogRecord`.

The message must be declared in the root schema itself, not in one of its imported files. If it isn't found, the plugin fails at startup with `cannot compile registered Protobuf schema: protobuf_message must select a message in the registered root schema`.

Fluent Bit converts each record to JSON and then encodes it with the selected message, so record field names must match the fields declared in the schema. Encoded messages use the Confluent wire format selected by `schema_registry_framing`, which prefixes the payload with the schema ID and the message index path.

The following example encodes records with a Protobuf schema resolved from a registry subject:

{% tabs %}
{% tab title="protobuf-fluent-bit.yaml" %}

```yaml
pipeline:
  outputs:
    - name: kafka
      match: '*'
      brokers: 192.168.1.3:9092
      topics: test
      format: protobuf
      schema_registry_url: 'https://registry-1:8081'
      schema_registry_subject: fluent-bit-logs-value
      protobuf_message: com.example.logs.LogRecord
```

{% endtab %}

{% tab title="protobuf-fluent-bit.conf" %}

```
[OUTPUT]
  Name                     kafka
  Match                    *
  Brokers                  192.168.1.3:9092
  Topics                   test
  Format                   protobuf
  Schema_Registry_Url      https://registry-1:8081
  Schema_Registry_Subject  fluent-bit-logs-value
  Protobuf_Message         com.example.logs.LogRecord
```

{% endtab %}
{% endtabs %}

### Resolve schemas from a registry

Resolving schemas from a Confluent Schema Registry is available in Fluent Bit version 5.1 and greater. It requires either the Avro encoder build option (`-DFLB_AVRO_ENCODER=On`) or, in version 5.1.3 and greater, the Protobuf encoder build option (`-DFLB_PROTOBUF_ENCODER=On`). Neither option is included in official release builds. See [Avro support](#avro-support) and [Protobuf support](#protobuf-support).

With `format` set to `avro`, the registry is an alternative to setting `schema_str` and `schema_id` in your configuration. With `format` set to `protobuf`, the registry is the only source of the schema. Set `schema_registry_url` to enable registry resolution. For Avro, if `schema_str` and `schema_id` are both set, Fluent Bit uses them and never contacts the registry.

Fluent Bit resolves the schema in one of two ways:

* If `schema_registry_subject` is set, it requests `/subjects/{subject}/versions/{version}`, where the version comes from `schema_registry_version` and defaults to `latest`.
* Otherwise, it requests `/schemas/ids/{id}` using the value of `schema_id`.

The schema is fetched once and reused for the lifetime of the plugin instance. If the registry can't be reached or returns an error, the chunk is retried.

To authenticate, set either `schema_registry_http_user` and `schema_registry_http_passwd` for HTTP basic authentication, or `schema_registry_bearer_token` for bearer token authentication.

For high availability, set `schema_registry_url` to a comma-separated list of registry URLs. During the initial schema fetch, Fluent Bit tries the next endpoint in the list when a request fails. Once the schema is successfully resolved, it's cached for the lifetime of the plugin instance, and the registry isn't contacted again.

Except for `schema_registry_framing`, each of these settings also accepts a dotted spelling that matches Confluent client configuration, such as `schema.registry.url` for `schema_registry_url`. The two spellings are interchangeable.

The following example resolves the latest version of the `fluent-bit-logs-value` subject from either of two registry endpoints:

{% tabs %}
{% tab title="avro-fluent-bit.yaml" %}

```yaml
pipeline:
  outputs:
    - name: kafka
      match: '*'
      brokers: 192.168.1.3:9092
      topics: test
      format: avro
      schema_registry_url: 'https://registry-1:8081,https://registry-2:8081'
      schema_registry_subject: fluent-bit-logs-value
      schema_registry_version: latest
      schema_registry_http_user: fluentbit
      schema_registry_http_passwd: ${SCHEMA_REGISTRY_PASSWORD}
```

{% endtab %}

{% tab title="avro-fluent-bit.conf" %}

```
[OUTPUT]
  Name                        kafka
  Match                       *
  Brokers                     192.168.1.3:9092
  Topics                      test
  Format                      avro
  Schema_Registry_Url         https://registry-1:8081,https://registry-2:8081
  Schema_Registry_Subject     fluent-bit-logs-value
  Schema_Registry_Version     latest
  Schema_Registry_Http_User   fluentbit
  Schema_Registry_Http_Passwd ${SCHEMA_REGISTRY_PASSWORD}
```

{% endtab %}
{% endtabs %}

### Kafka configuration file with `raw` format

This example Fluent Bit configuration file creates example records with the `payloadkey` and `msgkey` keys. The `msgkey` value is used as the Kafka message key, and the `payloadkey` value as the payload.

{% tabs %}
{% tab title="fluent-bit.yaml" %}

```yaml
pipeline:
  inputs:
    - name: dummy
      tag: example.data
      dummy: '{"payloadkey":"Data to send to kafka", "msgkey": "Key to use in the message"}'

  outputs:
    - name: kafka
      match: '*'
      brokers: 192.168.1.3:9092
      topics: test
      format: raw
      raw_log_key: payloadkey
      message_key_field: msgkey
```

{% endtab %}

{% tab title="fluent-bit.conf" %}

```
[INPUT]
  Name dummy
  Tag  example.data
  Dummy {"payloadkey":"Data to send to kafka", "msgkey": "Key to use in the message"}


[OUTPUT]
  Name        kafka
  Match       *
  Brokers     192.168.1.3:9092
  Topics      test
  Format      raw
  Raw_Log_Key       payloadkey
  Message_Key_Field msgkey
```

{% endtab %}
{% endtabs %}

## AWS MSK IAM authentication

Fluent Bit 4.0.4 and later supports authentication to Amazon MSK (Managed Streaming for Apache Kafka) clusters using AWS IAM for the Kafka output plugin. This lets you securely send data to MSK brokers with AWS credentials, leveraging IAM roles and policies for access control.

### Prerequisites

If you are compiling Fluent Bit from source, ensure the following requirements are met to enable AWS MSK IAM support:

* Build Requirements

| Platform        | Requirements                                                                                                                                                                        |
| --------------- | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
| **Linux/macOS** | The packages `libsasl2` and `libsasl2-dev` must be installed on your build environment.                                                                                             |
| **Windows**     | No additional SASL libraries required. Windows uses the built-in Security Support Provider Interface (SSPI) for SASL authentication, which only requires OpenSSL/TLS to be enabled. |

* Runtime Requirements:

  * Network Access: Fluent Bit must be able to reach your MSK broker endpoints (AWS VPC setup).
  * AWS Credentials: Provide credentials using any supported AWS method:
    * IAM roles (recommended for EC2, ECS, or EKS)
    * Environment variables (`AWS_ACCESS_KEY_ID`, `AWS_SECRET_ACCESS_KEY`)
    * AWS credentials file (`~/.aws/credentials`)
    * Instance metadata service (IMDS)

  These credentials are discovered by default when `aws_msk_iam` flag is enabled.
* IAM Permissions: The credentials must allow access to the target MSK cluster.

### AWS MSK IAM configuration parameters

See `aws_msk_iam` and `aws_msk_iam_cluster_arn` in the [configuration parameters](#configuration-parameters) table.

### Configuration example

{% tabs %}
{% tab title="fluent-bit.yaml" %}

```yaml
pipeline:
  inputs:
    - name: random

  outputs:
    - name: kafka
      match: '*'
      brokers: my-cluster.abcdef.c1.kafka.us-east-1.amazonaws.com:9098
      topics: my-topic
      aws_msk_iam: true
      aws_msk_iam_cluster_arn: arn:aws:kafka:us-east-1:123456789012:cluster/my-cluster/abcdef-1234-5678-9012-abcdefghijkl-s3
```

{% endtab %}
{% endtabs %}

### AWS IAM policy

IAM policies and permissions can be complex and can vary depending on your organization's security requirements. If you are unsure about the correct permissions or best practices, consult with your AWS administrator or an AWS expert who is familiar with MSK and IAM security.

The AWS credentials used by Fluent Bit must have permission to connect to your MSK cluster. Here is a minimal example policy:

```json
{
    "Version": "2012-10-17",
    "Statement": [
        {
            "Sid": "VisualEditor0",
            "Effect": "Allow",
            "Action": [
                "kafka-cluster:*",
                "kafka-cluster:DescribeCluster",
                "kafka-cluster:ReadData",
                "kafka-cluster:DescribeTopic",
                "kafka-cluster:Connect"
            ],
            "Resource": "*"
        }
    ]
}
```


---

# Agent Instructions
This documentation is published with GitBook. GitBook is the documentation platform designed so that both humans and AI agents can read, navigate, and reason over technical content effectively. Learn more at gitbook.com.

## Querying This Documentation
If you need additional information that is not directly available in this page, you can query the documentation dynamically by asking a question.

Perform an HTTP GET request on the current page URL with the `ask` query parameter, and the optional `goal` query parameter:

```
GET https://docs.fluentbit.io/manual/data-pipeline/outputs/kafka.md?ask=<question>&goal=<endgoal>
```

`ask` is the immediate question: it should be specific, self-contained, and written in natural language.
`goal` is optional and describes the broader end goal you are ultimately trying to accomplish on behalf of the user. GitBook uses it to tailor the answer towards what is most useful for that goal.

The response will contain a direct answer to the question and relevant excerpts and sources from the documentation.

Use this mechanism when the answer is not explicitly present in the current page, you need clarification or additional context, or you want to retrieve related documentation sections.
