Debezium 服务器
Debezium Server
Debezium Server 是一个开箱即用的应用程序,无需 Apache Kafka 即可将变更事件从源数据库直接流式传输到消息系统和存储系统。你通过配置源连接器和接收器(sink)来定义事件的来源以及事件的投递目标。
| 如果你在使用此功能时遇到任何问题,请告知我们。此外,如果你有 Debezium Server 需要支持特定接收器的需求,或者有兴趣贡献所需的实现,也请与我们联系。 |
|---|
安装
下载并解压 Debezium Server 发行包,即可在你的系统上安装该服务器。
操作步骤
将创建一个名为 debezium-server 的目录,其内容如下:
debezium-server/
|-- config
|-- connectors
|-- debezium-server-dist-3.6.3.Final-SNAPSHOT-runner.jar
|-- jmx
|-- lib
|-- lib_metrics
|-- lib_opt
|-- run.bat
|-- run.shlib 目录用于存放依赖项,config 目录包含配置文件。run.bat 或 run.sh 文件用于启动服务器。
在启动服务器之前,请先指定你想要流式传输的更改事件的来源,以及这些事件的接收目标。有关如何配置 Debezium Server 的详细信息,请参阅配置部分。
启动服务器
运行 Debezium Server 启动脚本,即可开始从已配置的源连接器向已配置的接收端流式传输更改事件。
前置条件
- 按照 Debezium Server 配置中的说明配置 Debezium Server。
- 要使用 Oracle 连接器,请在启动服务器之前将 Oracle JDBC 驱动程序复制到
lib目录。 - 要使用 Oracle XStream,请将 Oracle JDBC 驱动程序和所需的 XStream API 文件一并复制到
lib目录。有关更多信息,请参阅获取 Oracle JDBC 驱动程序和 XStream API 文件。
操作步骤
打开命令提示符或终端窗口。
切换到
<ServerInstallDirectory>/debezium-server目录。输入以下命令之一来启动 Debezium Server:
在 Linux 上,运行
run.sh脚本:./run.sh在 Windows 上,运行
run.bat脚本:run.bat
Debezium Server 配置
Debezium Server 使用 MicroProfile Configuration 进行配置。MicroProfile 配置是一种标准,用于从多种来源向 Java 应用程序提供配置值,包括配置文件、环境变量和系统属性。
包含 $ 字符的配置属性通常由 Quarkus 用于属性展开。如果属性值以 $ 字符作为前缀,例如 ${my.property},MicroProfile 配置解析器会判定该值是一个属性表达式,并尝试展开它。
为防止解析器尝试展开使用此语法的值,请使用第二个 $ 字符对 $ 前缀进行转义;即提供值为 $${my.property}。这样该属性的值就会以 ${my.property} 的形式传递给 Debezium。
主配置文件是 config/application.properties。其中包含以下几个配置部分:
debezium.source用于源连接器配置;Debezium Server 的每个实例恰好运行一个连接器debezium.sink用于接收端系统的配置debezium.format用于输出序列化格式的配置debezium.transforms用于消息转换的配置debezium.predicates用于消息转换谓词的配置
以下示例展示了配置文件的一个片段:
debezium.sink.type=kinesis
debezium.sink.kinesis.region=eu-central-1
debezium.source.connector.class=io.debezium.connector.postgresql.PostgresConnector
debezium.source.offset.storage.file.filename=data/offsets.dat
debezium.source.offset.flush.interval.ms=0
debezium.source.database.hostname=localhost
debezium.source.database.port=5432
debezium.source.database.user=postgres
debezium.source.database.password=postgres
debezium.source.database.dbname=postgres
debezium.source.topic.prefix=tutorial
debezium.source.schema.include.list=inventory上例中的配置文件指定了以下设置:
- 接收端(sink)配置为位于
eu-central-1区域的 AWS Kinesis - 源连接器配置为 PostgreSQL,并使用默认的 Debezium decoderbufs 插件。如果使用 PostgreSQL 内置的
pgoutput插件,请设置debezium.source.plugin.name=pgoutput - 源连接器从名为
inventory的模式(schema)中捕获事件。如果要捕获数据库中的所有变更,请删除此行。否则,请修改此行,使其与你所需的模式或表相对应 - 源偏移量存储在
data目录中名为offsets.dat的文件里。为避免启动错误,你指定的目录必须在启动服务器之前已存在
服务器启动后,会生成一系列类似于以下示例的日志消息:
__ ____ __ _____ ___ __ ____ ______
--/ __ \/ / / / _ | / _ \/ //_/ / / / __/
-/ /_/ / /_/ / __ |/ , _/ ,< / /_/ /\ \
--\___\_\____/_/ |_/_/|_/_/|_|\____/___/
2020-05-15 11:33:12,189 INFO [io.deb.ser.kin.KinesisChangeConsumer] (main) Using 'io.debezium.server.kinesis.KinesisChangeConsumer$$Lambda$119/0x0000000840130c40@f58853c' stream name mapper
2020-05-15 11:33:12,628 INFO [io.deb.ser.kin.KinesisChangeConsumer] (main) Using default KinesisClient 'software.amazon.awssdk.services.kinesis.DefaultKinesisClient@d1f74b8'
2020-05-15 11:33:12,628 INFO [io.deb.ser.DebeziumServer] (main) Consumer 'io.debezium.server.kinesis.KinesisChangeConsumer' instantiated
2020-05-15 11:33:12,754 INFO [org.apa.kaf.con.jso.JsonConverterConfig] (main) JsonConverterConfig values:
converter.type = key
decimal.format = BASE64
schemas.cache.size = 1000
schemas.enable = true
2020-05-15 11:33:12,757 INFO [org.apa.kaf.con.jso.JsonConverterConfig] (main) JsonConverterConfig values:
converter.type = value
decimal.format = BASE64
schemas.cache.size = 1000
schemas.enable = false
2020-05-15 11:33:12,763 INFO [io.deb.emb.EmbeddedEngine$EmbeddedConfig] (main) EmbeddedConfig values:
access.control.allow.methods =
access.control.allow.origin =
admin.listeners = null
bootstrap.servers = [localhost:9092]
client.dns.lookup = default
config.providers = []
connector.client.config.override.policy = None
header.converter = class org.apache.kafka.connect.storage.SimpleHeaderConverter
internal.key.converter = class org.apache.kafka.connect.json.JsonConverter
internal.value.converter = class org.apache.kafka.connect.json.JsonConverter
key.converter = class org.apache.kafka.connect.json.JsonConverter
listeners = null
metric.reporters = []
metrics.num.samples = 2
metrics.recording.level = INFO
metrics.sample.window.ms = 30000
offset.flush.interval.ms = 0
offset.flush.timeout.ms = 5000
offset.storage.file.filename = data/offsets.dat
offset.storage.partitions = null
offset.storage.replication.factor = null
offset.storage.topic =
plugin.path = null
rest.advertised.host.name = null
rest.advertised.listener = null
rest.advertised.port = null
rest.extension.classes = []
rest.host.name = null
rest.port = 8083
ssl.client.auth = none
task.shutdown.graceful.timeout.ms = 5000
topic.tracking.allow.reset = true
topic.tracking.enable = true
value.converter = class org.apache.kafka.connect.json.JsonConverter
2020-05-15 11:33:12,763 INFO [org.apa.kaf.con.run.WorkerConfig] (main) Worker configuration property 'internal.key.converter' is deprecated and may be removed in an upcoming release. The specified value 'org.apache.kafka.connect.json.JsonConverter' matches the default, so this property can be safely removed from the worker configuration.
2020-05-15 11:33:12,763 INFO [org.apa.kaf.con.run.WorkerConfig] (main) Worker configuration property 'internal.value.converter' is deprecated and may be removed in an upcoming release. The specified value 'org.apache.kafka.connect.json.JsonConverter' matches the default, so this property can be safely removed from the worker configuration.
2020-05-15 11:33:12,765 INFO [org.apa.kaf.con.jso.JsonConverterConfig] (main) JsonConverterConfig values:
converter.type = key
decimal.format = BASE64
schemas.cache.size = 1000
schemas.enable = true
2020-05-15 11:33:12,765 INFO [org.apa.kaf.con.jso.JsonConverterConfig] (main) JsonConverterConfig values:
converter.type = value
decimal.format = BASE64
schemas.cache.size = 1000
schemas.enable = true
2020-05-15 11:33:12,767 INFO [io.deb.ser.DebeziumServer] (main) Engine executor started
2020-05-15 11:33:12,773 INFO [org.apa.kaf.con.sto.FileOffsetBackingStore] (pool-3-thread-1) Starting FileOffsetBackingStore with file data/offsets.dat
2020-05-15 11:33:12,835 INFO [io.deb.con.com.BaseSourceTask] (pool-3-thread-1) Starting PostgresConnectorTask with configuration:
2020-05-15 11:33:12,837 INFO [io.deb.con.com.BaseSourceTask] (pool-3-thread-1) connector.class = io.debezium.connector.postgresql.PostgresConnector
2020-05-15 11:33:12,837 INFO [io.deb.con.com.BaseSourceTask] (pool-3-thread-1) offset.flush.interval.ms = 0
2020-05-15 11:33:12,838 INFO [io.deb.con.com.BaseSourceTask] (pool-3-thread-1) database.user = postgres
2020-05-15 11:33:12,838 INFO [io.deb.con.com.BaseSourceTask] (pool-3-thread-1) database.dbname = postgres
2020-05-15 11:33:12,838 INFO [io.deb.con.com.BaseSourceTask] (pool-3-thread-1) offset.storage.file.filename = data/offsets.dat
2020-05-15 11:33:12,838 INFO [io.deb.con.com.BaseSourceTask] (pool-3-thread-1) database.hostname = localhost
2020-05-15 11:33:12,838 INFO [io.deb.con.com.BaseSourceTask] (pool-3-thread-1) database.password = ********
2020-05-15 11:33:12,839 INFO [io.deb.con.com.BaseSourceTask] (pool-3-thread-1) name = kinesis
2020-05-15 11:33:12,839 INFO [io.deb.con.com.BaseSourceTask] (pool-3-thread-1) topic.prefix = tutorial
2020-05-15 11:33:12,839 INFO [io.deb.con.com.BaseSourceTask] (pool-3-thread-1) database.port = 5432
2020-05-15 11:33:12,839 INFO [io.deb.con.com.BaseSourceTask] (pool-3-thread-1) schema.include.list = inventory
2020-05-15 11:33:12,908 INFO [io.quarkus] (main) debezium-server 1.2.0-SNAPSHOT (powered by Quarkus 1.4.1.Final) started in 1.198s. Listening on: http://0.0.0.0:8080
2020-05-15 11:33:12,911 INFO [io.quarkus] (main) Profile prod activated.
2020-05-15 11:33:12,911 INFO [io.quarkus] (main) Installed features: [cdi, smallrye-health]源配置
源配置属性用于指定 Debezium Server 运行的连接器,并配置位于 Kafka Connect 之外的位点(offset)与 schema 历史存储。
大多数源配置属性都与各个连接器文档中所记录的属性相对应。若要在 Debezium Server 中使用这些属性,请在连接器属性名称前加上 debezium.source. 前缀。Debezium Server 还提供了一些额外的源配置属性,用于设置 Kafka Connect 通常会自动处理的项。例如,debezium.source.offset.storage 和 debezium.source.offset.storage.file.filename 用于配置基于文件的位点存储。
属性 默认值 说明
debezium.source.connector.class
实现源连接器的 Java 类的名称。
debezium.source.offset.storage
org.apache.kafka.connect.storage.FileOffsetBackingStore
在非 Kafka 部署中用于存储和检索位点的类。可选项如下:
org.apache.kafka.connect.storage.FileOffsetBackingStore,用于非 Kafka 部署org.apache.kafka.connect.storage.MemoryOffsetBackingStore,用于测试环境的易失性(内存)存储io.debezium.storage.jdbc.offset.JdbcOffsetBackingStore,用于通过 JDBC 访问的数据库io.debezium.storage.redis.offset.RedisOffsetBackingStore,用于 Redis 部署
debezium.source.offset.storage.file.filename
若使用文件位点存储(默认),在非 Kafka 部署中用于存储连接器位点的文件路径。
debezium.source.offset.flush.interval.ms
定义位点写入文件的频率。
debezium.source.offset.storage.redis.address
(可选)若使用 Redis 存储位点,指定 Redis 目标流的地址,格式为 host:port。如果未提供,将尝试读取 debezium.sink.redis.address。
debezium.source.offset.storage.redis.user
(可选)若使用 Redis 存储位点,指定用于与 Redis 通信的用户名。如果未提供 redis.address 配置,且 redis.address 取自 Redis sink,则会尝试从 debezium.sink.redis.user 加载该值。
debezium.source.offset.storage.redis.password
(可选)如果使用 Redis 存储偏移量,则用于与 Redis 通信的密码(对应用户的密码)。如果设置了用户,则必须设置密码。如果未提供 redis.address 配置,且 redis.address 取自 Redis sink,则会尝试从 debezium.sink.redis.password 加载该值。
debezium.source.offset.storage.redis.ssl.enabled
(可选)如果使用 Redis 存储偏移量,则表示与 Redis 通信时是否使用 SSL。如果未提供 redis.address 配置,且 redis.address 取自 Redis sink,则会尝试从 debezium.sink.redis.ssl.enabled 加载该值。默认为 false。
debezium.source.offset.storage.redis.ssl.hostname.verification.enabled
(可选)如果使用 Redis 存储偏移量,则表示与 Redis 通信时是否启用主机名验证。如果未提供 redis.address 配置,且 redis.address 取自 Redis sink,则会尝试从 debezium.sink.redis.ssl.hostname.verification.enabled 加载该值。默认为 false。
debezium.source.offset.storage.redis.ssl.truststore.path
(可选)如果使用 Redis 并启用 SSL 存储偏移量,则表示信任库(trust store)文件的路径。如果设置该值,Redis 连接将优先使用此属性,而不是其他配置或系统属性。
debezium.source.offset.storage.redis.ssl.truststore.password
(可选)如果使用 Redis 并启用 SSL 存储偏移量,则表示信任库文件的密码。如果设置该值,Redis 连接将优先使用此属性,而不是其他配置或系统属性。
debezium.source.offset.storage.redis.ssl.truststore.type
JKS
(可选)如果使用 Redis 并启用 SSL 存储偏移量,则表示信任库文件的类型。如果设置该值,Redis 连接将优先使用此属性,而不是其他配置或系统属性。
debezium.source.offset.storage.redis.ssl.keystore.path
(可选)如果使用 Redis 存储偏移量并启用了 SSL,则为密钥库文件的路径。如果设置了该属性,Redis 连接将优先使用此属性,而不是其他配置或系统属性。
debezium.source.offset.storage.redis.ssl.keystore.password
(可选)如果使用 Redis 存储偏移量并启用了 SSL,则为密钥库文件的密码。如果设置了该属性,Redis 连接将优先使用此属性,而不是其他配置或系统属性。
debezium.source.offset.storage.redis.ssl.keystore.type
JKS
(可选)如果使用 Redis 存储偏移量并启用了 SSL,则为密钥库文件的类型。如果设置了该属性,Redis 连接将优先使用此属性,而不是其他配置或系统属性。
debezium.source.offset.storage.redis.key
(可选)如果使用 Redis 存储偏移量,用于定义 Redis 中的哈希键。如果没有提供 redis.key 配置,则默认值为 metadata:debezium:offsets
debezium.source.offset.storage.redis.wait.enabled
false
如果使用 Redis 存储偏移量,启用对副本的等待。当 Redis 配置了副本分片时,这允许验证数据已写入副本。有关更多信息,请参阅 Redis 的 WAIT 命令。
debezium.source.offset.storage.redis.wait.timeout.ms
1000
如果使用 Redis 存储偏移量,用于定义等待副本时的超时时间(毫秒)。该值必须为正数。
debezium.source.offset.storage.redis.wait.retry.enabled
false
如果使用 Redis 存储偏移量,启用在等待副本失败时进行重试。
debezium.source.offset.storage.redis.wait.retry.delay.ms
1000
如果使用 Redis 存储偏移量,用于定义在等待副本失败时重试的延迟时间。
debezium.source.offset.storage.redis.cluster.enabled
false
如果你将 Debezium Server 配置为把偏移量存储在 Redis 中,可通过此属性指定是否使用 Redis 集群模式。将其设置为 true,可配置 Debezium 使用 JedisCluster 客户端将偏移数据路由到 Redis 节点。
debezium.source.schema.history.internal
io.debezium.storage.kafka.history.KafkaSchemaHistory
部分连接器(如 MySQL、SQL Server、Db2、Oracle)会跟踪数据库 schema 随时间的演进,并将该数据存储在数据库 schema 历史记录中。默认情况下基于 Kafka 存储,同时还提供其他选项:
io.debezium.storage.file.history.FileSchemaHistory,用于非 Kafka 部署io.debezium.relational.history.MemorySchemaHistory,用于测试环境的易失性存储io.debezium.storage.redis.history.RedisSchemaHistory,用于 Redis 部署io.debezium.storage.rocketmq.history.RocketMqSchemaHistory,用于 RocketMQ 部署io.debezium.storage.azure.blob.history.AzureBlobSchemaHistory,用于 Azure Blob Storage 部署
debezium.source.schema.history.internal.file.filename
FileSchemaHistory 持久化其数据所使用的文件的名称和位置。
debezium.source.schema.history.internal.redis.address
使用 RedisSchemaHistory 时要连接的 Redis 主机:端口。
debezium.source.schema.history.internal.redis.user
使用 RedisSchemaHistory 时要使用的 Redis 用户。
debezium.source.schema.history.internal.redis.password
使用 RedisSchemaHistory 时要使用的 Redis 密码。
debezium.source.schema.history.internal.redis.ssl.enabled
使用 RedisSchemaHistory 时是否使用 SSL 连接。
debezium.source.schema.history.internal.redis.ssl.hostname.verification.enabled
使用 RedisSchemaHistory 时是否启用主机名验证。
debezium.source.schema.history.internal.redis.ssl.truststore.path
(可选)如果使用 Redis 存储模式历史记录并启用了 SSL,信任库文件的路径。如果设置了该属性,Redis 连接将优先使用此属性,而不是其他配置或系统属性。
debezium.source.schema.history.internal.redis.ssl.truststore.password
(可选)如果使用 Redis 存储模式历史记录并启用了 SSL,信任库文件的密码。如果设置了该属性,Redis 连接将优先使用此属性,而不是其他配置或系统属性。
debezium.source.schema.history.internal.redis.ssl.truststore.type
JKS
(可选)如果使用 Redis 存储模式历史记录并启用了 SSL,信任库文件的类型。如果设置了该属性,Redis 连接将优先使用此属性,而不是其他配置或系统属性。
debezium.source.schema.history.internal.redis.ssl.keystore.path
(可选)如果使用 Redis 存储模式历史记录并启用了 SSL,密钥库文件的路径。如果设置了该属性,Redis 连接将优先使用此属性,而不是其他配置或系统属性。
debezium.source.schema.history.internal.redis.ssl.keystore.password
(可选)如果使用 Redis 存储模式历史记录并启用了 SSL,密钥库文件的密码。如果设置了该属性,Redis 连接将优先使用此属性,而不是其他配置或系统属性。
debezium.source.schema.history.internal.redis.ssl.keystore.type
JKS
(可选)如果使用 Redis 存储模式历史记录并启用了 SSL,密钥库文件的类型。如果设置了该属性,Redis 连接将优先使用此属性,而不是其他配置或系统属性。
debezium.source.schema.history.internal.redis.key
如果使用 RedisSchemaHistory,用于存储的 Redis 键。默认值:metadataschema_history
debezium.source.schema.history.internal.redis.retry.initial.delay.ms
如果使用 RedisSchemaHistory,在需要重连 Redis 时的初始延迟时间。默认值:300(毫秒)
debezium.source.schema.history.internal.redis.retry.max.delay.ms
使用 RedisSchemaHistory 时,重连 Redis 的最大延迟时间。默认值:10000(毫秒)
debezium.source.schema.history.internal.redis.retry.max.attempts
连接 Redis 的最大重试次数。默认值:10
debezium.source.schema.history.internal.redis.connection.timeout.ms
使用 RedisSchemaHistory 时,Redis 客户端的连接超时时间。默认值:2000(毫秒)
debezium.source.schema.history.internal.redis.socket.timeout.ms
使用 RedisSchemaHistory 时,Redis 客户端的套接字超时时间。默认值:2000(毫秒)
debezium.source.schema.history.internal.redis.wait.enabled
false
如果使用 Redis 存储架构历史记录,则启用等待副本功能。当 Redis 配置了副本分片时,该功能可用于确认数据已写入副本。有关更多信息,请参阅 Redis 的 WAIT 命令。
debezium.source.schema.history.internal.redis.wait.timeout.ms
1000
如果使用 Redis 存储架构历史记录,则定义等待副本时的超时时间(毫秒)。该值必须为正数。
debezium.source.schema.history.internal.redis.wait.retry.enabled
false
如果使用 Redis 存储架构历史记录,则在等待副本失败时启用重试。
debezium.source.schema.history.internal.redis.wait.retry.delay.ms
1000
如果使用 Redis 存储架构历史记录,则定义等待副本失败时的重试延迟时间。
debezium.source.schema.history.internal.redis.cluster.enabled
false
如果你配置 Debezium Server 将架构历史记录存储在 Redis 中,可通过该属性指定是否使用 Redis Cluster 模式。将其设置为 true,可配置 Debezium 使用 JedisCluster 客户端将历史数据路由到 Redis 节点。
debezium.source.schema.history.internal.rocketmq.topic
用于数据库 Schema 历史的 RocketMQ 主题名称。
debezium.source.schema.history.internal.rocketmq.name.srv.addr
localhost:9876
RocketMQ 服务发现 NameServer 地址配置。
debezium.source.schema.history.internal.rocketmq.acl.enabled
false
RocketMQ 访问控制启用配置,默认为 false。
debezium.source.schema.history.internal.rocketmq.access.key
RocketMQ access key。如果 debezium.source.schema.history.internal.rocketmq.acl.enabled 为 true,则该值不能为空。
debezium.source.schema.history.internal.rocketmq.secret.key
RocketMQ secret key。如果 debezium.source.schema.history.internal.rocketmq.acl.enabled 为 true,则该值不能为空。
debezium.source.schema.history.internal.rocketmq.recovery.attempts
60
恢复数据库 Schema 历史的最大尝试次数。
debezium.source.schema.history.internal.rocketmq.recovery.poll.interval.ms
1000
恢复过程中轮询已持久化数据时的等待毫秒数。
debezium.source.schema.history.internal.rocketmq.store.record.timeout.ms
60000
向 RocketMQ 发送消息的超时时间。
debezium.source.schema.history.internal.azure.storage.account.connectionstring
Azure Blob Storage 账户连接字符串。
debezium.source.schema.history.internal.azure.storage.account.name
Azure Blob Storage 账户名称。仅当 debezium.source.schema.history.internal.azure.storage.account.connectionstring 为空时才应设置该值,此时将使用 Azure Active Directory 进行身份验证。
debezium.source.schema.history.internal.azure.storage.account.blob.endpoint
可选的 Azure Blob 存储账户终结点 URL。当使用主权云(如 Azure Government(https://<account>.blob.core.usgovcloudapi.net)或 Azure China(https://<account>.blob.core.chinacloudapi.cn))时,可用它覆盖默认的公共终结点(https://<account>.blob.core.windows.net)。仅在未设置 debezium.source.schema.history.internal.azure.storage.account.connectionstring 时使用。不应与 debezium.source.schema.history.internal.azure.storage.account.name 同时设置。
debezium.source.schema.history.internal.azure.storage.account.container.name
Azure Blob 存储账户容器名称。
debezium.source.schema.history.internal.azure.storage.blob.name
用于持久化架构历史数据的 Azure Blob 存储 blob 名称。
格式配置
配置 Debezium Server 如何序列化消息键和消息值。你可以使用默认的 JSON 格式或自定义的 Kafka Connect Converter 实现,独立配置键和值的格式。
| 属性 | 默认值 | 描述 |
|---|---|---|
debezium.format.key | json | 键的输出格式名称,可选值为 json/jsonbytearray/avro/protobuf/simplestring/binary。 |
debezium.format.key.* | 传递给键转换器的配置属性。 | |
debezium.format.value | json | 值的输出格式名称,可选值为 json/jsonbytearray/avro/protobuf/cloudevents/simplestring/binary。 |
debezium.format.value.* | 传递给值转换器的配置属性。 | |
debezium.format.header | json | 头部的输出格式名称,可选值为 json/jsonbytearray。 |
debezium.format.header.* | 传递给头部转换器的配置属性。 |
转换配置
配置 Kafka Connect 单消息转换(SMT),Debezium Server 会在将消息投递到接收器之前应用这些转换。请指定要应用的转换以及每种转换的配置选项。
该服务器支持 Kafka Connect 定义的单消息转换。
| 属性 | 默认值 | 说明 [id="debezium-transforms"] |
|---|---|---|
debezium.transforms | 以逗号分隔的转换(transformation)符号名称列表。 | |
debezium.transforms.<name>.type | 实现名为 <name> 的转换的 Java 类名称。 | |
debezium.transforms.<name>.* | 传递给名为 <name> 的转换的配置属性。 | |
debezium.transforms.<name>.predicate | 要应用于名为 <name> 的转换的谓词名称。 | |
debezium.transforms.<name>.negate | false | 确定是否对应用于名为 <name> 的转换的谓词结果取反。 |
谓词配置
配置用于评估消息的谓词,以确定 Debezium Server 何时应用消息转换。Debezium Server 支持 Kafka Connect 的用于条件性单消息转换(SMT)的谓词框架。谓词属性用于指定要使用的谓词、其实现类及其配置选项。
| 属性 | 默认值 | 说明 [id="debezium-predicates"] |
|---|---|---|
debezium.predicates | 以逗号分隔的谓词符号名称列表。 | |
debezium.predicates.<name>.type | 实现名为 <name> 的谓词的 Java 类名称。 | |
debezium.predicates.<name>.* | 传递给名为 <name> 的谓词的配置属性。 |
启用消息过滤
你可以在 Debezium Server 中启用 filter SMT,根据你定义的基于表达式的条件,有选择地传递或丢弃变更事件。
由于过滤表达式是在 Debezium Server 进程内运行的脚本引擎中执行的,因此只有在服务器配置的访问权限仅限于授权运维人员的环境中,启用此功能才是安全的。脚本支持因此默认处于禁用状态,必须显式启用。
要启用消息过滤,请在启动 Debezium Server 之前将环境变量 ENABLE_DEBEZIUM_SCRIPTING 设置为 true。启用此设置后,启动脚本会将 debezium-scripting JAR 以及 JSR 223 脚本引擎实现(Groovy 和 GraalVM JavaScript)从 Debezium Server 发行版的 opt_lib 目录添加到服务器类路径中。这些库之所以保留在 opt_lib 中而不是放在默认类路径上,正是为了确保在运维人员明确决定引入它们之前,脚本功能始终不可用。
其他资源
异步引擎属性
默认情况下,Debezium Server 使用异步嵌入式引擎(AsyncEmbeddedEngine)来处理变更事件。使用以下属性可配置记录处理线程、任务生命周期管理以及引擎关闭行为。
| 属性 | 默认值 | 描述 |
|---|---|---|
record.processing.threads | 根据工作负载和可用 CPU 核心数按需分配线程。 | 可用于处理变更事件记录的线程数。如果未指定值(即默认情况),引擎会使用 Java 的 ThreadPoolExecutor 根据当前工作负载动态调整线程数,最大线程数为给定机器的 CPU 核心数。如果指定了值,引擎会使用 Java 的 固定线程池 方法创建具有指定线程数的线程池。要使用给定机器上的所有可用核心,请设置占位符值 AVAILABLE_CORES。 |
record.processing.shutdown.timeout.ms | 1000 | 引擎在任务关闭被调用后,允许处理待处理记录的最长时间,以毫秒为单位。 |
task.management.timeout.ms | 180,000(3 分钟) | 引擎等待任务生命周期管理操作(启动和停止)完成的时间,以毫秒为单位。 |
附加配置
Debezium Server 运行在 Quarkus 框架之上。Quarkus 暴露的所有配置属性在 Debezium Server 中均可用。下表列出了用于配置 HTTP 端点访问和日志行为的常用 Quarkus 属性。
| 属性 | 默认值 | 说明 [id="debezium-quarkus-http-port"] |
|---|---|---|
quarkus.http.port | 8080 | Debezium 暴露 Microprofile Health 端点及其他状态信息所使用的端口。可通过 http://host:8080/q/health 访问健康检查信息。 |
quarkus.log.level | INFO | 每个日志类别的默认日志级别。 |
quarkus.log.console.json | true | 决定是否启用 JSON 控制台格式化扩展,该扩展会禁用"常规"控制台格式化。 |
可以通过在 config/application.properties 文件中设置 quarkus.log.console.json=false 来禁用 JSON 日志。详情请参见该文件中的示例。
Sink 配置
Debezium Server 支持将变更事件投递到不同 sink 目的地的各种实现。请配置 Debezium Server 将变更事件投递到的目的地。通过设置 debezium.sink.type 属性来指定 sink 实现,然后配置适用于所选 sink 的各项属性。
Amazon Kinesis
Amazon Kinesis 是一种数据流式传输系统的实现,支持流分片等实现高扩展性的技术。Kinesis 暴露了一组 REST API,并提供(不仅限于)用于实现该 sink 的 Java SDK。
| 属性 | 默认值 | 描述 |
|---|---|---|
debezium.sink.type | 必须设置为 kinesis。 | |
debezium.sink.kinesis.region | Kinesis 目标流所在的区域名称。 | |
debezium.sink.kinesis.endpoint | 由 aws sdk 决定的端点 | (可选)提供 Kinesis 目标流的端点 URL。 |
debezium.sink.kinesis.credentials.profile | (可选)通过默认凭据配置文件与 Amazon API 通信时使用的凭据配置文件名称。若未提供,则使用默认凭据提供程序链。它将按以下顺序查找凭据:环境变量、Java 系统属性、Web 身份令牌凭据、默认凭据配置文件、Amazon ECS 容器凭据以及实例配置文件凭据。 | |
debezium.sink.kinesis.null.key | default | Kinesis 不支持没有键的消息。因此,对于没有主键的表产生的消息,将使用该字符串作为消息键。 |
注入点
可以通过自定义逻辑为特定功能提供替代实现,从而修改 Kinesis 接收器的行为。当不存在替代实现时,将使用默认实现。
| 接口 | CDI 分类器 | 描述 |
|---|---|---|
software.amazon.awssdk.services.kinesis.KinesisClient | @CustomConsumerBuilder | 自定义配置的 KinesisClient 实例,用于向目标流发送消息。 |
io.debezium.server.StreamNameMapper | 自定义实现将计划的目标(主题)名称映射为物理 Kinesis 流名称。默认情况下使用相同的名称。 |
Google Cloud Pub/Sub
Google Cloud Pub/Sub 是一种消息传递/事件系统实现,专为可扩展的批处理和流处理应用而设计。Pub/Sub 提供了一组 REST API,并提供了一个(不仅限于)Java 的 SDK,接收器即通过该 SDK 实现。
| 属性 | 默认值 | 描述 |
|---|---|---|
debezium.sink.type | 必须设置为 pubsub。 | |
debezium.sink.pubsub.project.id | 系统级默认项目 ID | 创建目标主题所在的项目名称。 |
debezium.sink.pubsub.ordering.enabled | true | Pub/Sub 可以选择使用消息键,以保证具有相同顺序键的消息按照与发送时相同的顺序 传递。该功能可以禁用。 |
debezium.sink.pubsub.null.key | default | 没有主键的表会发送键为 null 的消息。Pub/Sub 不支持这种情况,因此必须使用替代键。 |
debezium.sink.pubsub.batch.delay.threshold.ms | 100 | 在将未发送的消息发布到 Pub/Sub 之前,等待达到元素数量或请求字节阈值的最长时间。 |
debezium.sink.pubsub.batch.element.count.threshold | 100L | 一旦队列中积累了这么多消息,就将所有消息通过一次调用发送出去,即使延迟阈值尚未到期也是如此。 |
debezium.sink.pubsub.batch.request.byte.threshold | 10000000L | 一旦批处理请求中的字节数达到此阈值,就将所有消息通过一次调用发送出去,即使延迟阈值和消息数量阈值都尚未超过也是如此。 |
debezium.sink.pubsub.flowcontrol.enabled | false | 启用后,会为发布客户端配置流量控制,以限制发布请求的速率。 |
debezium.sink.pubsub.flowcontrol.max.outstanding.messages | Long.MAX_VALUE | (可选)若启用了流量控制,则在消息被阻止发布之前允许的最大消息数量 |
debezium.sink.pubsub.flowcontrol.max.outstanding.bytes | Long.MAX_VALUE | (可选)若启用了流量控制,则在消息被阻止发布之前允许的最大字节数 |
debezium.sink.pubsub.retry.total.timeout.ms | 60000 | 向 Pub/Sub 发布(包括重试)调用的总超时时间。 |
debezium.sink.pubsub.retry.initial.delay.ms | 5 | 重试请求前的初始等待时间。 |
debezium.sink.pubsub.retry.delay.multiplier | 2.0 | 将上一次的等待时间乘以该乘数,得出下一次的等待时间,直到达到最大值为止。 |
debezium.sink.pubsub.retry.max.delay.ms | Long.MAX_VALUE | 重试前的最长等待时间。即达到该值后,等待时间将不再按乘数递增。 |
debezium.sink.pubsub.retry.initial.rpc.timeout.ms | 10000 | 控制初始远程过程调用(RPC)的超时时间 |
debezium.sink.pubsub.retry.rpc.timeout.multiplier | 2.0 | 将上一次的 RPC 超时时间乘以该乘数,得出下一次的 RPC 超时值,直到达到最大值为止 |
debezium.sink.pubsub.retry.max.rpc.timeout.ms | 10000 | 向 Cloud Pub/Sub 发布单个请求的最大超时时间。 |
debezium.sink.pubsub.wait.message.delivery.timeout.ms | 30000 | 检索向 Cloud Pub/Sub 发布请求结果的最长等待时间。 |
debezium.sink.pubsub.concurrency.threads | 0 | 客户端库用于发布消息的线程数。设置为 0 时禁用。 |
debezium.sink.pubsub.compression.threshold.bytes | -1 | 超过该字节阈值的消息将被压缩后再传输。设置为 -1 时禁用。 |
debezium.sink.pubsub.address | Pub/Sub 模拟器的地址。仅用于开发或测试环境,需配合 pubsub 模拟器 使用。除非设置了该值,否则 debezium-server 将连接到运行在 GCP 项目中的 Cloud Pub/Sub 实例,这也是生产环境所期望的行为。 | |
debezium.sink.pubsub.region | 要连接到的 Google Cloud 区域(例如 us-central1、asia-northeast1)。指定后,Debezium 将使用格式为 {region}-pubsub.googleapis.com:443 的 Pub/Sub 区域端点,从而连接到区域端点而非全局端点。注意:如果指定了 debezium.sink.pubsub.address,则此参数将被忽略。 |
注入点
可以通过自定义逻辑为特定功能提供替代实现,从而修改 Pub/Sub 接收器的行为。当不存在替代实现时,将使用默认实现。
| 接口 | CDI 分类器 | 说明 |
|---|---|---|
io.debezium.server.pubsub.PubSubChangeConsumer.PublisherBuilder | @CustomConsumerBuilder | 提供自定义配置的 Publisher 实例的类,用于向指定主题发送消息。 |
io.debezium.server.StreamNameMapper | 自定义实现将计划的目标(主题)名称映射为物理 Pub/Sub 主题名称。默认情况下使用相同的名称。 |
Pub/Sub Lite
Google Cloud Pub/Sub Lite 是 Google Cloud Pub/Sub 的一种高性价比替代方案。Pub/Sub 暴露了一组 REST API,并提供了一个(不仅是 Java 的)SDK,本接收器正是基于该 SDK 实现的。
| 属性 | 默认值 | 描述 |
|---|---|---|
debezium.sink.type | 必须设置为 pubsublite | |
debezium.sink.pubsublite.project.id | 系统范围的默认项目 ID | 创建目标主题所在的项目名称或项目 ID。 |
debezium.sink.pubsublite.region | 创建主题所在的区域。例如 us-east1-b。 | |
debezium.sink.pubsublite.ordering.enabled | true | Pub/Sub Lite 可以选择使用消息键来保证具有相同键的消息被投递到同一分区。可以禁用此功能。 |
debezium.sink.pubsublite.null.key | default | 没有主键的表会发送 null 键的消息。Pub/Sub Lite 不支持这种方式,因此必须使用替代键。 |
debezium.sink.pubsublite.wait.message.delivery.timeout.ms | 30000 | 检索发布请求在 Cloud Pub/Sub 中的投递结果所允许的最长等待时间。 |
注入点
通过自定义逻辑为特定功能提供替代实现,可以修改 Pub/Sub Lite 接收器的行为。如果不存在替代实现,则使用默认实现。
| 接口 | CDI 分类器 | 说明 |
|---|---|---|
io.debezium.server.pubsub.PubSubLiteChangeConsumer.PublisherBuilder | @CustomConsumerBuilder | 该类为用于向专用主题发送消息的 Publisher 提供经过自定义配置的实例。 |
io.debezium.server.StreamNameMapper | 自定义实现将计划的目标(主题)名称映射为实际的 Pub/Sub Lite 主题名称。默认情况下使用相同的名称。 |
HTTP 客户端
配置 HTTP 客户端接收器,可将变更事件从 Debezium Server 传递到 HTTP 端点以供下游处理。HTTP 客户端接收器支持可选的 JWT 和 OAuth 2.0 客户端凭据认证,并提供了配置批量事件传递的选项。
您可以使用 HTTP 客户端接收器将 Debezium Server 与 Knative 集成,使其能够充当 Knative 事件源。
属性 默认值 说明
必须设置为 http
用于流式传输事件的 HTTP 服务器 URL。也可以通过定义 K_SINK 环境变量来设置,该变量由 Knative source 框架使用。
60000
在超时之前等待服务器响应的秒数。(默认为 60 秒)
5
抛出异常前的重试次数(默认 5 次)。
debezium.sink.http.retry.interval.ms
1000
发送记录失败后,再次尝试发送前等待的毫秒数(默认为 1 秒)。
debezium.sink.http.headers.prefix
X-DEBEZIUM-
请求头将加上此前缀(默认为 X-DEBEZIUM-)。
debezium.sink.http.headers.encode.base64
true
请求头的值将进行 base64 编码(默认为 true)。
debezium.sink.http.batch.enabled
false
设置为 true 时,会将一批中的所有变更事件聚合为一个 JSON 数组,并通过单个 HTTP POST 请求发送,而不是逐个发送每个事件。
debezium.sink.http.batch.max-size
200
启用批量模式时,每个 HTTP 请求所包含的最大事件数量。如果引擎投递的事件数量超过该限制,这些事件将被拆分为多个请求发送。
debezium.sink.http.authentication.type
指定 HTTP 客户端接收器连接 HTTP 服务器时使用的认证类型。支持以下选项之一:
jwt
JSON Web Token(JWT)认证。
oauth2
OAuth2 客户端凭据授权(RFC 6749 第 4.4 节)。
standard-webhooks
如果省略此属性,HTTP 客户端接收器在连接时不会使用认证请求头。
debezium.sink.http.authentication.jwt.username
指定 JWT 认证使用的用户名。
debezium.sink.http.authentication.jwt.password
指定 JWT 认证使用的密码。
debezium.sink.http.authentication.jwt.url
指定 JWT 认证使用的基础 URL(例如 http://myserver:8000/)。JWT 初始认证及刷新认证的 REST 请求会在其后追加 auth/authenticate 和 auth/refreshToken 路径。
debezium.sink.http.authentication.jwt.token_expiration
认证令牌过期前请求的时长(以分钟计)。
debezium.sink.http.authentication.jwt.refresh_token_expiration
刷新令牌过期前请求的时长(以分钟计)。
debezium.sink.http.authentication.webhook.secret
Debezium 用于为 Webhook 请求生成 HMAC-SHA256 签名的 Webhook 签名密钥。该密钥必须采用 Base64 编码,大小为 24 到 64 字节(192–512 位)。你可以在密钥前添加前缀 whsec_,以便将其与其他类型的密钥或令牌区分开来。有关实现或验证 Webhook 签名的更多信息,请参阅 Standard Webhooks 规范。
debezium.sink.http.authentication.oauth2.client_id
指定客户端凭据授予(client credentials grant)所使用的 OAuth2 客户端 ID。
debezium.sink.http.authentication.oauth2.client_secret
指定客户端凭据授予(client credentials grant)所使用的 OAuth2 客户端密钥。
debezium.sink.http.authentication.oauth2.token_url
OAuth2 令牌端点的 URL(例如 https://auth.example.com/oauth/token)。
debezium.sink.http.authentication.oauth2.scope
可选的 OAuth2 作用域列表,以空格分隔(例如 data:read data:write)。
debezium.sink.http.authentication.oauth2.client_auth_method
client_secret_basic
指定客户端凭据发送到令牌端点的方式。支持以下选项之一:
client_secret_basic
以 HTTP Basic 授权头的形式发送凭据(RFC 6749 第 2.3.1 节)。
client_secret_post
在 POST 请求体中以表单字段的形式发送 client_id 和 client_secret。
debezium.sink.http.authentication.oauth2.token_url.http_method
POST
用于令牌请求的 HTTP 方法。OAuth2 规范要求使用 POST,但有些提供方也接受通过查询参数使用 GET。仅当提供方不支持标准的 POST 方法时,才将此值设置为 GET。
debezium.sink.http.authentication.oauth2.params.*
可选的额外参数,包含在令牌请求中。例如,某些提供方要求提供 audience 或 resource 参数。使用 params. 前缀来设置这些参数,例如 debezium.sink.http.authentication.oauth2.params.audience=\https://my-api.example.com。
Apache Pulsar
Apache Pulsar 是一款高性能、低延迟的服务器间消息传递平台。Pulsar 提供了 REST API 和一个原生端点,并提供(不仅限于)Java 客户端,我们用它来实现 sink。
| 属性 | 默认值 | 描述 |
|---|---|---|
debezium.sink.type | 必须设置为 pulsar。 | |
debezium.sink.pulsar.timeout | 0 | 配置将一批消息发送到 Pulsar 并等待生产者刷新并持久化所有消息的超时时间(以毫秒为单位)。默认值为 0,表示不设置超时。请确保在生产者上正确配置了 maxPendingMessages 和 blockIfQueueFull。 |
debezium.sink.pulsar.client.* | Pulsar 模块支持透传配置。客户端的配置属性在去除前缀后传递给客户端。至少必须提供 serviceUrl。 | |
debezium.sink.pulsar.producer.* | Pulsar 模块支持透传配置。消息生产者的配置属性在去除前缀后传递给生产者。topic 由 Debezium 设置。 | |
debezium.sink.pulsar.producer.batcherBuilder | DEFAULT | 指定生产者的批处理器构建器。生产者使用该批处理器构建器来创建批量消息容器。此设置仅在启用批量处理时有效。可选值为 DEFAULT 或 KEY_BASED,后者用于 KeyShared 订阅。 |
debezium.sink.pulsar.null.key | default | 没有主键的表会发送键为 null 的消息。Pulsar 不支持这种情况,因此必须使用替代键。 |
debezium.sink.pulsar.tenant | public | 用于投递消息的目标租户。 |
debezium.sink.pulsar.namespace | default | 用于投递消息的目标命名空间。 |
注入点
Pulsar 接收器的行为可以通过自定义逻辑进行修改,即为特定功能提供替代实现。当不存在替代实现时,则使用默认实现。
| 接口 | CDI 分类器 | 说明 |
|---|---|---|
io.debezium.server.StreamNameMapper | 自定义实现将计划中的目标(topic)名称映射为物理 Pulsar topic 名称。默认情况下使用相同的名称。 |
Azure Event Hubs
Azure Event Hubs 是一个大数据流平台和事件摄取服务,每秒可以接收并处理数百万个事件。发送到事件中心的数据可以通过任意实时分析提供程序或批处理/存储适配器进行转换和存储。
属性 默认值 说明
必须设置为 eventhubs。
debezium.sink.eventhubs.connectionstring
与 Event Hubs 通信所需的连接字符串。格式为:Endpoint=sb://<NAMESPACE>/;SharedAccessKeyName=<ACCESS_KEY_NAME>;SharedAccessKey=<ACCESS_KEY_VALUE>
debezium.sink.eventhubs.hubname
Event Hub 的名称
debezium.sink.eventhubs.dynamicpartitionrouting
一个可选设置,用于控制动态分区路由的行为。有三个可能的值:
default:如果记录键不为空,则将其用作批次分区键。否则,如果记录的分区 ID 不为空,则对事件进行批处理并路由到指定的分区 ID。否则,采用轮询路由。这是未提供任何值时的默认方式。key:如果记录键不为空,则将其用作批次分区键。否则,采用轮询路由。partitionid:如果记录的分区 ID 不为空,则对事件进行批处理并路由到指定的分区 ID。否则,采用轮询路由。
debezium.sink.eventhubs.dynamicpartitionrouting 配置仅在以下两个选项均未设置时才会生效:
debezium.sink.eventhubs.partitioniddebezium.sink.eventhubs.partitionkey
debezium.sink.eventhubs.partitionid
(可选)事件将发送到的 Event Hub 分区的标识符。如果您希望 Debezium 接收到的所有变更事件都发送到 Event Hubs 中的特定分区,请使用此配置。如果您已指定 debezium.sink.eventhubs.partitionkey,请勿使用此配置。
debezium.sink.eventhubs.partitionkey
(可选)用于对事件进行哈希计算的分区键。如果您希望 Debezium 接收到的所有变更事件都发送到 Event Hubs 中的特定分区,请使用此配置。如果您已指定 debezium.sink.eventhubs.partitionid,请勿使用此配置。
debezium.sink.eventhubs.maxbatchsize
设置事件批次的最大大小(以字节为单位)。
debezium.sink.eventhubs.hashmessagekeyfunction
无默认值
(可选)指定 Debezium 用于对 Azure Event Hubs 消息键进行加密的哈希函数。
请指定以下值之一:
javamd5sha1sha256
在 EventHubs 中使用分区
Azure Event Hubs 支持多种分区路由策略,用于控制 Debezium 如何在各分区之间分发变更事件。您可以指定固定的分区 ID、分区键,或使用 Partition Routing 转换器,以满足部署中的吞吐量和顺序要求。
默认情况下,当既未定义可选属性 debezium.sink.eventhubs.partitionid 也未定义 debezium.sink.eventhubs.partitionkey 时,EventHubs sink 将以轮询方式将事件发送到所有可用分区。
你可以通过设置 debezium.sink.eventhubs.partitionid 属性,强制将所有消息发送到单个固定的分区。或者,你也可以使用 debezium.sink.eventhubs.partitionkey 属性指定一个固定的分区键,EventHubs 将据此将所有事件路由到特定分区。
如果你有更复杂的路由需求,可以使用 分区路由 转换器。请确保转换器中 partition.topic.num 设置指定的分区数量小于或等于 EventHubs 命名空间中可用的分区数量,这样事件就不会被路由到不存在的分区 ID 上。例如,若要根据源架构名称将所有事件路由到 5 个分区,你可以在 application.properties 中进行如下设置:
# Uses a hash of `source.db` to calculate which partition to send the event to. Ensures all events from the same source schema are sent to the same partition.
debezium.transforms=PartitionRouter
debezium.transforms.PartitionRouter.type=io.debezium.transforms.partitions.PartitionRouting
debezium.transforms.PartitionRouter.partition.payload.fields=source.db
debezium.transforms.PartitionRouter.partition.topic.num=5注入点
默认的 sink 行为可以通过自定义逻辑修改,为特定功能提供替代实现。当替代实现不可用时,则使用默认实现。
| 接口 | CDI classifier | 描述 |
|---|---|---|
com.azure.messaging.eventhubs.EventHubProducerClient | @CustomConsumerBuilder | 用于发送消息的经过自定义配置的 EventHubProducerClient 实例。 |
Redis(Stream)
Redis sink 将变更事件存储为 Redis stream 中的条目,同时支持单实例和集群部署。通过将 debezium.sink.type 属性设置为 redis 并指定连接、认证和存储选项来配置该 sink。
Redis 是一个开源(BSD 许可证)的内存数据结构存储,用作数据库、缓存和消息代理。Stream 是一种以更抽象的方式对日志数据结构建模的数据类型。它实现了强大的操作,以克服日志文件的局限性。
Debezium 可以使用单个 Redis 实例,也可以使用 Redis 集群模式执行 sink 操作:
单实例模式
连接到单个 Redis 服务器实例。
集群模式
连接到 Redis 集群,以实现高可用性和水平扩展。
要启用 Redis 集群模式,请将 debezium.sink.redis.cluster.enabled 属性设置为 true,并在 debezium.sink.redis.address 属性中提供以逗号分隔的 host:port 地址。
| 属性 | 默认值 | 描述 |
|---|---|---|
debezium.sink.type | 必须设置为 redis。 | |
debezium.sink.redis.address | 目标 Redis 流所在的地址,格式为 host:port。 | |
debezium.sink.redis.db.index | 0 | 0—15 范围内的数字,Debezium 用它来选择要使用的数据库。此设置仅适用于独立 Redis 连接;Redis 集群只使用数据库 0。 |
debezium.sink.redis.user | (可选)用于与 Redis 通信的用户名。 | |
debezium.sink.redis.password | (可选)用于与 Redis 通信的(对应用户的)密码。如果设置了用户,则必须设置密码。 | |
debezium.sink.redis.ssl.enabled | false | (可选)一个布尔值,指定连接 Redis 是否需要 SSL。 |
debezium.sink.redis.ssl.hostname.verification.enabled | false | (可选)一个布尔值,指定连接 Redis 时是否验证服务器的主机名。 |
debezium.sink.redis.ssl.truststore.path | (可选)如果在启用 SSL 的情况下使用 Redis sink,则为信任库文件的路径。如果设置,Redis 连接将优先使用此属性,而不是其他配置或系统属性。 | |
debezium.sink.redis.ssl.truststore.password | (可选)如果在启用 SSL 的情况下使用 Redis sink,则为信任库文件的密码。如果设置,Redis 连接将优先使用此属性,而不是其他配置或系统属性。 | |
debezium.sink.redis.ssl.truststore.type | JKS | (可选)如果在启用 SSL 的情况下使用 Redis sink,则为信任库文件的类型。如果设置,Redis 连接将优先使用此属性,而不是其他配置或系统属性。 |
debezium.sink.redis.ssl.keystore.path | (可选)如果在启用 SSL 的情况下使用 Redis sink,则为密钥库文件的路径。如果设置,Redis 连接将优先使用此属性,而不是其他配置或系统属性。 | |
debezium.sink.redis.ssl.keystore.password | (可选)如果在启用 SSL 的情况下使用 Redis sink,则为密钥库文件的密码。如果设置,Redis 连接将优先使用此属性,而不是其他配置或系统属性。 | |
debezium.sink.redis.ssl.keystore.type | JKS | (可选)如果在启用 SSL 的情况下使用 Redis sink,则为密钥库文件的类型。如果设置,Redis 连接将优先使用此属性,而不是其他配置或系统属性。 |
debezium.sink.redis.null.key | default | Redis 不支持没有键的数据概念,因此对于没有主键的记录,将使用此字符串作为键。 |
debezium.sink.redis.null.value | default | Redis 不支持空负载的概念,这与墓碑事件的情况类似。因此对于没有负载的记录,将使用此字符串作为值。 |
debezium.sink.redis.batch.size | 500 | 单次批量写入(流水线事务)中插入的变更记录数量。 |
debezium.sink.redis.retry.initial.delay.ms | 300 | 遇到 Redis 连接或内存不足(OOM)问题时的初始重试延迟。每次重试该值都会翻倍,但不会超过 debezium.sink.redis.retry.max.delay.ms。 |
debezium.sink.redis.retry.max.delay.ms | 10000 | 遇到 Redis 连接或内存不足(OOM)问题时的最大延迟。 |
debezium.sink.redis.connection.timeout.ms | 2000 | Redis 客户端的连接超时时间。 |
debezium.sink.redis.socket.timeout.ms | 2000 | Redis 客户端的套接字超时时间。 |
debezium.sink.redis.wait.enabled | false | 启用等待副本功能。如果 Redis 配置了副本分片,此功能可用来验证数据已经写入副本。更多信息请参阅 Redis WAIT 命令。 |
debezium.sink.redis.wait.timeout.ms | 1000 | 等待副本时的超时时间(毫秒)。必须为正值。 |
debezium.sink.redis.wait.retry.enabled | false | 启用等待副本失败时的重试。 |
debezium.sink.redis.wait.retry.delay.ms | 1000 | 等待副本失败时的重试延迟。 |
debezium.sink.redis.message.format | compact | 发送到 Redis 流的消息格式。可选值为 extended(较新的格式)和 compact(迄今为止的旧格式)。有关消息格式的更多信息,请参阅下文。 |
| debezium.sink.redis.memory.threshold.percentage | 85 | 如果 used_memory 占 Redis 配置的 maxmemory 的百分比高于或等于此阈值,sink 将停止消费记录。如果配置
消息格式
Redis 接收器支持两种消息格式:compact(紧凑)和 extended(扩展),它们的区别在于 Redis 流条目中键和值的存储方式不同。
设置 redis.message.format 属性,以指定向 Redis 接收器发送消息时使用哪种格式。
extended格式,使用两对内容 {1), 2)}={"key", "message key"} 和 {3), 4)}={"value", "message value"}:
1) 1) "1639304527499-0"
2) 1) "key"
2) "{\"schema\": {\"type\": \"struct\", \"fields\": [{\"type\": \"int32\", \"optional\": false, \"field\": \"empno\"}], \"optional\": false, \"name\": \"redislabs.dbo.emp.Key\"}, \"payload\": {\"empno\": 11}}"
3) "value"
4) "{\"schema\": {\"type\": \"struct\", \"fields\": [{\"type\": \"struct\", \"fields\": [{\"type\": \"int32\", \"optional\": false, \"field\": \"empno\"}, {\"type\": \"string\", \"optional\": true, \"field\": \"fname\"}, {\"type\": \"string\", \"optional\": true, \"field\": \"lname\"}, {\"type\": \"string\", \"optional\": true, \"field\": \"job\"}, {\"type\": \"int32\", \"optional\": true, \"field\": \"mgr\"}, {\"type\": \"int64\", \"optional\": true, \"name\": \"io.debezium.time.Timestamp\", \"version\": 1, \"field\": \"hiredate\"}, {\"type\": \"bytes\", \"optional\": true, \"name\": \"org.apache.kafka.connect.data.Decimal\", \"version\": 1, \"parameters\": {\"scale\": \"4\", \"connect.decimal.precision\": \"19\"}, \"field\": \"sal\"}, {\"type\": \"bytes\", \"optional\": true, \"name\": \"org.apache.kafka.connect.data.Decimal\", \"version\": 1, \"parameters\": {\"scale\": \"4\", \"connect.decimal.precision\": \"19\"}, \"field\": \"comm\"}, {\"type\": \"int32\", \"optional\": true, \"field\": \"dept\"}], \"optional\": true, \"name\": \"redislabs.dbo.emp.Value\", \"field\": \"before\"}, {\"type\": \"struct\", \"fields\": [{\"type\": \"int32\", \"optional\": false, \"field\": \"empno\"}, {\"type\": \"string\", \"optional\": true, \"field\": \"fname\"}, {\"type\": \"string\", \"optional\": true, \"field\": \"lname\"}, {\"type\": \"string\", \"optional\": true, \"field\": \"job\"}, {\"type\": \"int32\", \"optional\": true, \"field\": \"mgr\"}, {\"type\": \"int64\", \"optional\": true, \"name\": \"io.debezium.time.Timestamp\", \"version\": 1, \"field\": \"hiredate\"}, {\"type\": \"bytes\", \"optional\": true, \"name\": \"org.apache.kafka.connect.data.Decimal\", \"version\": 1, \"parameters\": {\"scale\": \"4\", \"connect.decimal.precision\": \"19\"}, \"field\": \"sal\"}, {\"type\": \"bytes\", \"optional\": true, \"name\": \"org.apache.kafka.connect.data.Decimal\", \"version\": 1, \"parameters\": {\"scale\": \"4\", \"connect.decimal.precision\": \"19\"}, \"field\": \"comm\"}, {\"type\": \"int32\", \"optional\": true, \"field\": \"dept\"}], \"optional\": true, \"name\": \"redislabs.dbo.emp.Value\", \"field\": \"after\"}, {\"type\": \"struct\", \"fields\": [{\"type\": \"string\", \"optional\": false, \"field\": \"version\"}, {\"type\": \"string\", \"optional\": false, \"field\": \"connector\"}, {\"type\": \"string\", \"optional\": false, \"field\": \"name\"}, {\"type\": \"int64\", \"optional\": false, \"field\": \"ts_ms\"}, {\"type\": \"string\", \"optional\": true, \"name\": \"io.debezium.data.Enum\", \"version\": 1, \"parameters\": {\"allowed\": \"true,last,false\"}, \"default\": \"false\", \"field\": \"snapshot\"}, {\"type\": \"string\", \"optional\": false, \"field\": \"db\"}, {\"type\": \"string\", \"optional\": true, \"field\": \"sequence\"}, {\"type\": \"string\", \"optional\": false, \"field\": \"schema\"}, {\"type\": \"string\", \"optional\": false, \"field\": \"table\"}, {\"type\": \"string\", \"optional\": true, \"field\": \"change_lsn\"}, {\"type\": \"string\", \"optional\": true, \"field\": \"commit_lsn\"}, {\"type\": \"int64\", \"optional\": true, \"field\": \"event_serial_no\"}], \"optional\": false, \"name\": \"io.debezium.connector.sqlserver.Source\", \"field\": \"source\"}, {\"type\": \"string\", \"optional\": false, \"field\": \"op\"}, {\"type\": \"int64\", \"optional\": true, \"field\": \"ts_ms\"}, {\"type\": \"struct\", \"fields\": [{\"type\": \"string\", \"optional\": false, \"field\": \"id\"}, {\"type\": \"int64\", \"optional\": false, \"field\": \"total_order\"}, {\"type\": \"int64\", \"optional\": false, \"field\": \"data_collection_order\"}], \"optional\": true, \"field\": \"transaction\"}], \"optional\": false, \"name\": \"redislabs.dbo.emp.Envelope\"}, \"payload\": {\"before\": {\"empno\": 11, \"fname\": \"Yossi\", \"lname\": \"Mague\", \"job\": \"PFE\", \"mgr\": 1, \"hiredate\": 1562630400000, \"sal\": \"dzWUAA==\", \"comm\": \"AYag\", \"dept\": 3}, \"after\": null, \"source\": {\"version\": \"1.6.0.Final\", \"connector\": \"sqlserver\", \"name\": \"redislabs\", \"ts_ms\": 1637859764960, \"snapshot\": \"false\", \"db\": \"RedisConnect\", \"sequence\": null, \"schema\": \"dbo\", \"table\": \"emp\", \"change_lsn\": \"0000003a:00002f50:0002\", \"commit_lsn\": \"0000003a:00002f50:0005\", \"event_serial_no\": 1}, \"op\": \"d\", \"ts_ms\": 1637859769370, \"transaction\": null}}"- 以及
compact格式,仅使用一组 {1)、2)}={"message key"、"message value"}:
1) 1) "1639304527499-0"
2) 1) "{\"schema\": {\"type\": \"struct\", \"fields\": [{\"type\": \"int32\", \"optional\": false, \"field\": \"empno\"}], \"optional\": false, \"name\": \"redislabs.dbo.emp.Key\"}, \"payload\": {\"empno\": 11}}"
2) "{\"schema\": {\"type\": \"struct\", \"fields\": [{\"type\": \"struct\", \"fields\": [{\"type\": \"int32\", \"optional\": false, \"field\": \"empno\"}, {\"type\": \"string\", \"optional\": true, \"field\": \"fname\"}, {\"type\": \"string\", \"optional\": true, \"field\": \"lname\"}, {\"type\": \"string\", \"optional\": true, \"field\": \"job\"}, {\"type\": \"int32\", \"optional\": true, \"field\": \"mgr\"}, {\"type\": \"int64\", \"optional\": true, \"name\": \"io.debezium.time.Timestamp\", \"version\": 1, \"field\": \"hiredate\"}, {\"type\": \"bytes\", \"optional\": true, \"name\": \"org.apache.kafka.connect.data.Decimal\", \"version\": 1, \"parameters\": {\"scale\": \"4\", \"connect.decimal.precision\": \"19\"}, \"field\": \"sal\"}, {\"type\": \"bytes\", \"optional\": true, \"name\": \"org.apache.kafka.connect.data.Decimal\", \"version\": 1, \"parameters\": {\"scale\": \"4\", \"connect.decimal.precision\": \"19\"}, \"field\": \"comm\"}, {\"type\": \"int32\", \"optional\": true, \"field\": \"dept\"}], \"optional\": true, \"name\": \"redislabs.dbo.emp.Value\", \"field\": \"before\"}, {\"type\": \"struct\", \"fields\": [{\"type\": \"int32\", \"optional\": false, \"field\": \"empno\"}, {\"type\": \"string\", \"optional\": true, \"field\": \"fname\"}, {\"type\": \"string\", \"optional\": true, \"field\": \"lname\"}, {\"type\": \"string\", \"optional\": true, \"field\": \"job\"}, {\"type\": \"int32\", \"optional\": true, \"field\": \"mgr\"}, {\"type\": \"int64\", \"optional\": true, \"name\": \"io.debezium.time.Timestamp\", \"version\": 1, \"field\": \"hiredate\"}, {\"type\": \"bytes\", \"optional\": true, \"name\": \"org.apache.kafka.connect.data.Decimal\", \"version\": 1, \"parameters\": {\"scale\": \"4\", \"connect.decimal.precision\": \"19\"}, \"field\": \"sal\"}, {\"type\": \"bytes\", \"optional\": true, \"name\": \"org.apache.kafka.connect.data.Decimal\", \"version\": 1, \"parameters\": {\"scale\": \"4\", \"connect.decimal.precision\": \"19\"}, \"field\": \"comm\"}, {\"type\": \"int32\", \"optional\": true, \"field\": \"dept\"}], \"optional\": true, \"name\": \"redislabs.dbo.emp.Value\", \"field\": \"after\"}, {\"type\": \"struct\", \"fields\": [{\"type\": \"string\", \"optional\": false, \"field\": \"version\"}, {\"type\": \"string\", \"optional\": false, \"field\": \"connector\"}, {\"type\": \"string\", \"optional\": false, \"field\": \"name\"}, {\"type\": \"int64\", \"optional\": false, \"field\": \"ts_ms\"}, {\"type\": \"string\", \"optional\": true, \"name\": \"io.debezium.data.Enum\", \"version\": 1, \"parameters\": {\"allowed\": \"true,last,false\"}, \"default\": \"false\", \"field\": \"snapshot\"}, {\"type\": \"string\", \"optional\": false, \"field\": \"db\"}, {\"type\": \"string\", \"optional\": true, \"field\": \"sequence\"}, {\"type\": \"string\", \"optional\": false, \"field\": \"schema\"}, {\"type\": \"string\", \"optional\": false, \"field\": \"table\"}, {\"type\": \"string\", \"optional\": true, \"field\": \"change_lsn\"}, {\"type\": \"string\", \"optional\": true, \"field\": \"commit_lsn\"}, {\"type\": \"int64\", \"optional\": true, \"field\": \"event_serial_no\"}], \"optional\": false, \"name\": \"io.debezium.connector.sqlserver.Source\", \"field\": \"source\"}, {\"type\": \"string\", \"optional\": false, \"field\": \"op\"}, {\"type\": \"int64\", \"optional\": true, \"field\": \"ts_ms\"}, {\"type\": \"struct\", \"fields\": [{\"type\": \"string\", \"optional\": false, \"field\": \"id\"}, {\"type\": \"int64\", \"optional\": false, \"field\": \"total_order\"}, {\"type\": \"int64\", \"optional\": false, \"field\": \"data_collection_order\"}], \"optional\": true, \"field\": \"transaction\"}], \"optional\": false, \"name\": \"redislabs.dbo.emp.Envelope\"}, \"payload\": {\"before\": {\"empno\": 11, \"fname\": \"Yossi\", \"lname\": \"Mague\", \"job\": \"PFE\", \"mgr\": 1, \"hiredate\": 1562630400000, \"sal\": \"dzWUAA==\", \"comm\": \"AYag\", \"dept\": 3}, \"after\": null, \"source\": {\"version\": \"1.6.0.Final\", \"connector\": \"sqlserver\", \"name\": \"redislabs\", \"ts_ms\": 1637859764960, \"snapshot\": \"false\", \"db\": \"RedisConnect\", \"sequence\": null, \"schema\": \"dbo\", \"table\": \"emp\", \"change_lsn\": \"0000003a:00002f50:0002\", \"commit_lsn\": \"0000003a:00002f50:0005\", \"event_serial_no\": 1}, \"op\": \"d\", \"ts_ms\": 1637859769370, \"transaction\": null}}"关于 Redis Streams 的更多信息,可参阅此处。
注入点
Redis sink 的行为可以通过自定义逻辑进行修改,即为特定功能提供替代实现。当替代实现不可用时,将使用默认实现。
| 接口 | CDI 分类器 | 说明 |
|---|---|---|
io.debezium.server.StreamNameMapper | 自定义实现将计划中的目标(topic)名称映射为物理 Redis 流名称。默认情况下使用相同的名称。 |
Redis 配置示例
以下示例展示了在单实例和集群两种部署模式下,将 Debezium Server 连接到 Redis sink 的完整 application.properties 配置。
单实例模式
当 sink、offset 存储和 schema 历史记录均连接到单个 Redis 实例时,请使用以下配置。
# Sink configuration
debezium.sink.type=redis
debezium.sink.redis.address=localhost:6379
debezium.sink.redis.password=password
debezium.sink.redis.cluster.enabled=false
# Offset storage configuration
debezium.source.offset.storage=io.debezium.storage.redis.offset.RedisOffsetBackingStore
debezium.source.offset.storage.redis.address=localhost:6379
debezium.source.offset.storage.redis.password=password
debezium.source.offset.storage.redis.cluster.enabled=false
# Schema history storage configuration
debezium.source.schema.history.internal=io.debezium.storage.redis.history.RedisSchemaHistory
debezium.source.schema.history.internal.redis.address=localhost:6379
debezium.source.schema.history.internal.redis.password=password
debezium.source.schema.history.internal.redis.cluster.enabled=false集群模式
将 Debezium Server 连接到 Redis Cluster 时,请使用以下配置,为 sink、offset 存储和 schema history 指定多个节点地址。
# Sink configuration
debezium.sink.type=redis
debezium.sink.redis.address=redis-node-1:7001,redis-node-2:7002,redis-node-3:7003
debezium.sink.redis.password=password
debezium.sink.redis.cluster.enabled=true
# Offset storage configuration
debezium.source.offset.storage=io.debezium.storage.redis.offset.RedisOffsetBackingStore
debezium.source.offset.storage.redis.address=redis-node-1:7001,redis-node-2:7002,redis-node-3:7003
debezium.source.offset.storage.redis.password=password
debezium.source.offset.storage.redis.cluster.enabled=true
# Schema history storage configuration
debezium.source.schema.history.internal=io.debezium.storage.redis.history.RedisSchemaHistory
debezium.source.schema.history.internal.redis.address=redis-node-1:7001,redis-node-2:7002,redis-node-3:7003
debezium.source.schema.history.internal.redis.password=password
debezium.source.schema.history.internal.redis.cluster.enabled=true| 如果将 Debezium 配置为使用 Redis 集群模式,请确保 Redis 集群已正确配置且可以访问。Debezium Server 实例必须能够与集群节点通信。 |
|---|
NATS Streaming
NATS Streaming 是一个由 NATS 提供支持、使用 Go 编程语言编写的数据流系统。
| 属性 | 默认值 | 描述 |
|---|---|---|
debezium.sink.type | 必须设置为 nats-streaming。 | |
debezium.sink.nats-streaming.url | 指向集群中一个或多个节点的 URL(或以逗号分隔的 URL 列表),格式为 nats://host:port。 | |
debezium.sink.nats-streaming.cluster.id | NATS Streaming 集群 ID。 | |
debezium.sink.nats-streaming.client.id | NATS Streaming 客户端 ID。 |
扩展注入点
NATS Streaming 接收器的行为可以通过自定义逻辑进行修改,即为特定功能提供替代实现。当替代实现不可用时,则会使用默认实现。
| 接口 | CDI 分类器 | 描述 |
|---|---|---|
io.nats.streaming.StreamingConnection | @CustomConsumerBuilder | 用于向目标 subject 发布消息的、经过自定义配置的 StreamingConnection 实例。 |
io.debezium.server.StreamNameMapper | 自定义实现,用于将计划中的目标(主题)名称映射为物理 NATS Streaming subject 名称。默认情况下使用相同的名称。 |
NATS JetStream
NATS 内置了一个名为 JetStream 的分布式持久化系统,它在基础的 Core NATS 功能和服务质量之上,提供了新的功能和更高的服务质量。
属性 默认值 描述
必须设置为 nats-jetstream。
debezium.sink.nats-jetstream.stream-name
DebeziumStream
Debezium 所创建的流的可选自定义名称。
debezium.sink.nats-jetstream.url
指向集群中一个或多个节点的 URL(或以逗号分隔的 URL 列表),格式为 nats://host:port。
debezium.sink.nats-jetstream.create-stream
如果为 true,将创建一个基础流。
debezium.sink.nats-jetstream.subjects
*.*.*
以逗号分隔的 JetStream subject 列表,也即消息通道名称。可以指定包含通配符的条目,例如 test.inventory.*。
要同时捕获架构变更事件和数据变更事件,必须同时指定主题前缀和通配符模式。例如,如果 debezium.source.topic.prefix 为 myapp,则将主题(subject)配置为 myapp,myapp.> 或 myapp,myapp.. |
|---|
数据变更事件会发布到针对特定表的主题(subject)上(例如 myapp.database.table)。值 myapp.> 可匹配任意数量的主题层级;类似地,myapp.. 恰好匹配前缀之后的两个层级。
debezium.sink.nats-jetstream.storage
memory
控制消息在流(stream)中的保存方式,可以是 memory 或 file。
debezium.sink.nats-jetstream.auth.jwt
无默认值
指定 NATS 服务器客户端的身份。将此属性添加到配置中即可启用与 NATS 的 JSON Web Token(JWT)身份验证。要使用 NATS 的 JWT 身份验证,必须指定 NKey 种子。如果启用了密码身份验证,则不要启用 JWT 身份验证。
debezium.sink.nats-jetstream.auth.seed
无默认值
当为 NATS 启用 JWT 身份验证时,使用此属性指定代表 Debezium 用户的 NKey 种子。Debezium 使用指定的 NKey 种子派生出一个私钥,然后使用该私钥对 NATS 服务器在身份验证过程中发出的随机数质询(nonce challenge)进行加密签名。Debezium 将签名后的随机数连同 debezium.sink.nats-jetstream.auth.jwt 客户端对应的公钥一起返回给服务器。
debezium.sink.nats-jetstream.auth.user
无默认值
指定已授权 NATS 用户的用户名。当配置中存在此属性时,即启用与 NATS 的密码身份验证。要使用 NATS 的密码身份验证,请指定 debezium.sink.nats-jetstream.auth.password。如果启用了 JWT 身份验证,则不要启用密码身份验证。
debezium.sink.nats-jetstream.auth.password
无默认值
指定启用密码认证时所使用的密码。
debezium.sink.nats-jetstream.async.enabled
true
(可选)布尔值,指定 Debezium 是否可以异步流式传输到 NATS JetStream 服务器。
debezium.sink.nats-jetstream.async.timeout.ms
5000
(可选)指定 Debezium 在发送一批消息进行异步处理后,等待 NATS 服务器确认的最大时间(以毫秒为单位)。在异步处理期间,每条消息都会以 asyncTimeoutMs 指定的超时时间进行发布。
如果你需要一个更具可配置性的流,可以使用 nats cli 创建。有关流的更多信息,请参阅:https://docs.nats.io/nats-concepts/jetstream/streams
注入点
可以通过自定义逻辑为特定功能提供替代实现,从而修改 NATS JetStream 接收器的行为。当替代实现不可用时,将使用默认实现。
| 接口 | CDI 限定符 | 说明 |
|---|---|---|
io.nats.client.JetStream | @CustomConsumerBuilder | 用于向目标 subject 发布消息的经过自定义配置的 JetStream 实例。 |
io.debezium.server.StreamNameMapper | 自定义实现,将计划的目标(主题)名称映射为物理 NATS JetStream subject 名称。默认情况下使用相同的名称。 |
Apache Fluss
Apache Fluss 是为实时分析构建的流式存储,提供统一的流处理与批处理能力,兼具高性能、低延迟和强一致性保证。Debezium Server 可以将变更事件发布到 Apache Fluss,既支持用于插入/更新/删除操作的主键表,也支持仅追加的仅日志表。
| Apache Fluss 要求连接器启用 schema 支持,不支持无 schema 的变更事件。 |
|---|
属性 默认值 描述
必须设置为 fluss。
debezium.sink.fluss.bootstrap.servers
以逗号分隔的 Fluss 协调者服务器地址列表,格式为 host:port。
debezium.sink.fluss.default.database
从事件目标解析表路径时使用的默认 Fluss 数据库。
debezium.sink.fluss.primary.key.mode
auto
根据主键是否存在来决定写入模式的选择。有效值:
auto
根据表是否具有主键,在追加模式和 upsert 模式之间自动选择。
upsert
要求表必须具有主键,并执行插入/更新/删除操作。
append
无论是否存在主键,都以仅追加模式写入所有记录。
debezium.sink.fluss.table.auto.create
false
设置为 true 时,如果 Fluss 表不存在则自动创建,并根据 Debezium 事件 schema 推导表结构。
debezium.sink.fluss.retries.max
5
针对瞬时写入失败的最大重试次数。
debezium.sink.fluss.retries.interval.ms
1000
初始重试间隔,单位为毫秒。
debezium.sink.fluss.retries.max.interval.ms
60000
最大重试间隔,单位为毫秒。
debezium.sink.fluss.retries.backoff.multiplier
2.0
指数退避期间应用于重试间隔的倍数。
注入点
可以通过自定义逻辑为特定功能提供替代实现,从而修改 Fluss sink 的行为。当没有提供替代实现时,将使用默认实现。
| 接口 | CDI classifier | 说明 |
|---|---|---|
io.debezium.server.StreamNameMapper | 自定义实现将计划的目标(topic)名称映射为物理 Fluss 表名。默认情况下使用相同的名称。 |
Apache Kafka
配置 Apache Kafka sink,将 Debezium Server 的变更事件发布到 Kafka topic。使用 sink 属性指定 broker 连接详情和其他 Kafka 特有的设置。
| 属性 | 默认值 | 描述 |
|---|---|---|
debezium.sink.type | 必须设置为 kafka。 | |
debezium.sink.kafka.producer.* | Kafka 接收端适配器支持透传配置。这意味着所有 Kafka 生产者配置属性都会在去除前缀后传递给生产者。至少必须提供 bootstrap.servers、key.serializer 和 value.serializer 属性。topic 由 Debezium 设置。 | |
debezium.sink.kafka.wait.message.delivery.timeout.ms | 30000 | 服务器等待某个请求完成并返回记录元数据的最长时间,单位为毫秒。指定的超时值同时也决定了服务器等待 Kafka 响应请求的间隔。将该值设置为 0 可禁用超时。 |
注入点
通过提供受支持扩展点的实现,可以自定义 Kafka sink 的行为。例如,您可以实现自定义的 StreamNameMapper,以控制 Debezium Server 如何将源流映射到 Kafka 主题。如果您未提供自定义实现,Debezium Server 将使用默认实现。
| 接口 | CDI 分类器 | 说明 |
|---|---|---|
io.debezium.server.StreamNameMapper | 自定义实现会将原始目标(主题)名称映射为另一个 Kafka 主题。默认情况下使用相同的名称。 |
Pravega
Pravega 是一个面向事件流和数据流的云原生存储系统。该 sink 提供两种模式:非事务模式和事务模式。非事务模式会将 Debezium 批次中的每个事件单独写入 Pravega。事务模式则会将 Debezium 批次写入一个 Pravega 事务,该事务在批次完成时提交。
Pravega sink 要求目标 scope 和 stream 已经创建。
| 属性 | 默认值 | 描述 |
|---|---|---|
debezium.sink.type | 必须设置为 pravega。 | |
debezium.sink.pravega.controller.uri | tcp://localhost:9090 | Pravega 集群中 Controller 的连接字符串。 |
debezium.sink.pravega.scope | 用于查找目标流的 scope 名称。 | |
debezium.sink.pravega.transaction | false | 设置为 true 时,sink 会为每个 Debezium 批次使用 Pravega 事务。 |
注入点
Pravega sink 的行为可以通过自定义逻辑进行修改,为特定功能提供替代实现。当替代实现不可用时,将使用默认实现。
| 接口 | CDI 分类器 | 描述 |
|---|---|---|
io.debezium.server.StreamNameMapper | 自定义实现将计划的目标(流)名称映射为物理 Pravega 流名称。默认情况下使用相同的名称。 |
Infinispan
Infinispan 是一个开源的内存数据网格,提供了丰富的缓存类型以及缓存存储。凭借极快的数据访问速度,Infinispan 除了其他用途外,还可作为各种数据处理与分析工具的数据源。
Infinispan sink 要求目标缓存已在 Infinispan 集群中定义并创建。
| 属性 | 默认值 | 说明 |
|---|---|---|
debezium.sink.type | 必须设置为 infinispan。 | |
debezium.sink.infinispan.server.host | Infinispan 集群中某台服务器的主机名(也可以是以逗号分隔的服务器列表)。 | |
debezium.sink.infinispan.server.port | 11222 | Infinispan 服务器的端口。 |
debezium.sink.infinispan.cache | 用于存储记录的(已存在的)缓存名称。 | |
debezium.sink.infinispan.user | (可选)用于连接 Infinispan 集群的用户名。 | |
debezium.sink.infinispan.password | (可选)用于连接 Infinispan 集群的密码。 |
注入点
可以通过自定义逻辑为特定功能提供替代实现,从而修改 Infinispan sink 的行为。当替代实现不可用时,将使用默认实现。
| 接口 | CDI 分类器 | 描述 |
|---|---|---|
org.infinispan.client.hotrod.RemoteCache | @CustomConsumerBuilder | 用于连接 Infinispan 集群并发送事件的 Hot Rod 缓存自定义实例。 |
Apache RocketMQ
Apache RocketMQ 是一个分布式消息与流处理平台,具有低延迟、高性能、高可靠、万亿级容量和灵活的扩展能力。Debezium Server 支持将捕获的变更事件发布到已配置的 RocketMQ 中。
| 属性 | 默认值 | 描述 |
|---|---|---|
debezium.sink.type | 必须设置为 rocketmq。 | |
debezium.sink.rocketmq.producer.name.srv.addr | Apache RocketMQ 的 NameServer 地址。 | |
debezium.sink.rocketmq.producer.group | Apache RocketMQ 的生产者组。 | |
debezium.sink.rocketmq.producer.max.message.size | 4M,建议小于 4 MB。 | (可选)所发送消息体的最大字节数。 |
debezium.sink.rocketmq.producer.send.msg.timeout | 3000ms | (可选)发送消息的超时时长,即客户端本地同步调用的等待时间。请根据实际应用场景设置合适的值,以避免线程长时间阻塞。 |
debezium.sink.rocketmq.producer.acl.enabled | false | (可选)用于启用访问授权的配置。 |
debezium.sink.rocketmq.producer.access.key | (可选)用于连接 Apache RocketMQ 集群的 Access Key。 | |
debezium.sink.rocketmq.producer.secret.key | (可选)用于连接 Apache RocketMQ 集群的 Access Secret。 |
注入点
可以通过自定义逻辑为特定功能提供替代实现,从而修改 RocketMQ 接收器的行为。如果不存在替代实现,则使用默认实现。
| 接口 | CDI 分类器 | 描述 |
|---|---|---|
org.apache.rocketmq.client.producer.DefaultMQProducer | @CustomConsumerBuilder | RocketMQ 的自定义配置实例,用于向目标主题发布消息。 |
io.debezium.server.StreamNameMapper | 自定义实现,将计划中的目标(流)名称映射为 RocketMQ 主题名称。默认情况下使用相同的名称。 |
RabbitMQ Stream
RabbitMQ 是一个开源消息代理,支持多种消息传递协议,可以以分布式和联合配置进行部署,以满足高扩展性、高可用性的需求。RabbitMQ 支持消息队列和流。Debezium Server 支持将捕获的变更事件发布到已配置的 RabbitMQ Stream。
属性 默认值 说明
必须设置为 rabbitmq。
debezium.sink.rabbitmq.connection.host
localhost
RabbitMQ 服务器的主机。
debezium.sink.rabbitmq.connection.port
5672
RabbitMQ 服务器的端口。
debezium.sink.rabbitmq.connection.*
RabbitMQ 模块支持直通配置。连接配置属性会在去掉前缀后传递给 RabbitMQ 客户端。
debezium.sink.rabbitmq.ackTimeout
30000
定义发布消息后等待代理(broker)确认的最大时间,以毫秒为单位。
debezium.sink.rabbitmq.exchange
主题名称
(可选)发布消息时使用的交换机(exchange)名称。
debezium.sink.rabbitmq.routingKey
空字符串
(可选)发布消息时使用的静态路由键(routing key)。
debezium.sink.rabbitmq.autoCreateRoutingKey
false
(可选)若为 true,将自动创建不存在的路由键。
debezium.sink.rabbitmq.routingKeyDurable
true
(可选)若为 true,目标队列中的内容将在 RabbitMQ 服务器重启后仍然保留。
debezium.sink.rabbitmq.routingKeyFromTopicName
false
(可选)已弃用,请参阅 debezium.sink.rabbitmq.routingKey.source。
debezium.sink.rabbitmq.deliveryMode
2
(可选)消息在 RabbitMQ 服务器上的传输与存储方式
- 1 - 非持久化
- 2 - 持久化
debezium.sink.rabbitmq.null.value
default
RabbitMQ 不支持空负载(payload)的概念,而墓碑事件(tombstone event)正是这种情况。因此,该字符串将作为没有负载的记录的值。
debezium.sink.rabbitmq.routingKey.source
static
(可选)获取事件路由键的方式。
static(默认):路由键将取自debezium.sink.rabbitmq.routingKey。topic:路由键与交换机名称相同。key:路由键将取自记录键。
注入点
可以通过自定义逻辑为特定功能提供替代实现,从而修改 RabbitMQ 接收器(sink)的行为。当替代实现不可用时,将使用默认实现。
| 接口 | CDI classifier | 描述 |
|---|---|---|
io.debezium.server.StreamNameMapper | 自定义实现将预定的目标(流)名称映射为 RabbitMQ 交换机名称,并(在启用时)映射为路由键名称。默认情况下使用相同的名称。 |
RabbitMQ 原生 Stream
自 RabbitMQ 3.9 起,RabbitMQ 引入了 Streams,它采用了一种全新的极速协议,可与 AMQP 0.9.1 并行使用。Stream 非常适合大规模扇出、重放与时间回溯以及大型日志等场景,且吞吐量极高(每秒数百万条消息)。
Debezium Server 经过增强,支持借助 RabbitMQ Stream Java Client 将捕获的变更事件发布到原生 RabbitMQ Stream 中。
| 属性 | 默认值 | 说明 |
|---|---|---|
debezium.sink.type | 必须设置为 rabbitmqstream。 | |
debezium.sink.rabbitmqstream.connection.host | localhost | RabbitMQ 服务器的主机名。 |
debezium.sink.rabbitmqstream.connection.port | 5552 | RabbitMQ Stream 协议的端口。 |
debezium.sink.rabbitmqstream.connection.* | RabbitMQ 模块支持透传配置。连接相关的配置属性在去除前缀后会传递给 RabbitMQ 客户端。 | |
debezium.sink.rabbitmqstream.ackTimeout | 30000 | 定义发布消息后等待 broker 确认的最大时间(以毫秒为单位)。 |
debezium.sink.rabbitmqstream.null.value | default | RabbitMQ 不支持空负载的概念,墓碑事件即属于这种情况。因此,对于没有负载的记录,将使用该字符串作为其值。 |
Milvus
Milvus 是一个开源向量数据库,专为相似度搜索和高维数据检索而设计,例如来自机器学习模型的向量嵌入(如文本、图像和音频)。你可以使用 Milvus 处理从源数据库捕获的向量数据类型,也可以结合转换功能从消息字段计算向量,并将其作为嵌入使用。
Milvus sink 会摄取接收到的消息,并将每条消息的 after 部分写入(upsert)到集合中。集合名称不能包含点号,因此 sink 会将所有点号替换为下划线字符。当收到删除消息时,会从集合中移除匹配的记录。
| 属性 | 默认值 | 说明 |
|---|---|---|
debezium.sink.type | 指定 sink 的类型,必须设置为 milvus。 | |
debezium.sink.milvus.uri | http://localhost:19530 | (可选)用于访问 Milvus 数据库实例的 URL。 |
debezium.sink.milvus.database | default | (可选)包含目标集合的数据库名称。 |
注入点
你可以通过应用自定义逻辑,为特定函数提供替代实现,从而修改 Milvus sink 连接器的行为。如果替代实现不可用,连接器将使用默认实现。
| 接口 | CDI 分类器 | 描述 |
|---|---|---|
io.milvus.v2.client.MilvusClientV2.MilvusClientV2 | @CustomConsumerBuilder | 自定义 MilvusClientV2 客户端的一个实例,已配置为访问目标集合。 |
io.debezium.server.StreamNameMapper | 自定义实现,用于将计划中的目标主题名称映射到 Milvus 集合。默认情况下,名称中的点号会被替换为下划线。 |
Qdrant Sink
Qdrant 是一个开源向量数据库,针对向量相似度搜索进行了优化,并扩展了强大的过滤能力。它专为高负载应用设计,使你能够高效地存储、管理和搜索嵌入向量。你可以使用 Qdrant 处理直接从源数据库捕获的向量数据类型,也可以使用转换从消息字段计算出嵌入向量,然后将这些嵌入向量发送到数据库进行处理。
Qdrant sink 会摄取传入的消息,并将每条消息的 after 部分上插入(upsert)到集合中。当收到删除消息时,会从集合中移除匹配的记录。
该 sink 遵循以下规则:
- 每个 Debezium 集合或表都映射到一个 Qdrant 集合。
- 必须提供主键,并将其用作 Qdrant 的 point ID(仅支持
INT64和UUID)。 FloatVector和DoubleVector数据可以作为 Qdrant 向量的来源。- 非主键且非向量的字段会映射到 Qdrant 的 payload 中。
| 属性 | 默认值 | 说明 |
|---|---|---|
debezium.sink.type | 无默认值。 | 指定接收器(sink)的类型。必须显式将其设置为 qdrant。 |
debezium.sink.qdrant.host | localhost | (可选)用于访问 Qdrant 数据库实例的主机名。 |
debezium.sink.qdrant.port | 6333 | (可选)用于访问 Qdrant 数据库实例的端口。 |
debezium.sink.qdrant.api.key | (可选)向 Qdrant 数据库实例进行身份验证所需的 API 密钥。 | |
debezium.sink.qdrant.vector.field.names | 无默认值。 | (可选)以逗号分隔的 collection-name:field-name 对列表,用于显式定义每个集合要使用的向量字段。对于包含多个向量字段的源表和集合,此字段为必填项。 |
debezium.sink.qdrant.field.include.list.<collection-name> | 无默认值。 | (可选)以逗号分隔的列表,用于指定集合中代表 Qdrant 集合载荷(payload)的字段名子集。 |
注入点
要修改连接器行为,你可以应用自定义逻辑,为某些函数指定替代实现。如果你指定的实现不可用,连接器会使用默认实现。
| 接口 | CDI 分类器 | 说明 |
|---|---|---|
io.qdrant.client.QdrantClient | @CustomConsumerBuilder | 自定义 QdrantClient 客户端的实例,已配置为访问目标集合。 |
io.debezium.server.StreamNameMapper | 无默认值。 | 自定义实现,用于将计划中的目标主题名称映射到 Qdrant 集合。 |
InstructLab
InstructLab 是一个由社区驱动的项目,用于扩充大语言模型(LLM),以便将其应用于生成式人工智能应用。在使用 InstructLab 时,发现基础模型能力存在差距的用户可以协作开发分类体系来扩充该模型,每位贡献者都提供其特定的专业知识和技能。
Debezium Server 支持你通过配置基于 InstructLab 问答(qna.yml)文件的数据接收器,实现将技能和知识添加到分类体系过程的自动化。该接收器配置定义了一系列映射,用于从事件流中派生问题、答案和上下文的取值。这些映射可以直接从事件载荷中的字段、标头或静态配置的常量获取取值。之后,用户可以定期使用 InstructLab 训练模型,以利用添加到分类体系中的新技能和知识。
| 属性 | 默认值 | 描述 |
|---|---|---|
debezium.sink.type | 必须设置为 instructlab。 | |
debezium.sink.instructlab.taxonomy.base.path | 存储 InstructLab 分类技能与知识的根目录的绝对路径。该值与分类领域属性配合使用,用于构造 qna.yml 文件的完整路径。 | |
debezium.sink.instructlab.taxonomies | 以逗号分隔的分类映射符号名称列表。 | |
debezium.sink.instructlab.taxonomy.<name>.topic | .* | 用于匹配主题的正则表达式,以确定是否应用 <name> 分类映射。 |
debezium.sink.instructlab.taxonomy.<name>.question | 指定一个映射定义,用作 qna.yml 文件中种子示例的 question 属性。此项为必填。 | |
debezium.sink.instructlab.taxonomy.<name>.answer | 指定一个映射定义,用作 qna.yml 文件中种子示例的 answer 属性。此项为必填。 | |
debezium.sink.instructlab.taxonomy.<name>.context | 指定一个映射定义,用作 qna.yml 文件中种子示例的 context 属性。此项为可选。 | |
debezium.sink.instructlab.taxonomy.<name>.domain | 指定分类领域,即从分类基础路径到 qna.yml 之间以 / 分隔的一系列目录(不含分类基础路径)。例如,值为 a/b、基础路径为 /taxonomy 时,表示 /taxonomy/a/b/qna.yml。 |
映射定义
在 InstructLab sink 配置中,你可以设置属性,以指定 Debezium Server 如何将事件消息中的字段映射到 InstructLab qna.yml 文件中的 question、answer 和 context 属性。
对于要在 qna.yml 文件中填充的每种属性类型,你需要指定一个映射前缀,以确定 Debezium 从哪个消息字段中提取值。可以指定以下前缀值:
值
如果映射定义以字符串 value: 作为前缀,Debezium 将从传入事件的负载中提取指定字段的值。例如,要使用事件负载中 abc 字段的值填充 InstructLab 的 question 属性,请将 debezium.sink.instructlab.taxonomy.<name>.question 属性设置为 value:abc。Debezium 处理该消息时,会获取负载字段 abc 的值,并将其作为 question 属性添加到 qna.yml 文件中。
对于具有 Debezium 结构化负载的事件,Debezium 会从负载的 after 部分提取指定字段。如果事件是扁平化的,则直接从事件的值中获取该字段。
Header
如果映射定义以字符串 header: 作为前缀,Debezium 将提取传入事件消息中指定 header 字段的值。例如,如果指定映射 header:h1,当 Debezium 在源消息中检测到名为 h1 的 header 时,它会提取 h1 字段的值。
常量
如果映射定义中不包含 header: 或 value: 前缀,当 Debezium 在传入消息中检测到指定的值时,会将其视为常量,并按原样使用。当你希望将某个特定静态常量的值映射到 qna.yml 文件中的属性时,请使用此选项。
JDBC
JDBC sink 使用 JDBC 将变更事件直接写入关系型数据库。它在底层借助了 Debezium JDBC 连接器,支持多种数据库方言,包括 Db2、MySQL、Oracle、PostgreSQL 和 SQL Server。
该 sink 通过使用 upsert 语义支持幂等写入,通过基本的模式演进自动创建或修改目标表,并支持删除传播。
| 属性 | 默认值 | 描述 |
|---|---|---|
debezium.sink.type | 必须设置为 jdbc。 | |
debezium.sink.jdbc.* | 所有 JDBC 连接器配置属性都可以使用 debezium.sink.jdbc. 前缀进行设置。有关可用属性的完整列表(包括连接、运行时以及特定于方言的选项),请参阅 Debezium JDBC 连接器配置属性。 |
JDBC 接收器要求目标数据库对应的 JDBC 驱动程序位于 classpath 中。Debezium Server 的发行包中不包含任何面向目标数据库的 JDBC 驱动程序,因此你必须手动将目标数据库的驱动程序添加到 lib/ 目录中。 |
|---|
| 尽管 JDBC sink 复用了 Debezium JDBC sink 连接器,但 Debezium Server 目前并不提供与在 Kafka Connect 等其他运行时中运行该连接器时相同的交付保证。具体而言,依赖运行时来实现偏移量管理、精确一次(exactly-once)语义或自动错误处理与重试的功能,其行为可能有所不同,甚至可能不可用。 |
|---|
JDBC sink 配置示例
debezium.sink.type=jdbc
debezium.sink.jdbc.connection.url=jdbc:mysql://localhost:3306/target_db
debezium.sink.jdbc.connection.username=root
debezium.sink.jdbc.connection.password=password
debezium.sink.jdbc.insert.mode=upsert
debezium.sink.jdbc.primary.key.mode=record_key
debezium.sink.jdbc.schema.evolution=basic
debezium.sink.jdbc.delete.enabled=trueApache Iceberg
Debezium Apache Iceberg 连接器是运行在 Debezium Server 中的接收端(sink)。它从源 Debezium 连接器消费数据库变更事件,并将其写入 Apache Iceberg 表。这样便可将数据库(变更数据捕获事件)直接复制到云存储或 HDFS 上的 Iceberg 表中,从而无需 ETL 工具、Spark 或其他流处理平台等中间系统。
Debezium Iceberg 连接器支持 upsert 和 append 两种数据复制模式。当启用事件与键的模式(schema)信息时,目标 Iceberg 表会在首次启动时自动创建。
该连接器提供以下功能:
- Upsert 与 Append 模式。
- 批量大小优化。
- 可自定义的表命名。
- 将位点(Offset)与 Schema 历史存储在 Iceberg 中。
- 自动模式演进。
- 建议使用
ExtractNewRecordState转换(事件扁平化)。upsert模式和主键自动检测等关键特性依赖于该 SMT 生成的扁平化事件结构。没有它,连接器只能以功能更受限的仅追加(append-only)模式运行。 - 面向 Apache Iceberg 的 Debezium Server 接收端是一个由社区维护的开源项目。请注意,它在单独的仓库中维护,地址如下:Debezium Server Iceberg GitHub
| 属性 | 默认值 | 说明 |
|---|---|---|
debezium.sink.type | 必须设置为 iceberg。 | |
debezium.sink.iceberg.catalog-name | default | 用户自定义的 Iceberg catalog 名称。 |
debezium.sink.iceberg.warehouse | 必填。Iceberg 数据仓库的根路径(例如 s3a://my-bucket/warehouse)。 | |
debezium.sink.iceberg.io-impl | org.apache.iceberg.io.ResolvingFileIO | 用于写入数据文件的 Iceberg FileIO 实现。 |
debezium.sink.iceberg.table-namespace | default | 在 catalog 中用于存放表的命名空间。 |
debezium.sink.iceberg.table-prefix | `` | 可选前缀,会添加到所有创建的 Iceberg 表名之前。 |
debezium.sink.iceberg.destination-regexp | `` | 用于重写目标表名的正则表达式。 |
debezium.sink.iceberg.destination-regexp-replace | `` | destination-regexp 表达式的替换字符串。 |
debezium.sink.iceberg.destination-uppercase-table-names | false | 若为 true,则创建大写的 Iceberg 表名。 |
debezium.sink.iceberg.destination-lowercase-table-names | false | 若为 true,则创建小写的 Iceberg 表名。 |
debezium.sink.iceberg.table-mapper | default-mapper | 源表名到目标表名的映射策略。目前仅提供 default-mapper。 |
debezium.sink.iceberg.upsert | false | 若为 true,则启用 upsert 模式;若为 false,则使用仅追加(append-only)模式。 |
debezium.sink.iceberg.upsert-keep-deletes | true | 在 upsert 模式下,若为 true,删除的行会被保留并标记为 __deleted=true(软删除);若为 false,则被物理删除。 |
debezium.sink.iceberg.upsert-dedup-column | __source_ts_ns | 在 upsert 模式下用于对记录去重的列。保留该列取值最高的记录。 |
debezium.sink.iceberg.upsert-op-field | __op | 包含操作类型(c、u、d、r)的字段名,用于去重逻辑。 |
debezium.sink.iceberg.write.format.default | parquet | Iceberg 表的默认文件格式,可以是 parquet、avro 或 orc。 |
debezium.sink.iceberg.nested-as-variant | false | 若为 true,所有嵌套数据都会存储在 Iceberg 的 variant 字段中,从而在无需进行 schema 演进的情况下吸收 schema 变更。 |
debezium.sink.iceberg.allow-field-addition | true | 若为 true,则当源事件中出现新字段时,允许连接器自动向 Iceberg 表中添加新列。 |
debezium.sink.iceberg.create-identifier-fields | true | 若为 false,连接器将不会为 Iceberg 表中的主键创建标识字段。在以仅追加方式消费嵌套事件时,必须将其设置为 false。 |
debezium.sink.iceberg.preserve-required-property | false | 若为 true,Iceberg 表中的列将保留源端原有的 required/optional 属性。默认情况下,仅主键列被标记为必需。 |
debezium.sink.batch.batch-size-wait | NoBatchSizeWait | 用于控制提交频率的批大小等待策略。 |
debezium.sink.batch.concurrent-uploads | 1 | 上传数据到 Iceberg 时使用的并行线程数。 |
debezium.sink.batch.concurrent-uploads.timeout-minutes | 60 | 等待所有并行上传完成的超时时间(分钟)。 |
debezium.source.offset.storage | io.debezium.server.iceberg.offset.IcebergOffsetBackingStore | 将偏移量存储设置为使用 Iceberg 表。 |
debezium.source.offset.storage.iceberg.table-name | _debezium_offset_storage | 用于存储连接器偏移量的 Iceberg 表名。 |
debezium.source.schema.history.internal | io.debezium.server.iceberg.history.IcebergSchemaHistory | 将 schema 历史记录设置为使用 Iceberg 表。 |
| debezium.source.schema.history.internal.iceberg.table-name | _debezium_database_history_storage | 用于存储数据库 schema 历史记录的 Iceberg �
扩展
Debezium Server 使用 Quarkus 框架,并依赖依赖注入机制,使开发者能够扩展其行为。请注意,仅支持 Quarkus 的 JVM 模式,不支持通过 GraalVM 进行本地执行。服务器可以通过两种方式提供自定义逻辑来进行扩展:
- 实现新的 sink
- 定制现有 sink,即非标准配置
实现新的 sink
新的 sink 可以实现为一个 CDI Bean,该 Bean 实现 DebeziumEngine.ChangeConsumer 接口,并带有 @Named 注解(指定唯一名称)和 @Dependent 作用域。该 Bean 的名称用作 debezium.sink.type 选项的值。
sink 需要使用 Microprofile Config API 读取配置。执行路径必须将消息传递到目标系统,并定期提交已传递/已处理的消息。
更多细节请参阅 Kinesis sink 的实现。
定制现有 sink
某些 sink 暴露了依赖注入点,允许用户提供自己的 Bean 来修改 sink 的行为。典型的例子包括对目标客户端设置、目标命名等进行微调。
更多细节请参阅自定义 主题命名策略 的实现示例。
评论
登录后参与评论
KnowForge