microsoft_sql_server_cdc
Enables Change Data Capture by consuming from Microsoft SQL Server’s change tables.
-
Common
-
Advanced
input:
label: ""
microsoft_sql_server_cdc:
connection_string: "" # No default (required)
stream_snapshot: false
max_parallel_snapshot_tables: 1
snapshot_max_batch_size: 1000
include: [] # No default (required)
exclude: [] # No default (optional)
checkpoint_cache: "" # No default (optional)
checkpoint_cache_table_name: rpcn.CdcCheckpointCache
checkpoint_cache_connection_string: "" # No default (optional)
checkpoint_cache_key: microsoft_sql_server_cdc
checkpoint_limit: 1024
stream_backoff_interval: 5s
auto_replay_nacks: true
batching:
count: 0
byte_size: 0
period: ""
check: ""
input:
label: ""
microsoft_sql_server_cdc:
connection_string: "" # No default (required)
stream_snapshot: false
max_parallel_snapshot_tables: 1
snapshot_max_batch_size: 1000
include: [] # No default (required)
exclude: [] # No default (optional)
checkpoint_cache: "" # No default (optional)
checkpoint_cache_table_name: rpcn.CdcCheckpointCache
checkpoint_cache_connection_string: "" # No default (optional)
checkpoint_cache_key: microsoft_sql_server_cdc
checkpoint_limit: 1024
stream_backoff_interval: 5s
auto_replay_nacks: true
batching:
count: 0
byte_size: 0
period: ""
check: ""
processors: [] # No default (optional)
Streams changes from a Microsoft SQL Server database for Change Data Capture (CDC).
Additionally, if stream_snapshot is set to true, then the existing data in the database is also streamed too.
Metadata
This input adds the following metadata fields to each message:
-
database_schema(The database schema for the table where the message originates from) -
schema(The table schema in benthos common schema format, compatible with processors like parquet_encode) -
table(Name of the table that the message originated from) -
operation(Type of operation that generated the message: "read", "delete", "insert", or "update_before" and "update_after". "read" is from messages that are read in the initial snapshot phase.) -
lsn(The commit Log Sequence Number of the change, from the change table column__$start_lsn. Not present on snapshot (read) messages.) -
seqval(The position of the change in the transaction log, from the change table column__$seqval, as a hexadecimal string with a0xprefix. Not present on snapshot (read) messages.) -
command_id(The order of the statement within its transaction, from the change table column__$command_id. Not present on snapshot (read) messages.)
Permissions
To use the default Microsoft SQL Server cache, the user must have permissions to create tables and stored procedures. Refer to checkpoint_cache_table_name for additional details.
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
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 current Log Sequence Number (LSN) position. The cache stores the highest LSN that has been successfully delivered downstream, which allows Redpanda Connect to resume from the last processed position after a restart, rather than consume the entire state of the change table. If not set, the default Microsoft SQL Server based cache is used. See checkpoint_cache_table_name for more information.
Type: string
checkpoint_cache_connection_string
An optional connection string for a remote Microsoft SQL Server to use for the checkpoint cache. When set, this creates the checkpoint cache table on the remote server instead of the source database. If checkpoint_cache is also set, that takes precedence.
Type: string
# Examples:
checkpoint_cache_connection_string: sqlserver://username:password@remotehost/instance?param1=value¶m2=value
checkpoint_cache_key
The key to use to store the snapshot position in checkpoint_cache. An alternative key can be provided if multiple CDC inputs share the same cache.
Type: string
Default: microsoft_sql_server_cdc
checkpoint_cache_table_name
The multipart identifier for the checkpoint cache table name. If no checkpoint_cache field is specified, this input will automatically create a table and stored procedure under the rpcn schema to act as a checkpoint cache. This table stores the latest processed Log Sequence Number (LSN) that has been successfully delivered, allowing Redpanda Connect to resume from that point upon restart rather than reconsume the entire change table.
Type: string
Default: rpcn.CdcCheckpointCache
# Examples:
checkpoint_cache_table_name: dbo.checkpoint_cache
checkpoint_limit
The maximum number of messages that can be processed concurrently before applying back pressure. Higher values enable more parallel processing and batching at the output level, but increase memory usage. To preserve at-least-once delivery guarantees, a given Log Sequence Number (LSN) is only acknowledged after all messages under that offset are delivered.
Type: int
Default: 1024
connection_string
The connection string for the Microsoft SQL Server database. Use the format sqlserver://username:password@host/instance?param1=value¶m2=value. For Windows Authentication, use sqlserver://host/instance?trusted_connection=yes. Include additional parameters like TrustServerCertificate=true for self-signed certificates or encrypt=disable to disable encryption.
Type: string
# Examples:
connection_string: sqlserver://username:password@host/instance?param1=value¶m2=value
exclude[]
Regular expressions for tables to exclude from CDC streaming. Use this to filter out specific tables from the include patterns. Table names should follow the schema.table format. Exclude patterns are applied after include patterns, allowing you to include broad patterns while excluding specific tables.
Type: array<string>
# Examples:
exclude: dbo.privatetable
include[]
Regular expressions for tables to include in CDC streaming. Specify table names using the format schema.table (such as dbo.orders, sales.customers). Each pattern is treated as a regular expression, allowing wildcards and pattern matching. All specified tables must have CDC enabled in SQL Server.
Type: array<string>
# Examples:
include: dbo.products
max_parallel_snapshot_tables
The maximum number of tables to read in parallel during the initial snapshot. Each table is read by its own reader.
Type: int
Default: 1
snapshot_max_batch_size
The maximum number of rows to stream in a single batch during the initial snapshot phase. Larger batch sizes can improve throughput for initial data loads but may increase memory usage. This setting only applies when stream_snapshot is enabled.
Type: int
Default: 1000
stream_backoff_interval
The interval to wait before checking for new changes after a pass over the change tables completes. Each pass drains changes up to the maximum LSN observed as the pass began, then waits for this interval while SQL Server’s capture job continues to publish changes. For low-traffic tables, increasing this value reduces query load on the server. For high-traffic tables, a longer interval directly reduces throughput, because the input sits idle for the full interval between passes. Consider lowering it towards 500ms, which matches the default poll interval of comparable CDC systems. Use Go duration format such as 500ms, 5s, or 1m.
Type: string
Default: 5s
# Examples:
stream_backoff_interval: 500ms
# ---
stream_backoff_interval: 5s
# ---
stream_backoff_interval: 1m
stream_snapshot
Whether to stream a snapshot of all existing data before streaming CDC changes. When set to true, the connector first queries all existing table data, then switches to streaming incremental changes from the change tables. When set to false, no snapshot is taken and, on first run, streaming begins from the start of each table’s existing change table: every change retained by SQL Server’s CDC capture and cleanup jobs (three days by default) is replayed, not only changes from the current LSN onward. To begin from the present on a table that already holds change history, disable and re-enable CDC on the table immediately before starting the pipeline so that its change table starts empty.
Type: bool
Default: false
# Examples:
stream_snapshot: true