概述
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.user 和 scylla.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.prefix和heartbeat.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 - 重要级别:高
- 有效值 / 注意事项:应显式设置为
true或false;设置为true后才读取其他 TLS 配置。
scylla.ssl.provider
选择 TLS 实现提供方。
- 类型:
string - 默认值:
jdk - 重要级别:低
- 有效值 / 注意事项:可选
jdk、openssl或openssl_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 - 重要级别:高
- 有效值 / 注意事项:可选
legacy或advanced。legacy已弃用并计划在 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 - 重要级别:中
- 有效值 / 注意事项:可选
none、full或only-updated,只适用于advanced。full和only-updated要求所有现有目标表启用 CDC preimage。
cdc.include.after
控制 Advanced 记录中的变更后镜像。
- 类型:
string - 默认值:
none - 重要级别:中
- 有效值 / 注意事项:可选
none、full或only-updated,只适用于advanced。full和only-updated要求所有现有目标表启用 CDC postimage。
cdc.include.primary-key.placement
设置主键在输出记录中的位置。
- 类型:
list - 默认值:
kafka-key,payload-after,payload-before - 重要级别:中
- 有效值 / 注意事项:Legacy 格式只接受默认组合;Advanced 格式可使用
kafka-key、payload-after、payload-before、payload-key和kafka-headers。payload-diff虽可通过校验,但 2.0.6 未实现实际输出。ConfigDef 声明的默认值仍为kafka-key,payload-after,payload-before,但 2.0.6 使用advanced格式时必须显式设置该组合,才能生成非null的结构化 Kafka Key;省略该配置可能得到 required Key Schema 与nullKey 值。依赖按 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 - 重要级别:低
- 有效值 / 注意事项:可使用
c、u、d、t或none。默认跳过 Truncate;该配置不能使 Connector 支持本身无法转为行事件的多行分区删除或范围删除。
event.processing.failure.handling.mode
设置 Debezium 事件处理失败时的行为。
- 类型:
string - 默认值:
fail - 重要级别:中
- 有效值 / 注意事项:可选
fail、warn或ignore。fail停止任务;warn和ignore可能跳过有问题的事件,应仅在业务可接受数据缺口时使用。
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 - 重要级别:低
- 有效值 / 注意事项:使用
1到100。2.0.6 只校验正整数,未强制上限,但不应配置超过100。
worker.max.retries
设置 Scylla Worker 的总执行尝试次数。
- 类型:
int - 默认值:
20 - 重要级别:低
- 有效值 / 注意事项:接受正整数或
-1。1表示只执行一次且不重试,-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 的
STRUCTSchema 以及所选输出格式的数据类型。
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 事件包含完整的before 和 after,而 INSERT 和 DELETE 仍遵循各自可用的镜像边界。
配置示例:
only-updated。镜像记录缺失并超过超时后,相应逻辑变更会被丢弃,因此应结合错误日志和下游数据完整性要求设置超时。
首次接入时限定 CDC 历史回溯范围
适用业务场景:目标表已经持续产生 CDC 日志,但首次上线只需要最近一段时间的变化,不希望从首个仍可见的 CDC Generation 开始追赶。预期结果是没有已存 Offset 的任务从当前时间向前一小时附近开始读取,而不是生成基表快照。 配置示例:scylla.initial.lookback.ms 只在对应 Task 没有已存 Offset 时生效,不能覆盖已有恢复位置,也不能读取已经超过 Scylla CDC 保留期的数据。先根据业务允许的历史范围和预计追赶吞吐确定回溯时长;若需要基表现有行,应使用独立的数据初始化流程,并与后续 CDC 事件做好去重和衔接。
在控制 Scylla 连接压力的前提下增加并发
适用业务场景:监控显示 Connector 持续积压,Kafka 和 Worker 仍有处理余量,但单 Task 无法及时消费当前 vnode 或 tablet 分配。预期结果是在当前可分配工作单元范围内增加并发,同时让同一 Worker JVM 的 Task 共享 Session,避免连接数简单按 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.before 或 cdc.include.after,还要确认表分别启用了 Preimage 或 Postimage,并检查 Connector 与 Task 状态、Scylla 权限和错误日志。
为什么重启后看到重复事件?
Connector 使用至少一次交付。若 Kafka 已确认记录但对应 Source Offset 尚未成功刷新,重启会从上一次持久化位置恢复并再次发送部分事件。保留完整主键作为 Kafka Key,在下游使用主键、操作类型和源端时间等字段实现幂等写入或去重,并监控 Offset 提交失败与延迟;不要依赖 Connector 提供 Exactly-once。为什么部分删除操作没有对应 Kafka 事件?
先确认skipped.operations 没有包含 d。普通行删除会输出删除事件,但带聚簇列的多行分区删除和行范围删除不会被展开为逐行事件。若业务必须捕获这些删除,应改为显式逐行删除,或通过应用层事件补充不可表示的批量删除语义。
为什么启用 before 或 after 后出现变更缺失?
Advanced 格式需要将主变更与对应 Preimage 或 Postimage 组合。确认所有目标表已经启用所需的 CDC 选项,并检查是否存在超过 cdc.incomplete.task.timeout.ms 的不完整组合错误。不要仅通过增大超时掩盖持续缺失的镜像;应先检查 Scylla CDC 配置、日志保留、读取延迟和 Connector 错误。