Connect

microsoft_sql_server_cdc

Enables Change Data Capture by consuming from Microsoft SQL Server’s change tables.

Introduced in version 4.67.5.

  • 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 a 0x prefix. 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.

Requires version 4.87.0 or later.

Type: string

# Examples:
checkpoint_cache_connection_string: sqlserver://username:password@remotehost/instance?param1=value&param2=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.

Requires version 4.67.0 or later.

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.

Requires version 4.67.0 or later.

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&param2=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&param2=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.

Requires version 4.69.0 or later.

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