Cloud
beta

redpanda_migrator

Unified Kafka consumer for migrating data between Kafka/Redpanda clusters. Use this input with the redpanda_migrator output to safely transfer topic data, ACLs, schemas, and consumer group offsets between clusters. This component is designed for migration scenarios.

  • Common

  • Advanced

input:
  label: ""
  redpanda_migrator:
    seed_brokers: [] # No default (required)
    topics: [] # No default (optional)
    regexp_topics_include: [] # No default (optional)
    regexp_topics_exclude: [] # No default (optional)
    transaction_isolation_level: read_uncommitted
    consumer_group: "" # No default (optional)
    schema_registry:
      url: "" # No default (required)
      timeout: 5s
    auto_replay_nacks: true
input:
  label: ""
  redpanda_migrator:
    seed_brokers: [] # No default (required)
    client_id: redpanda-connect
    tls:
      enabled: false
      skip_cert_verify: false
      enable_renegotiation: false
      root_cas: ""
      root_cas_file: ""
      client_certs: []
    sasl: [] # No default (optional)
    metadata_max_age: 1m
    request_timeout_overhead: 10s
    conn_idle_timeout: 20s
    tcp:
      connect_timeout: 0s
      keep_alive:
        idle: 15s
        interval: 15s
        count: 9
      tcp_user_timeout: 0s
    topics: [] # No default (optional)
    regexp_topics_include: [] # No default (optional)
    regexp_topics_exclude: [] # No default (optional)
    rack_id: ""
    instance_id: ""
    rebalance_timeout: 45s
    session_timeout: 1m
    heartbeat_interval: 3s
    start_offset: earliest
    fetch_max_bytes: 50MiB
    fetch_max_wait: 5s
    fetch_min_bytes: 1B
    fetch_max_partition_bytes: 1MiB
    transaction_isolation_level: read_uncommitted
    consumer_group: "" # No default (optional)
    commit_period: 5s
    partition_buffer_bytes: 1MB
    topic_lag_refresh_period: 5s
    max_yield_batch_bytes: 32KB
    schema_registry:
      url: "" # No default (required)
      timeout: 5s
      tls:
        enabled: false
        skip_cert_verify: false
        enable_renegotiation: false
        root_cas: ""
        root_cas_file: ""
        client_certs: []
      oauth:
        enabled: false
        consumer_key: ""
        consumer_secret: ""
        access_token: ""
        access_token_secret: ""
      basic_auth:
        enabled: false
        username: ""
        password: ""
      jwt:
        enabled: false
        private_key_file: ""
        signing_method: ""
        claims: {}
        headers: {}
    auto_replay_nacks: true

The redpanda_migrator input:

  • Reads a batch of messages from a broker.

  • Waits for the redpanda_migrator output to acknowledge the writes before updating the Kafka consumer group offset.

  • Provides the same delivery guarantees and ordering semantics as the redpanda input.

Specify a consumer group to make this input consume one or more topics and automatically balance the topic partitions across any other connected clients with the same consumer group. Otherwise, topics are consumed in their entirety or with explicit partitions.

This input requires a corresponding redpanda_migrator output in the same pipeline. Each pipeline must have both input and output components configured. For capabilities, guarantees, scheduling, and examples, see the output documentation.

Requirements

  • Must be paired with a redpanda_migrator output in the same pipeline.

  • Requires access to a source Kafka or Redpanda cluster.

  • Consumer group configuration is recommended for partition balancing.

  • When the source cluster enforces ACLs, a consumer ACL alone is not enough for the source principal: it needs at minimum topic READ and DESCRIBE_CONFIGS, plus consumer group and cluster permissions. A READ ACL grants DESCRIBE but not DESCRIBE_CONFIGS, so the migrator consumes messages but fails to create topics with TOPIC_AUTHORIZATION_FAILED. See Required permissions.

Multiple migrator pairs

When using multiple migrator pairs in a single pipeline, coordination is based on the label field. The label of the input and output must match exactly for correct pairing. If labels do not match, migration fails for that pair.

Performance tuning for high throughput

For workloads with high message rates or large messages, adjust the following settings to optimize throughput:

On this input component:

  • partition_buffer_bytes: Set to 2MB to increase per-partition buffer size

  • max_yield_batch_bytes: Set to 1MB to allow larger batches to be yielded

On the paired redpanda_migrator output component:

  • max_in_flight: Set to the total number of partitions being copied in parallel (up to all partitions in the cluster)

Setting max_yield_batch_bytes over 1MB is counter-productive unless you change the broker settings to allow bigger messages or batches. The partition_buffer_bytes setting allows for partition readahead.

Metrics

This input emits an input_redpanda_migrator_lag metric with topic and partition labels for each consumed topic. This metric records the number of produced messages that remain to be read from each topic/partition pair by the specified consumer group. Monitor this metric to track migration progress and detect bottlenecks.

Metadata

This input adds the following metadata fields to each message:

  • kafka_key

  • kafka_topic

  • kafka_partition

  • kafka_offset

  • kafka_lag

  • kafka_timestamp_ms

  • kafka_timestamp_unix

  • All record headers

Fields

auto_replay_nacks

Whether to automatically replay rejected messages (negative acknowledgements, or nacks) at the output level. If the cause of rejections persists, leaving this option enabled can result in back pressure.

Set auto_replay_nacks to false to delete rejected messages. Disabling auto replays can greatly improve memory efficiency of high throughput streams, as the original shape of the data is discarded immediately upon consumption and mutation.

Type: bool

Default: true

client_id

An identifier for the client connection. This identifier appears in broker logs and metrics, which helps you identify the Redpanda Connect instance that is connecting.

Type: string

Default: redpanda-connect

commit_period

The period of time between each commit of the current partition offsets to the consumer group. Offsets are always committed during shutdown.

Type: string

Default: 5s

conn_idle_timeout

The approximate amount of time that connections can remain idle before they are closed. In the worst case, a connection can stay idle for up to twice this value. This field accepts Go duration format strings such as 100ms, 1s, or 5s.

Type: string

Default: 20s

consumer_group

An optional consumer group. When you specify this value:

  • The partitions of any topics specified in the topics field are automatically distributed across consumers that share the consumer group.

  • Partition offsets are automatically committed and resumed under this name.

Consumer groups are not supported when you specify explicit partitions to consume from in the topics field.

Type: string

fetch_max_bytes

The maximum number of bytes that a broker tries to send during a fetch.

If individual records are larger than the fetch_max_bytes value, brokers still send them.

This field is equivalent to the Java setting fetch.max.bytes.

Type: string

Default: 50MiB

fetch_max_partition_bytes

The maximum number of bytes that are consumed from a single partition in a fetch request. This field is equivalent to the Java setting fetch.max.partition.bytes.

If a single batch is larger than the fetch_max_partition_bytes value, the batch is still sent so that the client can make progress.

Type: string

Default: 1MiB

fetch_max_wait

The maximum period of time a broker can wait for a fetch response to reach the required minimum number of bytes (fetch_min_bytes). This field is equivalent to the Java setting fetch.max.wait.ms.

Type: string

Default: 5s

fetch_min_bytes

The minimum number of bytes that a broker tries to send during a fetch. This field is equivalent to the Java setting fetch.min.bytes.

Type: string

Default: 1B

heartbeat_interval

When you specify a consumer_group, heartbeat_interval sets how frequently a consumer group member should send heartbeats to Apache Kafka. Apache Kafka uses heartbeats to make sure that a group member’s session is active.

This value must be lower than session_timeout, and should be no higher than one-third of session_timeout.

This field is equivalent to the Java heartbeat.interval.ms setting and accepts Go duration format strings such as 10s or 2m.

Type: string

Default: 3s

instance_id

When you specify a consumer_group, set instance_id to define the group’s static membership, which can prevent unnecessary rebalances during reconnections. The value must be unique per consumer within the group.

When you assign an instance ID, the client does not leave the consumer group when it closes. To remove the client from the group, you must use an external admin command on behalf of the instance ID.

Type: string

Default: ""

max_yield_batch_bytes

The maximum size (in bytes) for each batch yielded by this input. This value must be less than or equal to the partition_buffer_bytes. If using Redpanda output, this value should not be greater than the max_message_bytes option value (1MB by default), and for high-throughput scenarios they should be equal.

Type: string

Default: 32KB

metadata_max_age

The maximum period of time after which metadata is refreshed. This field accepts Go duration format strings such as 100ms, 1s, or 5s.

Lower values provide more responsive topic and partition discovery but may increase broker load. Higher values reduce broker queries but can delay detection of topology changes.

This interval also controls how frequently regex topic patterns are re-evaluated to discover new matching topics.

Type: string

Default: 1m

partition_buffer_bytes

A buffer size (in bytes) for each consumed partition, which allows the internal queuing of records before they are flushed. Increasing this value may improve throughput but results in higher memory utilization.

Each buffer can grow slightly beyond this value.

Type: string

Default: 1MB

rack_id

A rack specifies where the client is physically located, and changes fetch requests to consume from the closest replica as opposed to the leader replica.

Type: string

Default: ""

rebalance_timeout

When you specify a consumer_group, rebalance_timeout sets how long consumer group members can take to complete their work and commit offsets after a rebalance has begun. The time that a member takes to detect the rebalance (from a heartbeat) counts against this timeout. This field accepts Go duration format strings such as 100ms, 1s, or 5s.

Type: string

Default: 45s

regexp_topics_exclude[]

A list of regular expression patterns for excluding topics when regex mode is enabled (using regexp_topics_include or the deprecated regexp_topics boolean). Topics matching any of these patterns will be excluded from consumption, even if they match include patterns.

Each pattern is a full regular expression evaluated against the complete topic name. Patterns are not anchored by default, so use ^ and $ for exact matching. Exclude patterns are applied after include patterns, providing fine-grained control over topic selection.

Example: regexp_topics_exclude: ["^_", ".-temp$", ".-test.*"] excludes topics starting with underscore, ending with -temp, or containing -test.

Type: array<string>

regexp_topics_include[]

A list of regular expression patterns for matching topics to consume from. When specified, the client will periodically refresh the list of matching topics based on the metadata_max_age interval.

Each pattern is a full regular expression evaluated against the complete topic name. Patterns are not anchored by default, so logs_. matches my-logs_events and logs_errors. Use ^logs_.$ to match only topics starting with logs_.

This field enables regex mode (replacing the deprecated regexp_topics boolean) and cannot be used together with explicit topics lists. Use regexp_topics_exclude to filter out specific patterns from the matched topics.

Example: regexp_topics_include: ["events_.", "logs_."] consumes from all topics starting with events_ or logs_.

Type: array<string>

# Examples:
regexp_topics_include:
  - logs_.*
  - metrics_.*

# ---

regexp_topics_include:
  - "events_[0-9]+"

request_timeout_overhead

Additional time to apply as overhead when calculating request deadlines. For most requests, the deadline is this overhead alone. For requests that define their own timeout field, the overhead is added on top of that timeout, which helps prevent premature timeouts.

This field is roughly equivalent to Apache Kafka’s request.timeout.ms parameter, but grants extra time to requests that have timeout fields.

Type: string

Default: 10s

sasl[]

Specify one or more methods or mechanisms of SASL authentication. They are tried in order. If the broker supports the first SASL mechanism, all connections use it. If the first mechanism fails, the client picks the first supported mechanism. If the broker does not support any client mechanisms, all connections fail.

Type: array<object>

# Examples:
sasl:
  - mechanism: SCRAM-SHA-512
    password: bar
    username: foo

sasl[].aws

Contains AWS-specific fields for when sasl.mechanism is set to AWS_MSK_IAM.

Type: object

sasl[].aws.credentials

Manually configure the AWS credentials to use (optional). For more information, see the Amazon Web Services guide.

Type: object

sasl[].aws.credentials.from_ec2_role

Use the credentials of a host EC2 machine configured to assume an IAM role associated with the instance.

Type: bool

sasl[].aws.credentials.id

The ID of the AWS credentials to use.

Type: string

sasl[].aws.credentials.profile

The profile from ~/.aws/credentials to use.

Type: string

sasl[].aws.credentials.role

The ARN of the role to assume.

Type: string

sasl[].aws.credentials.role_external_id

An external ID to use when assuming a role.

Type: string

sasl[].aws.credentials.secret

The secret for the AWS credentials in use.

This field contains sensitive information that usually shouldn’t be added to a configuration directly. For more information, see Manage Secrets before adding it to your configuration.

Type: string

sasl[].aws.credentials.token

The token for the AWS credentials in use. Required only when using short-term credentials.

Type: string

sasl[].aws.endpoint

A custom endpoint URL for AWS API requests. Use this to connect to AWS-compatible services or local testing environments instead of the standard AWS endpoints.

Type: string

sasl[].aws.region

The AWS region in which your resources are hosted.

Type: string

sasl[].aws.tcp

Configure TCP socket-level settings to optimize network performance and reliability. These low-level controls are useful for:

  • Unresponsive hosts: Set connect_timeout to limit how long a connection attempt can take (the default 0s sets no limit)

  • Long-lived connections: Configure keep_alive settings to detect and recover from stale connections

  • Unstable networks: Tune keep-alive probes to balance between quick failure detection and avoiding false positives

  • Linux systems with specific requirements: Use tcp_user_timeout (Linux 2.6.37+) to control data acknowledgment timeouts

Most users should keep the default values. Only modify these settings if you’re experiencing connection stability issues or have specific network requirements.

Type: object

sasl[].aws.tcp.connect_timeout

Maximum amount of time a dial will wait for a connect to complete. Zero disables.

Type: string

Default: 0s

sasl[].aws.tcp.keep_alive

TCP keep-alive probe configuration.

Type: object

sasl[].aws.tcp.keep_alive.count

Maximum unanswered keep-alive probes before dropping the connection. Zero defaults to 9.

Type: int

Default: 9

sasl[].aws.tcp.keep_alive.idle

Duration the connection must be idle before sending the first keep-alive probe. Zero defaults to 15s. Negative values disable keep-alive probes.

Type: string

Default: 15s

sasl[].aws.tcp.keep_alive.interval

Duration between keep-alive probes. Zero defaults to 15s.

Type: string

Default: 15s

sasl[].aws.tcp.tcp_user_timeout

Maximum time to wait for acknowledgment of transmitted data before killing the connection. Linux-only (kernel 2.6.37+), ignored on other platforms. When enabled, keep_alive.idle must be greater than this value per RFC 5482. Zero disables.

Type: string

Default: 0s

sasl[].extensions

Key/value pairs to add to OAUTHBEARER authentication requests.

Type: object<string>

sasl[].mechanism

The SASL mechanism to use for authentication.

Type: string

Option Summary

AWS_MSK_IAM

AWS IAM based authentication as specified by the 'aws-msk-iam-auth' java library.

OAUTHBEARER

OAuth Bearer based authentication.

PLAIN

Plain text authentication.

REDPANDA_CLOUD_SERVICE_ACCOUNT

Redpanda Cloud Service Account authentication when running in Redpanda Cloud.

SCRAM-SHA-256

SCRAM based authentication as specified in RFC5802.

SCRAM-SHA-512

SCRAM based authentication as specified in RFC5802.

none

Disable sasl authentication

sasl[].password

The password to use for PLAIN or SCRAM-* authentication.

This field contains sensitive information that usually shouldn’t be added to a configuration directly. For more information, see Manage Secrets before adding it to your configuration.

Type: string

Default: ""

sasl[].token

The token to use for a single session’s OAUTHBEARER authentication.

Type: string

Default: ""

sasl[].username

The username to use for PLAIN or SCRAM-* authentication.

Type: string

Default: ""

schema_registry

Configuration for schema registry integration. Enables migration of schema subjects, versions, and compatibility settings between clusters.

Type: object

schema_registry.basic_auth

Configure basic authentication for requests from this component.

Type: object

schema_registry.basic_auth.enabled

Whether to use basic authentication in requests.

Type: bool

Default: false

schema_registry.basic_auth.password

The password to use for authentication. Used together with username for basic authentication.

This field contains sensitive information that usually shouldn’t be added to a configuration directly. For more information, see Manage Secrets before adding it to your configuration.

Type: string

Default: ""

schema_registry.basic_auth.username

The username of the account credentials to authenticate as. Used together with password for basic authentication.

Type: string

Default: ""

schema_registry.jwt

Beta

Configure JSON Web Token (JWT) authentication. This feature is in beta and may change in future releases. JWTs provide secure, stateless authentication between services.

Type: object

schema_registry.jwt.claims

A map of claims to include in the JWT. Claims pass the identity of the authenticated entity to the service provider.

Type: object

Default: {}

schema_registry.jwt.enabled

Whether to use JWT authentication in requests.

Type: bool

Default: false

schema_registry.jwt.headers

Additional key-value pairs to include in the JWT header (optional). These headers provide extra metadata for JWT processing.

Type: object

Default: {}

schema_registry.jwt.private_key_file

Path to a file containing the PEM-encoded private key using PKCS#1 or PKCS#8 format. The private key must be compatible with the algorithm specified in the signing_method field.

Type: string

Default: ""

schema_registry.jwt.signing_method

The cryptographic algorithm used to sign the JWT. Supported algorithms are RS256, RS384, RS512, and EdDSA. This algorithm must be compatible with the private key specified in the private_key_file field.

Type: string

Default: ""

schema_registry.oauth

Configure OAuth version 1.0 authentication for secure API access.

Type: object

schema_registry.oauth.access_token

The value used to gain access to the protected resources on behalf of the user.

Type: string

Default: ""

schema_registry.oauth.access_token_secret

The secret that establishes ownership of the access_token in OAuth 1.0 authentication.

This field contains sensitive information that usually shouldn’t be added to a configuration directly. For more information, see Manage Secrets before adding it to your configuration.

Type: string

Default: ""

schema_registry.oauth.consumer_key

The value used to identify this component or client to the service provider.

Type: string

Default: ""

schema_registry.oauth.consumer_secret

The secret that establishes ownership of the consumer key in OAuth 1.0 authentication.

This field contains sensitive information that usually shouldn’t be added to a configuration directly. For more information, see Manage Secrets before adding it to your configuration.

Type: string

Default: ""

schema_registry.oauth.enabled

Whether to enable OAuth version 1.0 authentication for requests.

Type: bool

Default: false

schema_registry.timeout

HTTP client timeout for schema registry requests.

Type: string

Default: 5s

schema_registry.tls

Configure Transport Layer Security (TLS) settings to secure network connections. This includes options for standard TLS as well as mutual TLS (mTLS) authentication where both client and server authenticate each other using certificates. Key configuration options include enabled to enable TLS, client_certs for mTLS authentication, root_cas/root_cas_file for custom certificate authorities, and skip_cert_verify for development environments.

Type: object

schema_registry.tls.client_certs[]

A list of client certificates for mutual TLS (mTLS) authentication. Configure this field to enable mTLS, authenticating the client to the server with these certificates.

You must set tls.enabled: true for the client certificates to take effect.

Certificate pairing rules: For each certificate item, provide either:

  • Inline PEM data using both cert and key or

  • File paths using both cert_file and key_file.

Mixing inline and file-based values within the same item is not supported.

Type: array<object>

Default: []

# Examples:
client_certs:
  - cert: foo
    key: bar


# ---

client_certs:
  - cert_file: ./example.pem
    key_file: ./example.key

schema_registry.tls.client_certs[].cert

The plaintext certificate to use for TLS authentication. Must be paired with the corresponding private key in the key field when using inline PEM data for mTLS client certificates.

Type: string

Default: ""

schema_registry.tls.client_certs[].cert_file

The path to a file containing the certificate to use for TLS authentication. Must be paired with the corresponding private key file in the key_file field when using file-based configuration for mTLS client certificates.

Type: string

Default: ""

schema_registry.tls.client_certs[].key

Private key for mTLS client certificate as inline PEM data. Must correspond to the client certificate specified in the cert field. Use this field together with cert when providing certificate data inline rather than through files.

This field contains sensitive information that usually shouldn’t be added to a configuration directly. For more information, see Manage Secrets before adding it to your configuration.

Type: string

Default: ""

schema_registry.tls.client_certs[].key_file

Path to private key file for mTLS client certificate in PEM format. Must correspond to the client certificate specified in the cert_file field. Use this field together with cert_file when loading certificate data from files.

Type: string

Default: ""

schema_registry.tls.client_certs[].password

The password to use for the private key (specified in the key or key_file fields), if it is password-protected. The PKCS#1 and PKCS#8 formats are supported. Supports environment variable interpolation for secure password management.

The pbeWithMD5AndDES-CBC algorithm is obsolete and not supported for the PKCS#8 format. This algorithm does not authenticate the ciphertext, making it vulnerable to padding oracle attacks that can let an attacker recover the plaintext.

This field contains sensitive information that usually shouldn’t be added to a configuration directly. For more information, see Manage Secrets before adding it to your configuration.

Type: string

Default: ""

# Examples:
password: foo

# ---

password: ${KEY_PASSWORD}

schema_registry.tls.enable_renegotiation

Whether to allow the remote server to repeatedly request renegotiation. Enable this option if you’re seeing the error message local error: tls: no renegotiation.

Type: bool

Default: false

schema_registry.tls.enabled

Whether to enable TLS for secure connections. Set to true to enable TLS encryption. Required to be true for other TLS options (like client_certs, root_cas, etc.) to take effect.

Type: bool

Default: false

schema_registry.tls.root_cas

Specify a root certificate authority to use (optional). This is a string that represents a certificate chain from the parent-trusted root certificate, through possible intermediate signing certificates, to the host certificate. Use either this field for inline certificate data or root_cas_file for file-based certificate loading.

This field contains sensitive information that usually shouldn’t be added to a configuration directly. For more information, see Manage Secrets before adding it to your configuration.

Type: string

Default: ""

# Examples:
root_cas: |-
  -----BEGIN CERTIFICATE-----
  ...
  -----END CERTIFICATE-----

schema_registry.tls.root_cas_file

Specify the path to a root certificate authority file (optional). This is a file, often with a .pem extension, which contains a certificate chain from the parent-trusted root certificate, through possible intermediate signing certificates, to the host certificate. Use either this field for file-based certificate loading or root_cas for inline certificate data.

Type: string

Default: ""

# Examples:
root_cas_file: ./root_cas.pem

schema_registry.tls.skip_cert_verify

Whether to skip server-side certificate verification. Set to true only for testing environments as this reduces security by disabling certificate validation. When using self-signed certificates or in development, this may be necessary, but should never be used in production. Consider using root_cas or root_cas_file to specify trusted certificates instead of disabling verification entirely.

Type: bool

Default: false

schema_registry.url

The base URL of the schema registry service. Required for schema migration functionality.

Type: string

# Examples:
url: http://localhost:8081

# ---

url: https://schema-registry.example.com:8081

seed_brokers[]

A list of broker addresses used to establish connections. If an item of the list contains commas, it is expanded into multiple addresses.

Type: array<string>

# Examples:
seed_brokers:
  - "localhost:9092"

# ---

seed_brokers:
  - "foo:9092"
  - "bar:9092"

# ---

seed_brokers:
  - "foo:9092,bar:9092"

session_timeout

When you specify a consumer_group, session_timeout sets the maximum interval between heartbeats sent by a consumer group member to the broker. If a broker doesn’t receive a heartbeat from a group member before the timeout expires, it removes the member from the consumer group and initiates a rebalance. This field accepts Go duration format strings such as 100ms, 1s, or 5s.

Type: string

Default: 1m

start_offset

Specify the offset from which this input starts or restarts consuming messages. Restarts occur when the OffsetOutOfRange error is seen during a fetch.

Type: string

Default: earliest

Option Summary

committed

Prevents consuming a partition in a group if the partition has no prior commits. Corresponds to Kafka’s auto.offset.reset=none option

earliest

Start from the earliest offset. Corresponds to Kafka’s auto.offset.reset=earliest option.

latest

Start from the latest offset. Corresponds to Kafka’s auto.offset.reset=latest option.

tcp

Configure TCP socket-level settings to optimize network performance and reliability. These low-level controls are useful for:

  • Unresponsive hosts: Set connect_timeout to limit how long a connection attempt can take (the default 0s sets no limit)

  • Long-lived connections: Configure keep_alive settings to detect and recover from stale connections

  • Unstable networks: Tune keep-alive probes to balance between quick failure detection and avoiding false positives

  • Linux systems with specific requirements: Use tcp_user_timeout (Linux 2.6.37+) to control data acknowledgment timeouts

Most users should keep the default values. Only modify these settings if you’re experiencing connection stability issues or have specific network requirements.

Type: object

tcp.connect_timeout

Maximum amount of time a dial will wait for a connect to complete. Zero disables.

Type: string

Default: 0s

tcp.keep_alive

TCP keep-alive probe configuration.

Type: object

tcp.keep_alive.count

Maximum unanswered keep-alive probes before dropping the connection. Zero defaults to 9.

Type: int

Default: 9

tcp.keep_alive.idle

Duration the connection must be idle before sending the first keep-alive probe. Zero defaults to 15s. Negative values disable keep-alive probes.

Type: string

Default: 15s

tcp.keep_alive.interval

Duration between keep-alive probes. Zero defaults to 15s.

Type: string

Default: 15s

tcp.tcp_user_timeout

Maximum time to wait for acknowledgment of transmitted data before killing the connection. Linux-only (kernel 2.6.37+), ignored on other platforms. When enabled, keep_alive.idle must be greater than this value per RFC 5482. Zero disables.

Type: string

Default: 0s

tls

Configure Transport Layer Security (TLS) settings to secure network connections. This includes options for standard TLS as well as mutual TLS (mTLS) authentication where both client and server authenticate each other using certificates. Key configuration options include enabled to enable TLS, client_certs for mTLS authentication, root_cas/root_cas_file for custom certificate authorities, and skip_cert_verify for development environments.

Type: object

tls.client_certs[]

A list of client certificates for mutual TLS (mTLS) authentication. Configure this field to enable mTLS, authenticating the client to the server with these certificates.

You must set tls.enabled: true for the client certificates to take effect.

Certificate pairing rules: For each certificate item, provide either:

  • Inline PEM data using both cert and key or

  • File paths using both cert_file and key_file.

Mixing inline and file-based values within the same item is not supported.

Type: array<object>

Default: []

# Examples:
client_certs:
  - cert: foo
    key: bar


# ---

client_certs:
  - cert_file: ./example.pem
    key_file: ./example.key

tls.client_certs[].cert

The plaintext certificate to use for TLS authentication. Must be paired with the corresponding private key in the key field when using inline PEM data for mTLS client certificates.

Type: string

Default: ""

tls.client_certs[].cert_file

The path to a file containing the certificate to use for TLS authentication. Must be paired with the corresponding private key file in the key_file field when using file-based configuration for mTLS client certificates.

Type: string

Default: ""

tls.client_certs[].key

Private key for mTLS client certificate as inline PEM data. Must correspond to the client certificate specified in the cert field. Use this field together with cert when providing certificate data inline rather than through files.

This field contains sensitive information that usually shouldn’t be added to a configuration directly. For more information, see Manage Secrets before adding it to your configuration.

Type: string

Default: ""

tls.client_certs[].key_file

Path to private key file for mTLS client certificate in PEM format. Must correspond to the client certificate specified in the cert_file field. Use this field together with cert_file when loading certificate data from files.

Type: string

Default: ""

tls.client_certs[].password

The password to use for the private key (specified in the key or key_file fields), if it is password-protected. The PKCS#1 and PKCS#8 formats are supported. Supports environment variable interpolation for secure password management.

The pbeWithMD5AndDES-CBC algorithm is obsolete and not supported for the PKCS#8 format. This algorithm does not authenticate the ciphertext, making it vulnerable to padding oracle attacks that can let an attacker recover the plaintext.

This field contains sensitive information that usually shouldn’t be added to a configuration directly. For more information, see Manage Secrets before adding it to your configuration.

Type: string

Default: ""

# Examples:
password: foo

# ---

password: ${KEY_PASSWORD}

tls.enable_renegotiation

Whether to allow the remote server to repeatedly request renegotiation. Enable this option if you’re seeing the error message local error: tls: no renegotiation.

Type: bool

Default: false

tls.enabled

Whether to enable TLS for secure connections. Set to true to enable TLS encryption. Required to be true for other TLS options (like client_certs, root_cas, etc.) to take effect.

Type: bool

Default: false

tls.root_cas

Specify a root certificate authority to use (optional). This is a string that represents a certificate chain from the parent-trusted root certificate, through possible intermediate signing certificates, to the host certificate. Use either this field for inline certificate data or root_cas_file for file-based certificate loading.

This field contains sensitive information that usually shouldn’t be added to a configuration directly. For more information, see Manage Secrets before adding it to your configuration.

Type: string

Default: ""

# Examples:
root_cas: |-
  -----BEGIN CERTIFICATE-----
  ...
  -----END CERTIFICATE-----

tls.root_cas_file

Specify the path to a root certificate authority file (optional). This is a file, often with a .pem extension, which contains a certificate chain from the parent-trusted root certificate, through possible intermediate signing certificates, to the host certificate. Use either this field for file-based certificate loading or root_cas for inline certificate data.

Type: string

Default: ""

# Examples:
root_cas_file: ./root_cas.pem

tls.skip_cert_verify

Whether to skip server-side certificate verification. Set to true only for testing environments as this reduces security by disabling certificate validation. When using self-signed certificates or in development, this may be necessary, but should never be used in production. Consider using root_cas or root_cas_file to specify trusted certificates instead of disabling verification entirely.

Type: bool

Default: false

topic_lag_refresh_period

The interval between consumer lag refreshes. During each cycle, this input asks the brokers for the consumer group’s committed offsets and the partition end offsets, and records the difference (the number of unread messages) for each topic partition in the consumer lag metric and the kafka_lag metadata field. Lag is only refreshed when consumer_group is set. This field accepts Go duration format strings such as 100ms, 1s, or 5s.

Type: string

Default: 5s

topics[]

A list of topics to consume from. You can list multiple comma-separated topics in a single element.

If you specify a consumer_group, partitions are automatically distributed across consumers of a topic. Otherwise, all partitions are consumed.

Alternatively, add a colon after the topic name to set the explicit partitions to consume. For example, foo:0 consumes the partition 0 of the topic foo. This syntax also supports ranges. For example, foo:0-10 consumes all partitions from 0 through to 10 inclusive.

Finally, add another colon after the partition to set an explicit offset to consume from. For example, foo:0:10 consumes the partition 0 of the topic foo starting from the offset 10. If the offset is not present (or remains unspecified) then the field start_offset determines which offset to start from.

Type: array<string>

# Examples:
topics:
  - foo
  - bar

# ---

topics:
  - things.*

# ---

topics:
  - "foo,bar"

# ---

topics:
  - "foo:0"
  - "bar:1"
  - "bar:3"

# ---

topics:
  - "foo:0,bar:1,bar:3"

# ---

topics:
  - "foo:0-5"

transaction_isolation_level

The isolation level for handling transactional messages. This setting determines how transactions are processed and affects data consistency guarantees.

Type: string

Default: read_uncommitted

Option Summary

read_committed

If set, only committed transactional records are processed.

read_uncommitted

If set, then uncommitted records are processed.

Troubleshooting

  • Ensure the input and output label fields match exactly.

  • Both input and output must be present in the pipeline.

  • Verify consumer group configuration for partition balancing.

  • Monitor the lag metric for stalled migration.