配置

存储 Debezium 状态

qianmoQqianmoQ· 更新于 2026-09-28· 阅读 110 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

存储 Debezium 连接器的状态

概述

Debezium 连接器需要持久化存储,以便在重启之间保留其状态。所有连接器都需要一种机制来为偏移量(offset)提供持久化存储。此外,Db2、MySQL、Oracle 和 SQL Server 等连接器还需要额外的存储空间,用于保存其所谓的内部模式历史(internal schema history),以记录数据库中表结构的变化。

在 Kafka Connect 运行时内部部署时,偏移量存储由以下机制之一自动提供:

Kafka 偏移量存储为 Kafka Connect distributed 提供存储。
文件偏移量存储为 Kafka Connect standalone 提供存储。

如果你在 Debezium Engine 或 Debezium Server 中运行连接器,则必须显式配置偏移量存储。对于基于模式(schema)的数据库连接器,你需要通过设置连接器属性来配置内部模式历史存储。

Kafka

Debezium 可以使用 Kafka 来存储其状态,包括源端偏移量和模式历史。连接器实现 KafkaOffsetBackingStore,将偏移量存储在 Kafka 主题(topic)中(例如 connect-offsets)。这些偏移量确保连接器重启后能够从正确的位置继续读取。连接器将其模式历史存储在一个单独的压缩主题中,例如 schema-changes.inventory。

偏移量存储

属性默认值说明
offset.storage无默认值必须设置为 org.apache.kafka.connect.storage.KafkaOffsetBackingStore。
offset.storage.topic无默认值指定连接器存储其偏移量的 Kafka 主题。为确保该主题保留最新的偏移量信息,必须为此主题启用日志压缩。
offset.storage.partitions25指定偏移量存储主题的分区数。请确保该设置的值与 Kafka 集群的分区策略保持一致。
offset.storage.replication.factor3设置偏移量存储主题的副本因子。将数据复制到多个 broker 可以提高容错能力。

内部模式历史存储

属性默认值说明
schema.history.internal无默认值必须设置为 io.debezium.storage.kafka.history.KafkaSchemaHistory。
schema.history.internal.kafka.topic无默认值存储数据库 Schema 历史记录的主题名称。
schema.history.internal.kafka.bootstrap.servers无默认值连接器用于建立与 Kafka 集群的初始连接、以获取其数据库 Schema 历史记录的主机和端口对列表。该值必须与 Kafka Connect 进程连接 Kafka 集群时所用的连接设置保持一致。
schema.history.internal.kafka.recovery.poll.interval.ms100指定连接器在恢复期间轮询持久化数据的请求之间等待的时间,单位为毫秒。
schema.history.internal.kafka.recovery.attempts100指定连接器允许的从 Kafka 连续获取 Schema 历史数据失败的最大次数。当尝试次数超过该值后,恢复尝试将停止。连接器在无法获取数据后等待的最长时间为 recovery.attempts × recovery.poll.interval.ms。
schema.history.internal.kafka.query.timeout.ms3指定 Kafka AdminClient 提交获取集群信息的请求后,连接器等待响应的最长时间(单位为毫秒),超过该时间请求即超时。
schema.history.internal.kafka.create.timeout.ms30指定 Kafka AdminClient 提交创建 Kafka 历史主题的请求后,连接器等待响应的最长时间(单位为毫秒),超过该时间请求即超时。
schema.history.internal.producer.*无默认值用于配置生产者客户端如何与 Schema 历史主题交互的透传属性前缀。
schema.history.internal.consumer.*无默认值用于配置消费者客户端如何与 Schema 历史主题交互的透传属性前缀。

文件

可以将连接器的位置(偏移量)持久化到磁盘上的本地文件中。这些偏移量可确保连接器重启后,Debezium 能从上次读取的位置继续读取。文件存储提供了一种简单、快速的偏移量存储机制,非常适合单节点应用或测试场景。

偏移量存储

属性默认值说明
offset.storage无默认值必须设置为 org.apache.kafka.connect.storage.FileOffsetBackingStore
offset.storage.file.filename无默认值Debezium 存储源连接器偏移量的文件路径。
offset.flush.interval.ms6000ms指定将当前偏移量状态刷新到已配置的偏移量文件之间的时间间隔,单位为毫秒。

内部模式历史存储

属性默认值说明
schema.history.internal无默认值必须设置为 io.debezium.storage.file.history.FileSchemaHistory
schema.history.internal.file无默认值Debezium 记录数据库 schema 历史的文件路径。

内存

MemoryOffsetBackingStore 是 Debezium Embedded 用于跟踪源偏移量的易失性内存存储。将偏移量存储在内存中,只会在应用程序运行期间保留偏移状态。如果连接器关闭或崩溃,偏移记录将会丢失。内存存储适用于测试或短生命周期的任务,但不适用于需要持久化偏移量的生产环境。

偏移存储

属性默认值说明
offset.storage无默认值必须设置为 org.apache.kafka.connect.storage.MemoryOffsetBackingStore

内部 schema 历史存储

属性默认值说明
schema.history.internal无默认值必须设置为 io.debezium.relational.history.MemorySchemaHistory

JDBC

该存储使用任意关系数据库来存储偏移数据。你必须为该数据库提供 JDBC 驱动程序。Debezium 可以将数据存储在它捕获事件的同一源数据库中,也可以配置为使用不同的数据库。

Debezium 提供了预先配置好的 DML 和 DDL 语句。你可以使用这些默认语句,也可以用自己的语句覆盖默认语句,以便与各种数据库方言保持兼容,或针对特定用例进行定制。

偏移存储

属性默认值说明
offset.storage无默认值必须设置为 io.debezium.storage.jdbc.offset.JdbcOffsetBackingStore。
offset.storage.jdbc.connection.url无默认值用于连接数据库的 JDBC 驱动连接字符串。
offset.storage.jdbc.connection.user无默认值(可选)Debezium 连接存储偏移量数据的数据库时所使用的用户名。
offset.storage.jdbc.connection.password无默认值(可选)由 offset.storage.jdbc.connection.user 指定的用户的密码。
offset.storage.jdbc.connection.wait.retry.delay.ms3 秒(可选)指定连接器在连接偏移量存储数据库的尝试失败后,等待多久(以毫秒为单位)再重试连接。
offset.storage.jdbc.connection.retry.max.attempts5(可选)指定连接失败后,Debezium 重试连接偏移量存储数据库的最大次数。
offset.storage.jdbc.table.namedebezium_offset_storageDebezium 存储偏移量的表名。
offset.storage.jdbc.table.ddl创建查询用于创建偏移量表的 DDL 语句。
offset.storage.jdbc.table.select查询语句Debezium 用于从表中读取偏移量值的 DML 语句。
offset.storage.jdbc.table.insert插入语句Debezium 用于向表中写入偏移量的 DML 语句。
offset.storage.jdbc.table.delete删除语句Debezium 用于从表中移除偏移量的 DML 语句。

3.2 之前的已弃用配置

属性默认值说明
offset.storage无默认值必须设置为 io.debezium.storage.jdbc.offset.JdbcOffsetBackingStore。
offset.storage.jdbc.url无默认值用于连接数据库的 JDBC 驱动连接字符串。
offset.storage.jdbc.user无默认值(可选)Debezium 连接存储偏移量数据的数据库时所使用的用户名。
offset.storage.jdbc.password无默认值(可选)由 offset.storage.jdbc.user 指定的用户的密码。
offset.storage.jdbc.wait.retry.delay.ms3 秒(可选)指定连接器在连接偏移量存储数据库失败后,重试连接前等待的时间(以毫秒为单位)。
offset.storage.jdbc.retry.max.attempts5(可选)指定连接失败后,Debezium 重试连接偏移量存储数据库的最大次数。
offset.storage.jdbc.offset.table.namedebezium_offset_storageDebezium 存储偏移量的表名。
offset.storage.jdbc.offset.table.ddl创建查询用于创建偏移量表的 DDL 语句。
offset.storage.jdbc.offset.table.select查询语句用于从表中读取已存储偏移量的 DML 语句。
offset.storage.offset.table.insert插入语句用于向表中写入偏移量的 DML 语句。
offset.storage.jdbc.offset.table.delete删除语句用于从表中移除偏移量的 DML 语句。

偏移表默认值

创建查询

CREATE TABLE %s (
id VARCHAR(36)      NOT NULL,
offset_key          VARCHAR(1255),
offset_val          VARCHAR(1255),
record_insert_ts    TIMESTAMP NOT NULL,
record_insert_seq   INTEGER NOT NULL),
PRIMARY KEY (id)

Select 查询

SELECT id, offset_key, offset_val FROM %s ORDER BY record_insert_ts, record_insert_seq

Debezium插入查询

INSERT INTO %s(id, offset_key, offset_val, record_insert_ts, record_insert_seq)
    VALUES ( ?, ?, ?, ?, ? )

删除查询

DELETE FROM %s

内置架构历史存储

属性默认值描述
schema.history.internal无默认值必须设置为 io.debezium.storage.jdbc.history.JdbcSchemaHistory。
schema.history.internal.jdbc.connection.url无默认值用于连接数据库的 JDBC 驱动连接字符串。
schema.history.internal.jdbc.connection.user无默认值(可选)Debezium 连接存储模式历史数据的数据库时所使用的用户名。
schema.history.internal.jdbc.connection.password无默认值(可选)由 schema.history.internal.jdbc.connection.user 指定的用户的密码。
schema.history.internal.jdbc.connection.retry.delay.ms3 秒(可选)指定在连接内部模式历史数据库失败后,连接器重试连接前等待的时间,以毫秒为单位。
schema.history.internal.jdbc.connection.retry.max.attempts5(可选)指定连接失败后,Debezium 重试连接内部模式历史数据库的最大次数。
schema.history.internal.jdbc.table.namedebezium_database_historyDebezium 存储内部模式历史的表名。
schema.history.internal.jdbc.table.ddl创建查询用于创建存储内部模式历史的表的 DDL 语句。
schema.history.internal.jdbc.table.select查询语句用于从内部模式历史表中读取模式变更的 SELECT 语句。
schema.history.internal.jdbc.table.exists数据存在性查询用于检查内部模式历史存储表是否存在的 SELECT 语句。
schema.history.internal.jdbc.table.insert插入查询用于记录内部模式历史表变更的 INSERT 语句。

3.2 之前的弃用配置

属性默认值说明
schema.history.internal无默认值必须设置为 io.debezium.storage.jdbc.history.JdbcSchemaHistory。
schema.history.internal.jdbc.url无默认值用于连接数据库的 JDBC 驱动连接字符串。
schema.history.internal.jdbc.user无默认值(可选)Debezium 连接存储内部 Schema 历史数据的数据库时所使用的用户名。
schema.history.internal.jdbc.password无默认值(可选)由 schema.history.internal.jdbc.user 指定的用户的密码。
schema.history.internal.jdbc.retry.delay.ms3 秒(可选)指定在连接内部 Schema 历史数据库的尝试失败后,连接器等待多久(以毫秒为单位)再重试连接。
schema.history.internal.jdbc.retry.max.attempts5(可选)指定在连接内部 Schema 历史数据库失败后,Debezium 重试连接的最大次数。
schema.history.internal.jdbc.schema.history.table.namedebezium_database_historyDebezium 存储内部 Schema 历史的表名。
schema.history.internal.jdbc.schema.history.table.ddl创建查询用于创建内部 Schema 历史存储表的 DDL 语句。
schema.history.internal.jdbc.schema.history.table.select查询语句用于从内部 Schema 历史表中读取 Schema 变更的 SELECT 语句。
schema.history.internal.jdbc.schema.history.table.exists数据存在性查询用于检查内部 Schema 历史存储表是否存在的 SELECT 语句。
schema.history.internal.jdbc.schema.history.table.insert插入查询用于记录内部 Schema 历史表变更的 INSERT 语句。

历史表默认值

创建查询

CREATE TABLE %s (
    id VARCHAR(36) NOT NULL,
    history_data VARCHAR(65000),
    history_data_seq INTEGER,
    record_insert_ts TIMESTAMP NOT NULL,
    record_insert_seq INTEGER NOT NULL,
    PRIMARY KEY (id, history_data_seq)
)

选择查询

SELECT id, history_data FROM %s
    ORDER BY record_insert_ts, record_insert_seq, id, history_data_seq

数据存在性查询

SELECT * FROM %s LIMIT 1

插入查询

INSERT INTO %s(id, history_data, history_data_seq, record_insert_ts, record_insert_seq) VALUES ( ?, ?, ?, ?, ? )

Redis

Debezium 可以使用 Jedis 客户端 将数据存储在 Redis 缓存中。

Debezium 既可以使用单个 Redis 实例,也可以使用 Redis 集群模式:

单实例模式

连接到单个 Redis 服务器实例。

集群模式

连接到 Redis 集群,以实现高可用性和水平扩展。

要启用 Redis 集群模式,请将 redis.cluster.enabled 属性设置为 true,并在 redis.address 属性中提供以逗号分隔的 host:port 地址。

偏移量存储

属性默认值描述
offset.storage无默认值必须设置为 io.debezium.storage.redis.offset.RedisOffsetBackingStore
offset.storage.redis.keymetadataoffsetsDebezium 用于存储偏移量的 Redis 键。
offset.storage.redis.address无默认值Debezium 用于连接 Redis 以存储偏移量数据的 URL。
offset.storage.redis.user无默认值Debezium 用于连接 Redis 以存储偏移量数据的用户账户。
offset.storage.redis.password无默认值Debezium 用于连接 Redis 以存储偏移量数据的用户账户的密码。
offset.storage.redis.db.index0Debezium 用于访问 Redis 以存储偏移量数据的数据库索引(0—​15)。
offset.storage.redis.ssl.enabledfalse指定 Debezium 在与 Redis 通信以存储偏移量数据时是否使用 SSL。
offset.storage.redis.ssl.hostname.verification.enabledfalse指定 Debezium 在与 Redis 通信以存储偏移量数据时是否启用主机名验证。
offset.storage.redis.ssl.truststore.path无默认值用于偏移量存储的 SSL/TLS 连接到 Redis 的信任库文件的路径。如果设置该属性,Redis 连接将优先使用此属性,而不是其他配置或系统属性。
offset.storage.redis.ssl.truststore.password无默认值用于偏移量存储的 SSL/TLS 连接到 Redis 的信任库文件的密码。如果设置该属性,Redis 连接将优先使用此属性,而不是其他配置或系统属性。
offset.storage.redis.ssl.truststore.typeJKS用于偏移量存储的 SSL/TLS 连接到 Redis 的信任库文件的类型。如果设置该属性,Redis 连接将优先使用此属性,而不是其他配置或系统属性。
offset.storage.redis.ssl.keystore.path无默认值用于偏移量存储的 SSL/TLS 连接到 Redis 的密钥库文件的路径。如果设置该属性,Redis 连接将优先使用此属性,而不是其他配置或系统属性。
offset.storage.redis.ssl.keystore.password无默认值用于偏移量存储的 SSL/TLS 连接到 Redis 的密钥库文件的密码。如果设置该属性,Redis 连接将优先使用此属性,而不是其他配置或系统属性。
offset.storage.redis.ssl.keystore.typeJKS用于偏移量存储的 SSL/TLS 连接到 Redis 的密钥库文件的类型。
offset.storage.redis.connection.timeout.ms2000指定 Debezium 在与 Redis 的连接超时之前等待建立连接的时间,单位为毫秒。
offset.storage.redis.socket.timeout.ms2000指定 Debezium 在套接字超时之前允许与 Redis 交换偏移量数据的时间间隔,单位为毫秒。如果在指定的时间间隔内未能传输数据包,Debezium 将关闭该套接字。
offset.storage.redis.retry.initial.delay.ms300指定首次尝试连接 Redis 失败后,Debezium 在重试连接之前等待的时间,单位为毫秒。
offset.storage.redis.retry.max.delay.ms10000指定尝试连接 Redis 失败后,Debezium 在重试连接之前等待的最长时间,单位为毫秒。
offset.storage.redis.retry.max.attempts10指定连接尝试失败后,Debezium 重试连接 Redis 的最大次数。
offset.storage.redis.wait.enabledfalse在配置为使用副本分片的 Redis 环境中,指定 Debezium 是否等待 Redis 确认已将数据写入副本。
offset.storage.redis.wait.timeout.ms1000指定 Debezium 等待 Redis 确认已将数据写入副本分片的超时时间,单位为毫秒。
offset.storage.redis.wait.retry.enabledfalse指定 Debezium 是否重试失败的请求,以确认数据是否已写入副本分片。
offset.storage.redis.wait.retry.delay.ms1000指定请求失败后,Debezium 在重新向 Redis 提交请求以确认数据已写入副本分片之前等待的时间,单位为毫秒。
offset.storage.redis.cluster.enabledfalse如果你配置 Debezium Server 将偏移量存储在 Redis 中,请设置此属性以指定是否使用 Redis Cluster 模式。将该值设置为 true,可配置 Debezium 使用 JedisCluster 客户端将偏移量数据路由到 Redis 节点。

内部模式历史存储

属性默认值说明
schema.history.internal无默认值必须设置为 io.debezium.storage.redis.history.RedisSchemaHistory
schema.history.internal.redis.keymetadataschema_historyDebezium 用于存储模式历史数据的 Redis 键。
schema.history.internal.redis.address无默认值Debezium 连接到 Redis 以存储模式历史数据所使用的 URL。
schema.history.internal.redis.user无默认值Debezium 连接到 Redis 以存储模式历史数据所使用的用户账户。
schema.history.internal.redis.password无默认值Debezium 连接到 Redis 以存储模式历史数据所使用用户账户的密码。
schema.history.internal.redis.db.index0Debezium 访问 Redis 以存储模式历史数据时使用的数据库索引(0—​15)。
schema.history.internal.storage.redis.ssl.enabledfalse指定 Debezium 在与 Redis 通信以存储模式历史数据时是否使用 SSL。
schema.history.internal.storage.redis.ssl.hostname.verification.enabledfalse指定 Debezium 在与 Redis 通信以存储模式历史数据时是否启用主机名验证。
schema.history.internal.storage.redis.ssl.truststore.path无默认值用于通过 SSL/TLS 连接到 Redis 以存储模式历史数据的信任库文件的路径。
schema.history.internal.storage.redis.ssl.truststore.password无默认值用于通过 SSL/TLS 连接到 Redis 以存储模式历史数据的信任库文件的密码。
schema.history.internal.storage.redis.ssl.truststore.typeJKS用于通过 SSL/TLS 连接到 Redis 以存储模式历史数据的信任库文件的类型。
schema.history.internal.storage.redis.ssl.keystore.path无默认值用于通过 SSL/TLS 连接到 Redis 以存储模式历史数据的密钥库文件的路径。
schema.history.internal.storage.redis.ssl.keystore.password无默认值用于通过 SSL/TLS 连接到 Redis 以存储模式历史数据的密钥库文件的密码。
schema.history.internal.storage.redis.ssl.keystore.typeJKS用于通过 SSL/TLS 连接到 Redis 以存储模式历史数据的密钥库文件的类型。
schema.history.internal.storage.redis.connection.timeout.ms2000指定 Debezium 在与 Redis 建立连接超时前等待的毫秒数。
schema.history.internal.storage.redis.socket.timeout.ms2000指定 Debezium 与 Redis 交换模式历史数据所允许的毫秒间隔。如果在指定间隔内未能传输数据包,Debezium 将关闭套接字。
schema.history.internal.storage.redis.retry.initial.delay.ms300指定在最初尝试连接 Redis 失败后,Debezium 重试连接前等待的毫秒数。
schema.history.internal.storage.redis.retry.max.delay.ms10000指定在尝试连接 Redis 失败后,Debezium 重试连接前等待的最大毫秒数。
schema.history.internal.storage.redis.retry.max.attempts10指定在连接尝试失败后,Debezium 重试连接 Redis 的最大次数。
schema.history.internal.storage.redis.wait.enabledfalse在配置为使用副本分片的 Redis 环境中,指定 Debezium 是否等待 Redis 确认已将数据写入副本。
schema.history.internal.storage.redis.wait.timeout.ms1000指定 Debezium 等待 Redis 确认数据已写入副本分片的毫秒数,超过该时间后请求将超时。
schema.history.internal.storage.redis.wait.retry.enabledfalse指定 Debezium 是否重试失败的请求,以确认数据是否已写入副本分片。
schema.history.internal.storage.redis.wait.retry.delay.ms1000指定在失败后、Debezium 重新向 Redis 提交请求以确认数据已写入副本分片之前等待的毫秒数。
schema.history.internal.storage.redis.cluster.enabledfalse如果你配置 Debezium Server 将模式历史存储在 Redis 中,请设置此属性以指定是否使用 Redis 集群模式。将该值设置为 true,可配置 Debezium 使用 JedisCluster 客户端将历史数据路由到 Redis 节点。

Redis 配置示例

单实例模式

# Offset storage configuration
offset.storage=io.debezium.storage.redis.offset.RedisOffsetBackingStore
offset.storage.redis.address=localhost:6379
offset.storage.redis.password=password
offset.storage.redis.cluster.enabled=false

# Schema history storage configuration
schema.history.internal=io.debezium.storage.redis.history.RedisSchemaHistory
schema.history.internal.storage.redis.address=localhost:6379
schema.history.internal.storage.redis.password=password
schema.history.internal.storage.redis.cluster.enabled=false

集群模式

# Offset storage configuration
offset.storage=io.debezium.storage.redis.offset.RedisOffsetBackingStore
offset.storage.redis.address=redis-node-1:7001,redis-node-2:7002,redis-node-3:7003
offset.storage.redis.password=password
offset.storage.redis.cluster.enabled=true

# Schema history storage configuration
schema.history.internal=io.debezium.storage.redis.history.RedisSchemaHistory
schema.history.internal.storage.redis.address=redis-node-1:7001,redis-node-2:7002,redis-node-3:7003
schema.history.internal.storage.redis.password=password
schema.history.internal.storage.redis.cluster.enabled=true
如果你将 Debezium 配置为使用 Redis Cluster 模式,请确保 Redis Cluster 已正确配置且可访问。Debezium Server 实例必须能够与集群节点通信。

Amazon S3

Debezium 可以使用 Amazon S3 对象存储服务。通常,当你使用 Amazon Managed Streaming for Apache Kafka (Amazon MSK) 部署 Debezium 时,会采用 S3 存储。

内部模式历史记录存储

属性 默认值 说明

schema.history.internal

无默认值

必须设置为 io.debezium.storage.s3.history.S3SchemaHistory。

schema.history.internal.s3.access.key.id

无默认值

(可选)Debezium 用于向 S3 进行身份验证的静态访问密钥标识符。

schema.history.internal.s3.secret.access.key

无默认值

(可选)Debezium 用于向 S3 进行身份验证的 Amazon Web Services (AWS) 秘密密钥。

schema.history.internal.s3.region.name

无默认值

(可选)指定托管 S3 存储桶的区域名称。

schema.history.internal.s3.bucket.name

无默认值

指定存储模式历史记录的 S3 存储桶名称。

schema.history.internal.s3.object.name

无默认值

指定存储模式历史记录的对象在存储桶中的名称。

schema.history.internal.s3.endpoint

无默认值

(可选)指定 Debezium 用于访问 S3 服务的自定义 URL。
请按以下格式提供 URL:http://<server>:<port>;

Azure Blob Storage

Debezium 可以使用 Azure Blob 存储服务来保存数据。通常,当你在 Apache Kafka in Azure HDInsight 服务中部署 Debezium 时,会采用 Azure Blob 存储。

内部模式历史记录存储

属性默认值描述
schema.history.internal无默认值必须设置为 io.debezium.storage.azure.blob.history.AzureBlobSchemaHistory。
schema.history.internal.azure.storage.account.connectionstring无默认值指定 Azure Blob 存储连接字符串。
schema.history.internal.azure.storage.account.name无默认值Debezium 用于连接 Azure 的账户名称。
schema.history.internal.azure.storage.account.blob.endpoint无默认值可选的 Azure Blob 存储账户终结点 URL。使用它可以覆盖主权云(如 Azure Government(https://<account>.blob.core.usgovcloudapi.net)或 Azure 中国(https://<account>.blob.core.chinacloudapi.cn))的默认公共终结点(https://<account>.blob.core.windows.net)。仅在未设置 schema.history.internal.azure.storage.account.connectionstring 时使用。不应与 schema.history.internal.azure.storage.account.name 同时设置。
schema.history.internal.azure.storage.account.container.name无默认值Debezium 存储数据所使用的 Azure 容器名称。
schema.history.internal.azure.storage.blob.name无默认值Debezium 存储数据所使用的 Blob 名称。

RocketMQ

Debezium 可以使用 RocketMqSchemaHistory 类在 Apache RocketMQ 中存储和检索数据库模式变更。

内部模式历史存储

属性默认值描述
schema.history.internal无默认值必须设置为 io.debezium.storage.rocketmq.history.RocketMqSchemaHistory。
schema.history.internal.rocketmq.topic无默认值Debezium 存储数据库 schema 历史的 RocketMQ topic 名称。
schema.history.internal.rocketmq.name.srv.addr无默认值指定 Apache RocketMQ NameServer 发现服务所在的主机和端口。
schema.history.internal.rocketmq.acl.enabledfalse指定是否在 RocketMQ 中启用访问控制列表。
schema.history.internal.rocketmq.access.key无默认值指定 RocketMQ 访问密钥。如果 schema.history.internal.rocketmq.acl.enabled 设置为 true,则此字段必须包含值。
schema.history.internal.rocketmq.secret.key无默认值指定 RocketMQ 密钥。如果 schema.history.internal.rocketmq.acl.enabled 设置为 true,则此字段必须包含值。
schema.history.internal.rocketmq.recovery.attempts无默认值指定在恢复完成之前,RocketMQ 连续返回空数据的尝试次数。
schema.history.internal.rocketmq.recovery.poll.interval.ms无默认值指定 Debezium 在每次轮询尝试恢复历史记录后等待的时间(毫秒)。
schema.history.internal.rocketmq.store.record.timeout.ms无默认值指定 Debezium 在写入 RocketMQ 的操作超时前等待完成的时间(毫秒)。

Chronicle Queue

Debezium 连接器在将变更事件传递给运行时框架之前,会先将其缓存在内部队列中。默认情况下,Debezium 使用内存中的有界队列。你可以通过设置 queue.provider.type 连接器属性,将默认队列替换为基于 Chronicle Queue 的实现。Chronicle Queue 将事件写入磁盘上的内存映射文件,从而在高吞吐量期间降低堆内存压力。

目前提供两种 Chronicle Queue 实现:

chronicle

将所有变更事件溢写到磁盘。这样可以将事件完全从 JVM 堆中移除,适用于预期持续高吞吐量并希望最大限度减少垃圾回收开销的场景。

hybrid_chronicle

以内存队列作为主要缓冲区,当缓冲区达到容量上限时将溢出数据写入磁盘。在低流量场景下,不会发生序列化或磁盘 I/O。

Chronicle Queue 不支持在任何基于网络的文件系统上运行,包括 NFS、AFS、基于 SAN 的存储或类似系统。原因是这些系统无法提供内存映射文件所需的所有底层原语。只要使文件可被主机访问需要任何网络参与,就无法使用 Chronicle Queue。

JVM 设置

使用 debezium/server 或 debezium/connect 容器镜像时,只需将环境变量 ENABLE_CHRONICLE_QUEUE 设置为 true,所有必需的 JVM 配置都会自动为你完成。

使用其他部署方式时,必须通过配置显式授予 Chronicle Queue 访问若干 JVM 内部 API 的权限。这需要指定以下 JVM 启动参数,这些参数授予 Chronicle Queue 反射访问 Java 核心内部类的权限,以支持其堆外内存管理和低延迟操作所需的机制。

--add-opens=java.base/sun.nio.ch=ALL-UNNAMED
--add-opens=java.base/java.lang=ALL-UNNAMED
--add-opens=java.base/java.lang.reflect=ALL-UNNAMED
--add-opens=java.base/java.io=ALL-UNNAMED
--add-opens=java.base/java.util=ALL-UNNAMED

Chronicle Queue 提供程序

要使用 chronicle 提供程序,请将 queue.provider.type 设置为 chronicle。

属性默认值说明
queue.provider.typememory必须设置为 chronicle。
chronicle.queue.path无默认值(可选)Chronicle Queue 存储其数据文件的目录路径。如果未设置,Debezium 会创建一个临时目录,该目录在连接器停止时会被清理。

混合 Chronicle Queue 提供程序

要使用 hybrid_chronicle 提供程序,请将 queue.provider.type 设置为 hybrid_chronicle。

内存缓冲区保存由 max.queue.size 定义的容量范围内的最新事件。当缓冲区达到容量上限时,在添加新事件之前,最旧的事件会被逐出到磁盘。在轮询时,Debezium 首先从 Chronicle Queue 中取出已逐出的事件,然后从内存缓冲区中取出事件,从而保持严格的 FIFO 顺序。

属性默认值说明
queue.provider.typememory必须设置为 hybrid_chronicle。
chronicle.queue.path无默认值(可选)Chronicle Queue 存储其数据文件的目录路径。如果未设置,Debezium 会创建一个临时目录,并在连接器停止时清理该目录。
hybrid_chronicle 提供程序的内存缓冲区容量由连接器现有的 max.queue.size 属性控制。

评论

登录后参与评论

正在加载评论…