Skip to main content

概述

Scylla CDC Source Connector 从启用了 CDC 的 Scylla 表读取行级变更,并将 INSERT、UPDATE 和受支持的 DELETE 事件写入 Kafka。每个源表默认对应一个形如 <topic-prefix>.<keyspace>.<table> 的 Topic,记录值使用 Debezium 变更事件结构,记录 Key 可保留完整的分区键和聚簇键。 Connector 只持续读取 Scylla CDC 日志,不会对基表执行初始快照。它适合将订单、账户、设备状态等表的增量变化接入事件流、缓存更新、搜索索引或审计链路。新部署应显式使用 advanced 输出格式,以获得直接字段值、复杂类型和可配置的变更前后镜像;默认的 legacy 格式已弃用。

前置条件

  • 开始采集前,目标 Scylla 表必须已创建并启用 CDC。本文快速开始和其他基础示例使用 cdc.include.after=only-updated,因此目标表还必须启用 CDC postimage;如需 before 镜像,还应启用 CDC preimage。CDC 日志保留时间应覆盖预期停机与恢复窗口。
  • Connector 使用的 Scylla 账号需要读取目标 Keyspace、表、CDC 日志和相关元数据的权限;启用 TLS 时,证书、信任库、密钥库或私钥文件必须可由 Connect Worker 读取。

授权许可

使用 Apache License 2.0。

快速开始

提前准备 Connect Cluster、Kafka 和一个已启用 CDC 的 Scylla 表,确认 Connect Worker 可以访问 Scylla 和 Kafka,并具有所需权限。具体准备和管理操作请参阅 管理 Connector
<scylla-host><keyspace><table> 替换为实际资源,并在启动 Connector 前为目标表启用 CDC postimage。若 Scylla 启用了认证,请同时添加 scylla.userscylla.password;生产环境通常还应启用 TLS。该配置让 INSERT 的 after 包含完整新行,让 UPDATE 的 after 包含主键和实际变更列;DELETE 通过 op=d 与 Kafka Key 识别。示例只读取启动后及 CDC 日志中可恢复的变更,不会发送基表中的现有行。

配置

Connector 标识与源表

connector.class

指定 Connector 实现类。
  • 类型string
  • 默认值:无
  • 重要级别:高
  • 有效值 / 注意事项:必须设置为 com.scylladb.cdc.debezium.connector.ScyllaConnector
  • 必填:是

scylla.table.names

指定要读取的 Scylla 表。
  • 类型list
  • 默认值null
  • 重要级别:高
  • 有效值 / 注意事项:必填,使用逗号分隔的 keyspace.table。未加引号的 Keyspace 和表名只能包含 1 到 48 个 ASCII 字母、数字或下划线;不支持带引号标识符。每个表都必须启用 CDC,缺失的表不会使校验失败,Connector 会等待其出现。
  • 必填:是

topic.prefix

设置数据 Topic 和心跳 Topic 的命名空间前缀。
  • 类型string
  • 默认值null
  • 重要级别:高
  • 有效值 / 注意事项:必填,并应在不同 Connector 实例间保持唯一。默认命名策略使用该前缀生成 <topic-prefix>.<keyspace>.<table>
  • 必填:是

topic.naming.strategy

指定 Debezium Topic 命名策略实现。
  • 类型class
  • 默认值io.debezium.schema.SchemaTopicNamingStrategy
  • 重要级别:中
  • 有效值 / 注意事项:自定义类必须实现 Debezium TopicNamingStrategy,并与 topic.prefixheartbeat.topics.prefix 一起决定 Topic 名称。

Scylla 连接与认证

scylla.cluster.ip.addresses

指定 Scylla CQL 联系点。
  • 类型list
  • 默认值null
  • 重要级别:高
  • 有效值 / 注意事项:必填,使用逗号分隔的 host:port;每项都必须显式提供数字端口并可从 Worker 访问。当前解析方式不支持包含冒号的 IPv6 字面量。
  • 必填:是

scylla.local.dc

指定优先访问的本地数据中心。
  • 类型string
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:名称必须与 Scylla 集群拓扑元数据一致;省略时由驱动选择。

scylla.user

指定 CQL 用户名。
  • 类型string
  • 默认值null
  • 重要级别:高
  • 有效值 / 注意事项:省略时使用无认证连接;设置后必须同时设置 scylla.password

scylla.password

指定 CQL 密码。
  • 类型password
  • 默认值null
  • 重要级别:高
  • 有效值 / 注意事项:设置后必须同时设置 scylla.user。请通过受控的 Secret 或 Config Provider 提供,不要写入普通配置文件或日志。

scylla.consistency.level

设置读取 CDC 日志时使用的一致性级别。
  • 类型string
  • 默认值quorum
  • 重要级别:中
  • 有效值 / 注意事项:使用底层 Scylla CDC 驱动支持的一致性级别;该值只影响 CDC 日志查询。

scylla.query.options.fetch.size

设置 CQL 查询分页大小。
  • 类型int
  • 默认值0
  • 重要级别:低
  • 有效值 / 注意事项:必须为非负整数。0 使用驱动默认值,正数设置每页读取量。

TLS

scylla.ssl.enabled

控制是否启用到 Scylla 的 TLS 连接。
  • 类型boolean
  • 默认值null
  • 重要级别:高
  • 有效值 / 注意事项:应显式设置为 truefalse;设置为 true 后才读取其他 TLS 配置。

scylla.ssl.provider

选择 TLS 实现提供方。
  • 类型string
  • 默认值jdk
  • 重要级别:低
  • 有效值 / 注意事项:可选 jdkopensslopenssl_refcnt,仅在 scylla.ssl.enabled=true 时生效。

scylla.ssl.truststore.path

指定 Java Truststore 文件路径。
  • 类型string
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:仅用于 TLS,文件必须存在且可由 Worker 读取;Connector 不预先校验路径。

scylla.ssl.truststore.password

指定 Java Truststore 密码。
  • 类型string
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:与 scylla.ssl.truststore.path 配合使用。虽然 ConfigDef 类型为 string,仍应按 Secret 管理。

scylla.ssl.keystore.path

指定 Java Keystore 文件路径。
  • 类型string
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:用于需要客户端证书的 TLS 连接,文件必须可由 Worker 读取;Connector 不预先校验路径。

scylla.ssl.keystore.password

指定 Java Keystore 密码。
  • 类型string
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:与 scylla.ssl.keystore.path 配合使用。虽然 ConfigDef 类型为 string,仍应按 Secret 管理。

scylla.ssl.cipherSuites

限制 TLS 密码套件。
  • 类型list
  • 默认值null
  • 重要级别:高
  • 有效值 / 注意事项:仅在 TLS 下生效;省略时使用所选 TLS Provider 的默认套件,显式值必须受 Provider 支持。

scylla.ssl.openssl.keyCertChain

指定 OpenSSL 客户端证书链路径。
  • 类型string
  • 默认值null
  • 重要级别:高
  • 有效值 / 注意事项:仅与 OpenSSL Provider 配合使用,文件必须可由 Worker 读取,并应与 scylla.ssl.openssl.privateKey 匹配。

scylla.ssl.openssl.privateKey

指定 OpenSSL 客户端私钥路径。
  • 类型string
  • 默认值null
  • 重要级别:高
  • 有效值 / 注意事项:仅与 OpenSSL Provider 配合使用。应限制文件权限并确保 Worker 可读;Connector 不校验路径或证书配对关系。

启动、读取窗口与缓冲

scylla.initial.lookback.ms

设置无已存 Offset 时首次读取的回溯时长。
  • 类型long
  • 默认值0
  • 重要级别:中
  • 有效值 / 注意事项:必须为非负数。0 从首个 CDC Generation 的起始窗口读取;正数从当前时间向前回溯相应时长。仅对没有已存 Offset 的具体任务生效,且不能恢复早于 CDC 保留期的数据。

streaming.delay.ms

设置进入流式读取前的等待时间。
  • 类型long
  • 默认值0
  • 重要级别:低
  • 有效值 / 注意事项:必须为非负数。Scylla 的快照阶段不读取数据,通常保持 0

scylla.query.time.window.size

设置每次 CDC 查询覆盖的时间窗口,单位毫秒。
  • 类型int
  • 默认值30000
  • 重要级别:低
  • 有效值 / 注意事项:必须为非负整数。较小窗口降低单次查询工作量但增加查询次数,也用于构造首次回溯窗口。

scylla.confidence.window.size

设置暂不读取最新 CDC 区间的安全窗口,单位毫秒。
  • 类型int
  • 默认值30000
  • 重要级别:低
  • 有效值 / 注意事项:必须为非负整数。较大值有助于避开尚未稳定的最新 CDC 数据,但会增加端到端延迟;0 取消该延迟。

scylla.minimal.wait.for.window.time

设置连续 CDC 查询窗口之间的最短等待时间,单位毫秒。
  • 类型int
  • 默认值0
  • 重要级别:低
  • 有效值 / 注意事项:必须为非负整数。0 不主动限速;正数可降低追赶期间的查询压力。

poll.interval.ms

设置从内部变更队列轮询记录的间隔,单位毫秒。
  • 类型long
  • 默认值500
  • 重要级别:中
  • 有效值 / 注意事项:必须为非负数。较小值可降低空闲后的响应延迟,但会增加轮询频率。

max.batch.size

设置每次 Poll 返回的最大记录数。
  • 类型int
  • 默认值2048
  • 重要级别:中
  • 有效值 / 注意事项:必须为正整数,并应小于 max.queue.size

max.queue.size

设置内部变更队列可容纳的最大记录数。
  • 类型int
  • 默认值8192
  • 重要级别:中
  • 有效值 / 注意事项:必须为正整数,并应大于 max.batch.size 以保留缓冲余量;队列满时读取端会受到背压。

max.queue.size.in.bytes

设置按字节计算的队列上限。
  • 类型long
  • 默认值0
  • 重要级别:中
  • 有效值 / 注意事项:该公开配置在 2.0.6 中未接入 Scylla 的队列构建逻辑,不能用于限制内存;请使用 max.queue.size 控制记录数量。

CDC 输出

cdc.output.format

选择 CDC 记录的输出格式。
  • 类型string
  • 默认值legacy
  • 重要级别:高
  • 有效值 / 注意事项:可选 legacyadvancedlegacy 已弃用并计划在 3.0.0 移除;新部署应显式使用 advanced,以支持直接字段值、复杂类型及前后镜像。
  • 弃用:默认值 legacy 已弃用,替代值为 advanced

experimental.preimages.enabled

控制 Legacy 格式是否读取 Preimage。
  • 类型boolean
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项:仅适用于 legacy;当 cdc.output.format=advanced 时设置为 true 会被拒绝。Advanced 格式请使用 cdc.include.before

cdc.include.before

控制 Advanced 记录中的变更前镜像。
  • 类型string
  • 默认值none
  • 重要级别:中
  • 有效值 / 注意事项:可选 nonefullonly-updated,只适用于 advancedfullonly-updated 要求所有现有目标表启用 CDC preimage。

cdc.include.after

控制 Advanced 记录中的变更后镜像。
  • 类型string
  • 默认值none
  • 重要级别:中
  • 有效值 / 注意事项:可选 nonefullonly-updated,只适用于 advancedfullonly-updated 要求所有现有目标表启用 CDC postimage。

cdc.include.primary-key.placement

设置主键在输出记录中的位置。
  • 类型list
  • 默认值kafka-key,payload-after,payload-before
  • 重要级别:中
  • 有效值 / 注意事项:Legacy 格式只接受默认组合;Advanced 格式可使用 kafka-keypayload-afterpayload-beforepayload-keykafka-headerspayload-diff 虽可通过校验,但 2.0.6 未实现实际输出。ConfigDef 声明的默认值仍为 kafka-key,payload-after,payload-before,但 2.0.6 使用 advanced 格式时必须显式设置该组合,才能生成非 null 的结构化 Kafka Key;省略该配置可能得到 required Key Schema 与 null Key 值。依赖按 Key 分区、顺序或日志压缩时应保留 kafka-key

cdc.include.primary-key.payload-key-name

设置 payload-key 主键对象的字段名。
  • 类型string
  • 默认值key
  • 重要级别:低
  • 有效值 / 注意事项:仅在主键位置包含 payload-key 时使用;应选择不会与其他顶层字段冲突的非空名称。

cdc.incomplete.task.timeout.ms

设置 Advanced 镜像组合等待超时,单位毫秒。
  • 类型long
  • 默认值15000
  • 重要级别:低
  • 有效值 / 注意事项:必须为正数。超过该时长仍缺少所需 Preimage 或 Postimage 的变更会被记录为错误并丢弃;清理依赖后续事件推进,不是精确定时器。

tombstones.on.delete

控制 DELETE 事件之后是否发送 Tombstone。
  • 类型boolean
  • 默认值true
  • 重要级别:中
  • 有效值 / 注意事项true 会在删除事件后发送同 Key、空 Value 的记录,便于日志压缩 Topic 清理旧值;移除 kafka-key 时无法获得常规按 Key 压缩语义。

skipped.operations

指定不输出的 Debezium 操作类型。
  • 类型list
  • 默认值t
  • 重要级别:低
  • 有效值 / 注意事项:可使用 cudtnone。默认跳过 Truncate;该配置不能使 Connector 支持本身无法转为行事件的多行分区删除或范围删除。

event.processing.failure.handling.mode

设置 Debezium 事件处理失败时的行为。
  • 类型string
  • 默认值fail
  • 重要级别:中
  • 有效值 / 注意事项:可选 failwarnignorefail 停止任务;warnignore 可能跳过有问题的事件,应仅在业务可接受数据缺口时使用。

transaction.metadata.factory

指定 Debezium 事务元数据工厂。
  • 类型class
  • 默认值io.debezium.pipeline.txmetadata.DefaultTransactionMetadataFactory
  • 重要级别:低
  • 有效值 / 注意事项:属于高级扩展点。2.0.6 未公开启用事务元数据的配置,因此更换该工厂不会使 Connector 产生事务元数据。

重试与错误恢复

retriable.restart.connector.wait.ms

设置 Debezium 在可重试异常后等待重启的时间,单位毫秒。
  • 类型long
  • 默认值10000
  • 重要级别:低
  • 有效值 / 注意事项:必须为正数;它与 Scylla Worker 的重试参数属于不同层级。

errors.max.retries

设置 Debezium 连接错误的最大重试次数。
  • 类型int
  • 默认值-1
  • 重要级别:低
  • 有效值 / 注意事项-1 表示不限次数,0 表示不重试,正数表示有限重试;不要与 worker.max.retries 混淆。

worker.retry.backoff.base

设置 Scylla Worker 指数退避的基础时长,单位毫秒。
  • 类型int
  • 默认值50
  • 重要级别:低
  • 有效值 / 注意事项:必须为非负整数,并应不大于 worker.maximum.backoff;Connector 不校验两者关系。

worker.maximum.backoff

设置 Scylla Worker 指数退避上限,单位毫秒。
  • 类型int
  • 默认值30000
  • 重要级别:低
  • 有效值 / 注意事项:必须为非负整数,并应不小于 worker.retry.backoff.base

worker.jitter.percentage

设置 Worker 重试退避的随机抖动百分比。
  • 类型int
  • 默认值20
  • 重要级别:低
  • 有效值 / 注意事项:使用 1100。2.0.6 只校验正整数,未强制上限,但不应配置超过 100

worker.max.retries

设置 Scylla Worker 的总执行尝试次数。
  • 类型int
  • 默认值20
  • 重要级别:低
  • 有效值 / 注意事项:接受正整数或 -11 表示只执行一次且不重试,-1 表示不限次数;0 和小于 -1 的值无效。

Task 与连接池

tasks.max

设置 Kafka Connect 可创建的最大 Task 数。
  • 类型int
  • 默认值1
  • 重要级别:高
  • 有效值 / 注意事项:必须至少为 1。实际 Task 数可能更少,取决于当前 CDC vnode 或 tablet 工作分配;增加该值不保证吞吐线性增长,并会影响连接与队列总量。

tasks.max.enforce

控制框架是否强制执行 tasks.max 上限。
  • 类型boolean
  • 默认值true
  • 重要级别:低
  • 有效值 / 注意事项:Kafka Connect 3.9.1 已将该配置标记为弃用并计划移除。保持 true,通过 tasks.max 规划并发。
  • 弃用:是;没有直接替代配置。

worker.shared.session.enabled

控制同一 Worker JVM 内的 Task 是否共享 Scylla Session。
  • 类型boolean
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项true 时,只有连接和池配置相同的 Task 才共享 Session。共享仅限单个 JVM,可减少连接数,但会增加同一池上的竞争。

worker.pooling.core.pool.local

设置每个本地 Scylla 节点的目标核心连接数。
  • 类型int
  • 默认值1
  • 重要级别:低
  • 有效值 / 注意事项:必须为非负整数,并应不大于 worker.pooling.max.pool.local;Connector 不校验两者关系。

worker.pooling.max.pool.local

设置每个本地 Scylla 节点的最大连接数。
  • 类型int
  • 默认值1
  • 重要级别:低
  • 有效值 / 注意事项:必须为非负整数,并应不小于 worker.pooling.core.pool.local

worker.pooling.max.requests.per.connection

设置单个连接允许的最大并发请求数。
  • 类型int
  • 默认值256
  • 重要级别:低
  • 有效值 / 注意事项:必须为非负整数。未启用 Session 共享时,每个 Task 使用独立连接池,整体并发上限会随实际 Task 数增长;启用共享后,同一 Worker JVM 中配置一致的 Task 共用 Session,该值作用于共享连接池中的每个连接。

worker.pooling.max.queue.size

设置每个连接池的请求排队上限。
  • 类型int
  • 默认值256
  • 重要级别:低
  • 有效值 / 注意事项:必须为非负整数。超过容量或等待超时的请求会被拒绝。未启用 Session 共享时,每个 Task 使用独立连接池;启用共享后,同一 Worker JVM 中配置一致的 Task 共用并竞争同一连接池。

worker.pooling.pool.timeout.ms

设置从主机连接池获取连接的最长等待时间,单位毫秒。
  • 类型int
  • 默认值5000
  • 重要级别:低
  • 有效值 / 注意事项:必须为非负整数,应结合请求队列、连接数和 Scylla 延迟调整。

心跳

heartbeat.interval.ms

设置心跳记录间隔,单位毫秒。
  • 类型int
  • 默认值30000
  • 重要级别:中
  • 有效值 / 注意事项:必须为正整数,不能设置为 0。心跳用于在低流量期间推进并持久化 CDC 读取位置。

heartbeat.topics.prefix

设置心跳 Topic 的名称前缀。
  • 类型string
  • 默认值__debezium-heartbeat
  • 重要级别:低
  • 有效值 / 注意事项:命名策略会将该值与 topic.prefix 组合生成心跳 Topic;应符合 Kafka Topic 命名规则。

Converter 与扩展处理

key.converter

设置 SourceRecord Key 的 Kafka Connect Converter。
  • 类型class
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:省略时继承 Worker 配置。Converter 必须可实例化,并支持结构化主键;若从主键位置中移除 kafka-key,记录 Key 可能为 null

value.converter

设置 SourceRecord Value 的 Kafka Connect Converter。
  • 类型class
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:省略时继承 Worker 配置。Converter 必须支持 Debezium Envelope 的 STRUCT Schema 以及所选输出格式的数据类型。

converters

设置 Debezium 自定义 Converter 别名。
  • 类型string
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:使用逗号分隔别名,并通过别名前缀配置实现类和选项;实现类必须存在于插件 Classpath。

post.processors

设置 Debezium Post Processor 别名。
  • 类型string
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:使用逗号分隔别名,并只启用与当前 Debezium 运行时兼容且已安装的实现。

信号与通知

signal.enabled.channels

设置启用的 Debezium 信号通道。
  • 类型list
  • 默认值source
  • 重要级别:中
  • 有效值 / 注意事项:只配置当前插件 Classpath 和 Debezium 运行时支持的通道。Source 通道默认启用。

signal.poll.interval.ms

设置已启用信号通道的轮询间隔,单位毫秒。
  • 类型long
  • 默认值5000
  • 重要级别:中
  • 有效值 / 注意事项:必须为正数。Scylla 数据库信号表未接入,因此该值不能启用数据库表信号。

signal.data.collection

指定数据库信号集合。
  • 类型string
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:2.0.6 的 Scylla Task 未接入数据库信号表映射,该配置不能用于启用 Scylla 表信号;应使用受支持的其他信号通道。

notification.enabled.channels

设置启用的 Debezium 通知通道。
  • 类型list
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:只配置当前插件 Classpath 中可用的通道;启用 Sink 通知通道时还需设置 notification.sink.topic.name

notification.sink.topic.name

指定 Sink 通知通道使用的 Kafka Topic。
  • 类型string
  • 默认值null
  • 重要级别:高
  • 有效值 / 注意事项:仅在通知通道包含 Sink 时必需;Topic 必须有效且 Connector 有写权限。

可观测性

custom.metric.tags

为 Debezium JMX Metric 添加自定义标签。
  • 类型list
  • 默认值null
  • 重要级别:低
  • 有效值 / 注意事项:使用 key=value 列表。标签值应稳定且不包含密码、Token 或其他敏感信息,并符合 JMX ObjectName 约束。

已注册但不支持的快照配置

snapshot.mode.custom.name

指定自定义 Snapshotter 名称。
  • 类型string
  • 默认值null
  • 重要级别:中
  • 有效值 / 注意事项:该字段虽公开注册,但 2.0.6 固定使用不读取 Schema 或表数据的快照流程,也没有公开 snapshot.mode;设置后没有实际作用。

snapshot.mode.configuration.based.snapshot.data

控制 Configuration-based 模式是否快照数据。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:2.0.6 不公开可用的 Configuration-based Snapshot 模式,该配置不能请求基表快照。

snapshot.mode.configuration.based.snapshot.schema

控制 Configuration-based 模式是否快照 Schema。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:2.0.6 的 Scylla 快照实现不采集 Schema,该配置没有实际作用。

snapshot.mode.configuration.based.start.stream

控制 Configuration-based 模式是否启动流式读取。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:2.0.6 未公开可用的 Configuration-based Snapshot 模式,该配置没有实际作用。

snapshot.mode.configuration.based.snapshot.on.schema.error

控制 Schema 错误后是否重新执行 Configuration-based 快照。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:Scylla 的快照实现不读取 Schema 或表数据,该配置没有实际作用。

snapshot.mode.configuration.based.snapshot.on.data.error

控制数据错误后是否重新执行 Configuration-based 快照。
  • 类型boolean
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项:Scylla 的快照实现不读取 Schema 或表数据,该配置没有实际作用。

incremental.snapshot.watermarking.strategy

设置增量快照水位策略。
  • 类型string
  • 默认值INSERT_INSERT
  • 重要级别:低
  • 有效值 / 注意事项:该字段由 Debezium 公开注册,但 2.0.6 不提供可用的 Scylla 增量快照能力,不能用它读取基表现有数据。

最佳实践

为更新事件保留变更前后镜像

适用业务场景:下游需要比较字段变化、生成审计记录或维护完整物化视图,且目标表可以在接入前启用 CDC preimage 和 postimage。预期结果是 UPDATE 事件包含完整的 beforeafter,而 INSERT 和 DELETE 仍遵循各自可用的镜像边界。 配置示例
关键说明:先在所有目标表上启用对应的 CDC preimage 和 postimage,再启动 Connector。完整镜像会增加 Scylla CDC 存储、读取流量和 Kafka 消息大小;如果只关心被更新的列,可改用 only-updated。镜像记录缺失并超过超时后,相应逻辑变更会被丢弃,因此应结合错误日志和下游数据完整性要求设置超时。

首次接入时限定 CDC 历史回溯范围

适用业务场景:目标表已经持续产生 CDC 日志,但首次上线只需要最近一段时间的变化,不希望从首个仍可见的 CDC Generation 开始追赶。预期结果是没有已存 Offset 的任务从当前时间向前一小时附近开始读取,而不是生成基表快照。 配置示例
关键说明:目标表需要启用 CDC postimage,以便回溯和实时事件携带可消费的非主键变更值。scylla.initial.lookback.ms 只在对应 Task 没有已存 Offset 时生效,不能覆盖已有恢复位置,也不能读取已经超过 Scylla CDC 保留期的数据。先根据业务允许的历史范围和预计追赶吞吐确定回溯时长;若需要基表现有行,应使用独立的数据初始化流程,并与后续 CDC 事件做好去重和衔接。

在控制 Scylla 连接压力的前提下增加并发

适用业务场景:监控显示 Connector 持续积压,Kafka 和 Worker 仍有处理余量,但单 Task 无法及时消费当前 vnode 或 tablet 分配。预期结果是在当前可分配工作单元范围内增加并发,同时让同一 Worker JVM 的 Task 共享 Session,避免连接数简单按 Task 倍增。 配置示例
关键说明:目标表需要启用 CDC postimage。实际 Task 数可能低于 tasks.max,并且拓扑或 CDC Generation 变化会触发重新分配,不能假设一个 Task 固定对应一个 vnode 或 tablet。共享 Session 只在同一 JVM 内生效,减少连接数的同时会增加池竞争;逐步提高并发,并同时观察 Scylla 请求延迟、连接池排队、Task 重试、Kafka 吞吐和端到端延迟。

监控

监控内容

监控 Connect Cluster、Worker、Connector 和 Task 的运行状态与重启次数,并关注 Source Record 吞吐、端到端延迟、积压趋势、Offset 提交延迟与失败、错误和重试,以及 Worker JVM 的 CPU、堆内存、GC 和线程状态;如部署显式启用了错误容忍与 DLQ,再同时监控 DLQ 写入量和失败情况。

导入 Grafana 大盘

下载 AutoMQ Connect Cluster Grafana Dashboard,确认 Grafana 已配置采集 Kafka Connect JMX 指标的数据源,并且 Connector、Task、Worker 等标签可用于筛选,然后在 Grafana 的 Dashboard 导入页面上传该 JSON 文件并选择对应数据源。

限制条件

  • Connector 不执行基表初始快照,也不支持增量快照;只能读取仍保留在 Scylla CDC 日志中的变更。
  • 交付语义为至少一次。Kafka 已接收记录但 Source Offset 尚未持久化时发生故障,重启后可能重复输出;Connector 不提供 Exactly-once 或自身去重。
  • 不保证跨表、vnode、tablet、CDC Generation、Task 或 Kafka 分区的全局顺序;移除 kafka-key 还会削弱同一主键通常依赖的分区顺序基础。
  • 带聚簇列的表发生多行分区删除时不会逐行输出删除事件,行范围删除也不受支持;只有不含聚簇列的表可将分区删除表示为单行删除。
  • Connector 不支持行过滤、列 Include/Exclude 或动态表发现,采集范围只能通过 scylla.table.names 显式列出。
  • legacy 输出格式不支持 List、Set、Map、Tuple、UDT 或 Postimage,并已弃用;其孤立 Preimage 没有超时清理机制。
  • Connector 不输出事务元数据,也不建立跨 Scylla 读取位置、Kafka 记录和 Connect Offset 的事务边界。

常见问题

为什么使用 advanced 格式时 Kafka Key 为空或 Key Converter 报 required 字段为 null

Scylla CDC Source Connector 2.0.6 在 advanced 格式下不能只依赖 cdc.include.primary-key.placement 的 ConfigDef 默认值。请显式设置 cdc.include.primary-key.placement=kafka-key,payload-after,payload-before,并保留支持结构化主键的 key.converter;重新启动 Connector 后,Kafka Key 将包含源表主键字段。还应确认目标表确实定义了主键,并检查生效的 Connector 配置中没有覆盖或移除 kafka-key

Connector 启动后为什么没有发送表中已有数据?

该 Connector 只读取 CDC 日志,不扫描基表。检查目标表是否已启用 CDC、是否在 scylla.table.names 中,以及 CDC 日志是否仍保留目标时间范围。首次接入需要近期历史时,在没有已存 Offset 的前提下设置 scylla.initial.lookback.ms;需要完整存量数据时,应先使用独立初始化流程,再从明确的 CDC 边界持续同步。

为什么配置校验通过后仍然收不到某个表的变更?

缺失的目标表不会使配置校验失败,Connector 会等待表出现。确认 keyspace.table 拼写、未使用带引号标识符,并检查表已经创建且启用 CDC。若配置了 cdc.include.beforecdc.include.after,还要确认表分别启用了 Preimage 或 Postimage,并检查 Connector 与 Task 状态、Scylla 权限和错误日志。

为什么重启后看到重复事件?

Connector 使用至少一次交付。若 Kafka 已确认记录但对应 Source Offset 尚未成功刷新,重启会从上一次持久化位置恢复并再次发送部分事件。保留完整主键作为 Kafka Key,在下游使用主键、操作类型和源端时间等字段实现幂等写入或去重,并监控 Offset 提交失败与延迟;不要依赖 Connector 提供 Exactly-once。

为什么部分删除操作没有对应 Kafka 事件?

先确认 skipped.operations 没有包含 d。普通行删除会输出删除事件,但带聚簇列的多行分区删除和行范围删除不会被展开为逐行事件。若业务必须捕获这些删除,应改为显式逐行删除,或通过应用层事件补充不可表示的批量删除语义。

为什么启用 beforeafter 后出现变更缺失?

Advanced 格式需要将主变更与对应 Preimage 或 Postimage 组合。确认所有目标表已经启用所需的 CDC 选项,并检查是否存在超过 cdc.incomplete.task.timeout.ms 的不完整组合错误。不要仅通过增大超时掩盖持续缺失的镜像;应先检查 Scylla CDC 配置、日志保留、读取延迟和 Connector 错误。