mysql_cdc
Streams data changes from a MySQL database, using MySQL’s binary log to capture data updates.
This input is built on the mysql-canal library but uses a custom approach for streaming historical data.
Introduced in version 4.46.0.
-
Common
-
Advanced
input:
label: ""
mysql_cdc:
flavor: mysql
dsn: "" # No default (required)
tables: [] # No default (required)
checkpoint_cache: "" # No default (required)
checkpoint_key: mysql_binlog_position
snapshot_max_batch_size: 1000
stream_snapshot: false # No default (required)
max_parallel_snapshot_tables: 1
auto_replay_nacks: true
checkpoint_limit: 1024
batching:
count: 0
byte_size: 0
period: ""
check: ""
input:
label: ""
mysql_cdc:
flavor: mysql
dsn: "" # No default (required)
tables: [] # No default (required)
checkpoint_cache: "" # No default (required)
checkpoint_key: mysql_binlog_position
snapshot_max_batch_size: 1000
max_reconnect_attempts: 10
stream_snapshot: false # No default (required)
max_parallel_snapshot_tables: 1
auto_replay_nacks: true
checkpoint_limit: 1024
tls:
skip_cert_verify: false
enable_renegotiation: false
root_cas: ""
root_cas_file: ""
client_certs: []
aws:
enabled: false
region: "" # No default (optional)
endpoint: "" # No default (required)
id: "" # No default (optional)
secret: "" # No default (optional)
token: "" # No default (optional)
role: "" # No default (optional)
role_external_id: "" # No default (optional)
roles: [] # No default (optional)
batching:
count: 0
byte_size: 0
period: ""
check: ""
processors: [] # No default (optional)
The mysql_cdc input uses MySQL’s binary log (binlog) to capture changes made to a MySQL database in real time and streams them to Redpanda Connect.
Redpanda Connect allows you to specify which database tables in your source database to receive changes from. There are also two replication modes to choose from.
Prerequisites
Choose a replication mode
You can run the mysql_cdc input in one of two modes, depending on whether you need a snapshot of existing data.
-
Snapshot mode: Redpanda Connect first captures a snapshot of all data in the selected tables and streams the contents before processing changes from the last recorded binlog position.
-
Streaming mode: Redpanda Connect skips the snapshot and processes only the most recent data changes, starting from the latest binlog position.
Snapshot mode
If you set the stream_snapshot field to true, Redpanda Connect connects to your MySQL database and does the following to capture a snapshot of all data in the selected tables:
-
Executes the
FLUSH TABLES WITH READ LOCKquery to write any outstanding table updates to disk, and locks the tables. -
Runs the
START TRANSACTION WITH CONSISTENT SNAPSHOTstatement to create a new transaction with a consistent view of all data, capturing the state of the database at the moment the transaction started. -
Reads the current binlog position.
-
Runs the
UNLOCK TABLESstatement to release the database. -
Preserves the initial transaction for data integrity.
If the pipeline restarts during this process, Redpanda Connect must start the snapshot capture from scratch to store the current binlog position in the checkpoint_cache.
|
After the snapshot is taken, the input executes SELECT statements to extract data from the selected tables in two stages:
-
The input finds the primary keys of a table.
-
It selects the data ordered by primary key.
Finally, the input uses the stored binlog position to catch up with changes that occurred during snapshot processing.
Streaming mode
If you set the stream_snapshot field to false, Redpanda Connect connects to your MySQL database and starts processing data changes from the latest binlog position. If the pipeline restarts, Redpanda Connect resumes processing updates from the last binlog position written to the checkpoint_cache.
Binlog rotation
While the mysql_cdc input is streaming changes to Redpanda Connect, your MySQL server may rotate the binlog file. When this occurs, Redpanda Connect flushes the existing message batch and stores the new binlog position so that it can resume processing using the latest offset.
Data mappings
The following table shows how selected MySQL data types are mapped to data types supported in Redpanda Connect. All other data types are mapped to string values.
| MySQL data type | Bloblang value |
|---|---|
TEXT, VARCHAR |
A string value, for example: |
BINARY, VARBINARY, TINYBLOB, BLOB, MEDIUMBLOB, LONGBLOB |
An array of byte values, for example: |
DECIMAL, NUMERIC, TINYINT, SMALLINT, MEDIUMINT, INT, BIGINT, YEAR |
A standard numeric type, for example: |
FLOAT, DOUBLE |
A 64-bit decimal ( |
DATETIME, TIMESTAMP |
A Bloblang timestamp, for example:
|
SET |
An array of strings, for example: |
JSON |
A map object of the JSON, for example: |
Metadata
This input adds the following metadata fields to each message:
-
operation: The type of operation (insert, update, delete, or read for snapshot messages) -
table: The name of the table -
binlog_position: The binlog position (for CDC messages only, not set for snapshot messages) -
schema: The table schema in benthos common schema format, compatible with processors like parquet_encode
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
aws
AWS IAM authentication configuration for MySQL instances. When enabled, IAM credentials are used to generate temporary authentication tokens instead of a static password.
Requires version 4.72.0 or later.
Type: object
aws.enabled
Enable AWS IAM authentication for MySQL. When enabled, an IAM authentication token is generated and used as the password. When using IAM authentication ensure max_reconnect_attempts is set to a low value to ensure it can refresh credentials.
Type: bool
Default: false
aws.endpoint
The MySQL endpoint hostname (for example, mydb.abc123.us-east-1.rds.amazonaws.com).
Type: string
aws.id
The AWS access key ID to authenticate with. When empty, the default AWS credential chain is used.
Type: string
aws.region
The AWS region where the MySQL instance is located. The region is used to sign the IAM authentication token and for STS calls when assuming roles. If no region is specified, the environment default is used.
Type: string
aws.role
Optional AWS IAM role ARN to assume for authentication. When roles is also set, this role is assumed first and the roles entries are assumed after it.
Type: string
aws.role_external_id
Optional external ID to use when assuming the role set in role. Each entry in roles sets its own external ID.
Type: string
aws.roles[]
Optional array of AWS IAM roles to assume for authentication. Roles are assumed in sequence, each using the credentials of the previous one, enabling chaining for purposes such as cross-account access. Each role can optionally specify an external ID. When role is also set, it is assumed before the first entry.
Type: array<object>
aws.secret
The AWS secret access key that pairs with id.
|
This field contains sensitive information that usually shouldn’t be added to a configuration directly. For more information, see Secrets. |
Type: string
aws.token
The AWS session token to use with id and secret. Required only when using short-term credentials.
Type: string
batching
Configure a batching policy.
Type: object
# Examples:
batching:
byte_size: 5000
count: 0
period: 1s
# ---
batching:
count: 10
period: 1s
# ---
batching:
check: this.contains("END BATCH")
count: 0
period: 1m
batching.byte_size
The maximum total size (in bytes) that a batch can reach before it is flushed. When the combined size of all messages in the batch reaches or exceeds this limit, the batch is immediately sent to the next stage (such as a processor or output).
Set to 0 to disable size-based batching. When disabled, messages are flushed based on other conditions (such as count or period).
Type: int
Default: 0
batching.check
A Bloblang query that returns a boolean value indicating whether a message should end a batch.
Type: string
Default: ""
# Examples:
check: this.type == "end_of_transaction"
batching.count
The number of messages at which the batch is flushed. Set to 0 to disable count-based batching.
Type: int
Default: 0
batching.period
The length of time after which an incomplete batch is flushed regardless of its size. This field accepts Go duration format strings such as 100ms, 1s, or 5s. Supported time units are ns, us, ms, s, m, and h.
Type: string
Default: ""
# Examples:
period: 1s
# ---
period: 1m
# ---
period: 500ms
batching.processors[]
A list of processors to apply to a batch as it is flushed. This allows you to aggregate and archive the batch however you see fit. All resulting messages are flushed as a single batch, so splitting the batch into smaller batches with these processors has no effect.
Type: array<processor>
# Examples:
processors:
- archive:
format: concatenate
# ---
processors:
- archive:
format: lines
# ---
processors:
- archive:
format: json_array
checkpoint_cache
A cache resource to store the binlog position of the most recent data update delivered by Redpanda Connect. After a restart, Redpanda Connect can continue processing changes from this last known position, avoiding the need to reprocess all table updates.
Type: string
checkpoint_key
The key identifier used to store the binlog position in checkpoint_cache. If you have multiple mysql_cdc inputs sharing the same cache, you can provide an alternative key.
Type: string
Default: mysql_binlog_position
checkpoint_limit
The maximum number of messages that this input can process at a given time. Increasing this limit enables parallel processing, and batching at the output level. To preserve at-least-once guarantees, any given binlog position is not acknowledged until all messages up to that position are delivered.
Type: int
Default: 1024
dsn
The data source name (DSN) of the MySQL database from which you want to stream updates. Use the format user:password@tcp(localhost:3306)/database.
Type: string
# Examples:
dsn: user:password@tcp(localhost:3306)/database
flavor
The type of MySQL database to connect to.
Requires version 4.48.0 or later.
Type: string
Default: mysql
| Option | Summary |
|---|---|
|
MariaDB flavored databases. |
|
MySQL flavored databases. |
max_parallel_snapshot_tables
Specifies the number of tables that will be snapshotted in parallel.
Requires version 4.90.0 or later.
Type: int
Default: 1
max_reconnect_attempts
The maximum number of attempts the MySQL driver will try to re-establish a broken connection before Connect attempts reconnection. A zero or negative number means infinite retry attempts.
Requires version 4.72.0 or later.
Type: int
Default: 10
snapshot_max_batch_size
The maximum number of table rows to fetch in each batch when taking a snapshot. This option is only available when stream_snapshot is set to true.
Type: int
Default: 1000
stream_snapshot
When set to true, this input streams a snapshot of all existing data in the source database before streaming data changes. To use this setting, all database tables that you want to replicate must have a primary key. When set to false, the input starts streaming from the current binlog position.
Type: bool
tables[]
A list of the database table names to stream changes from. Specify each table name as a separate item.
Type: array<string>
# Examples:
tables:
- table1
- table2
tls
Custom TLS settings for the MySQL connection. When enabled is true, these settings replace any tls parameter in the dsn, and the server name is set to the host from the DSN.
Requires version 4.72.0 or later.
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.
Certificate pairing rules: For each certificate item, provide either:
-
Inline PEM data using both
certandkeyor -
File paths using both
cert_fileandkey_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 Secrets. |
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 Secrets. |
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.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 Secrets. |
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