向 Debezium 发送信号
向 Debezium 连接器发送信号
概述
Debezium 信号机制提供了一种修改连接器行为或触发一次性操作的方式,例如发起对某张表的临时快照。要使用信号触发连接器执行指定操作,你可以配置连接器使用以下一个或多个通道:
SourceSignalChannel
你可以发出 SQL 命令,将信号消息添加到专用的信号数据集合中。该信号数据集合由你在源数据库上创建,专门用于与 Debezium 通信。每个连接器实例必须拥有各自唯一的信号数据集合。
KafkaSignalChannel
你可以将信号消息提交到可配置的 Kafka 主题。
JmxSignalChannel
你可以通过 JMX 的 signal 操作提交信号。
FileSignalChannel
你可以使用文件发送信号。
Custom
你可以将信号提交到自行实现的自定义通道。当 Debezium 检测到通道中有新的日志记录或临时快照记录被添加时,它会读取该信号并启动所请求的操作。
以下 Debezium 连接器支持信号功能:
- CockroachDB
- Db2
- MariaDB
- MongoDB
- MySQL
- Oracle
- PostgreSQL
- SQL Server
你可以通过设置 signal.enabled.channels 配置属性来指定启用哪些通道。该属性列出已启用通道的名称。默认情况下,Debezium 提供以下通道:source 和 kafka。source 通道默认启用,因为增量快照信号需要使用它。
错误处理
除 source 通道外,Debezium 信号通道均不实现重试策略。发起信号后,请务必验证信号是否成功完成。
你可以通过配置连接器发送通知,使连接器自动报告增量快照或阻塞快照的进度。
启用 source 信号通道
默认情况下,Debezium 的 source 信号通道已启用。
对于每个要使用信号功能的连接器,你必须显式配置信号功能。
操作步骤
在源数据库中创建一个信号数据集合表,用于向连接器发送信号。有关信号数据集合所需结构的信息,请参阅信号数据集合的结构。
对于实现了原生变更数据捕获(CDC)机制的源数据库(如 Db2 或 SQL Server),请为信号表启用 CDC。
将信号数据集合的名称添加到 Debezium 连接器配置中:
在连接器配置中,添加属性
signal.data.collection,并将其值设置为你在步骤 1 中创建的信号数据集合的完全限定名称。该数据集合的名称区分大小写。例如,
signal.data.collection = inventory.debezium_signals
有关设置 signal.data.collection 属性的更多信息,请参阅相应连接器的配置属性表。
如果使用 table.include.list 属性,则无需在其中包含信号数据集合。但是,如果使用 column.include.list 属性,则必须包含信号表的各列(id、type 和 data)。 |
|---|
其他资源
指定数据集合完全限定名的格式
| CockroachDB | <databaseName>.<schemaName>.<tableName> |
|---|---|
| Db2 | <schemaName>.<tableName> |
| MariaDB | <databaseName>.<tableName> |
| MongoDB | <databaseName>.<collectionName> |
| MySQL | <databaseName>.<tableName> |
| Oracle | <databaseName>.<schemaName>.<tableName> |
| PostgreSQL | <schemaName>.<tableName> |
| SQL Server | <databaseName>.<schemaName>.<tableName> |
信号数据集合的结构
信号数据集合(也称信号表)用于存储你向连接器发送的信号,以触发指定的操作。信号表的结构必须符合以下标准格式。
- 包含三个字段(列)。
- 字段按特定顺序排列,如表 1 所示。
| 字段 | 类型 | 描述 |
|---|---|---|
id(必填) | string | 用于标识信号实例的任意唯一字符串。你为提交到信号表的每个信号分配一个 id。通常,该 ID 是一个 UUID 字符串。你可以将信号实例用于日志记录、调试或去重。当信号触发 Debezium 执行增量快照时,Debezium 会生成一条带有任意 id 字符串的信号消息。生成的消息所包含的 id 字符串与所提交信号中的 id 字符串无关。 |
type(必填) | string | 指定要发送的信号类型。某些信号类型可用于任何支持信号传递的连接器,而另一些信号类型则仅适用于特定的连接器。 |
data(可选) | string | 指定要传递给信号动作的 JSON 格式参数。每种信号类型都需要特定的数据集合。 |
表 1. 信号数据集合的必填结构
信号数据集合必须包含名为 id、type 和 data 的列。列名中不要包含引号。如果你为这些列指定了其他名称,连接器将无法处理信号。 |
|---|
创建信号数据集合
你可以通过向源数据库提交标准 SQL DDL 查询来创建信号表。
前置条件
- 你拥有在源数据库上创建表的足够访问权限。
操作步骤
向源数据库提交 SQL 查询,创建符合所需结构的表,如下例所示:
CREATE TABLE <tableName> (id VARCHAR(<varcharValue>) PRIMARY KEY, type VARCHAR(<varcharValue>) NOT NULL, data VARCHAR(<varcharValue>) NULL);
你为 id 变量的 VARCHAR 参数分配的空间必须足以容纳发送到信号表的信号 ID 字符串的长度。如果 ID 的长度超过可用空间,连接器将无法处理该信号。 |
|---|
以下示例展示了一条 CREATE TABLE 命令,用于创建包含三列的 debezium_signal 表:
CREATE TABLE debezium_signal (id VARCHAR(42) PRIMARY KEY, type VARCHAR(32) NOT NULL, data VARCHAR(2048) NULL);Debezium 数据库用户必须拥有 debezium_signal 表的 INSERT 权限。 |
|---|
启用 Kafka 信号通道
你可以通过将 Kafka 信号通道添加到 signal.enabled.channels 配置属性中来启用它,然后将接收信号的主题名称添加到 signal.kafka.topic 属性中。启用信号通道后,会创建一个 Kafka 消费者,用于消费发送到所配置信号主题的信号。
消费者可用的其他配置
- Db2 连接器 Kafka 信号配置属性
- Informix 连接器 Kafka 信号配置属性
- MariaDB 连接器 Kafka 信号配置属性
- MongoDB 连接器 Kafka 信号配置属性
- MySQL 连接器 Kafka 信号配置属性
- Oracle 连接器 Kafka 信号配置属性
- PostgreSQL 连接器 Kafka 信号配置属性
- SQL Server 连接器 Kafka 信号配置属性
消息格式
Kafka 消息的键必须与 topic.prefix 连接器配置选项的值相匹配。
值是一个包含 type 和 data 字段的 JSON 对象。
当信号类型设置为 execute-snapshot 时,data 字段必须包含下表所列的字段:
表 2. 执行快照的数据字段
字段 默认值
type
incremental
要运行的快照类型。目前 Debezium 支持 incremental 和 blocking 类型。
data-collections
N/A
一个正则表达式数组(以逗号分隔),用于匹配要包含在快照中的数据集合的完全限定名称。命名格式 取决于数据库。
additional-conditions
N/A
一个可选数组,用于指定一组附加条件,连接器会评估这些条件,以确定要包含在快照中的记录子集。每个附加条件都是一个对象,用于指定过滤临时快照所捕获数据的条件。你可以为每个附加条件设置以下属性:
data-collection
过滤器所应用的数据集合的完全限定名称。你可以为每个数据集合应用不同的过滤器。
filter
指定数据库记录中必须存在、快照才会包含该记录的列值,例如 "color='blue'"。快照进程会针对 filter 值评估数据集合中的记录,只捕获包含匹配值的记录。
你为 filter 属性指定的具体值取决于临时快照的类型:
- 对于增量快照,你需要指定一个搜索条件片段,例如
"color='blue'",快照会将其追加到查询的条件子句中。 - 对于阻塞快照,你需要指定完整的
SELECT语句,例如你可能在snapshot.select.statement.overrides属性中设置的语句。
以下示例展示了一条典型的 execute-snapshot Kafka 消息:
Key = `test_connector`
Value = `{"type":"execute-snapshot","data": {"data-collections": ["schema1.table1", "schema1.table2"], "type": "INCREMENTAL"}}`启用 JMX 信号通道
要启用 JMX 信号,只需在连接器配置的 signal.enabled.channels 属性中添加 jmx,然后启用 JMX MBean Server 以暴露信号 Bean。
发送 JMX 信号
操作步骤
使用你喜欢的 JMX 客户端(例如 JConsole 或 JDK Mission Control)连接到 MBean 服务器。
查找 MBean
debezium.<connector-type>.management.signals.<server>。该 MBean 暴露了signal操作,接受以下输入参数:p0
信号的 ID。
p1
信号的类型,例如
execute-snapshot。p2
一个 JSON 数据字段,其中包含与指定信号类型相关的附加信息。
通过为输入参数提供值来发送
execute-snapshot信号。在 JSON 数据字段中,包含下表所列的信息:表 3. 执行快照的数据字段
字段 默认值
typeincremental要运行的快照类型。目前 Debezium 支持
incremental和blocking两种类型。data-collections无
一个由逗号分隔的正则表达式数组,用于匹配要包含在快照中的表的全限定名。
additional-conditions无
一个可选数组,用于指定连接器求值的一组附加条件,以确定要包含在快照中的记录子集。每个附加条件都是一个对象,用于指定即席快照所捕获数据的过滤标准。每个附加条件可以设置以下属性:
data-collection过滤器所应用的数据集合的全限定名。你可以为每个数据集合应用不同的过滤器。
filter指定数据库记录中必须存在的列值,快照才会包含该记录,例如
"color='blue'"。快照进程会使用filter的值对数据集合中的记录进行求值,仅捕获包含匹配值的记录。你为
filter属性指定的具体值取决于即席快照的类型:- 对于增量快照,你需要指定一个搜索条件片段,例如
"color='blue'",快照会将其追加到查询的条件子句中。 - 对于阻塞快照,你需要指定完整的
SELECT语句,例如你可能会设置在snapshot.select.statement.overrides属性中的语句。
- 对于增量快照,你需要指定一个搜索条件片段,例如
下图展示了使用 JConsole 发送信号的示例:

启用文件信号通道
你可以通过在连接器配置的 signal.enabled.channels 属性中添加 file 来启用 File 信号通道。启用信号通道后,你必须配置连接器以从文件中读取信号。默认情况下,信号文件创建在连接器 classpath 的根目录下,文件名为 file-signals.txt。如果你想使用其他文件,请在连接器配置中设置 signal.file 属性,并指定文件名和路径。文件路径必须在连接器运行环境中可访问。
消息格式
信号文件中的信号以 JSON 对象的形式表示,由 id、type 和 data 字段组成。
id 字段是信号的唯一标识符,通常是一个 UUID 字符串。
当信号类型设置为 execute-snapshot 时,data 字段必须包含下表中列出的字段:
表 4. 执行快照的数据字段
| 字段 | 默认值 | 说明 |
|---|---|---|
type | incremental | 要执行的快照类型。目前 Debezium 支持 incremental 和 blocking 两种类型。 |
data-collections | N/A | 一个正则表达式数组(以逗号分隔),用于匹配要包含在快照中的数据集合的完全限定名称。命名格式取决于数据库。 |
additional-conditions | N/A | 一个可选的数组,用于指定连接器评估的一组附加条件,以确定要包含在快照中的记录子集。每个附加条件是一个对象,用于指定临时快照所捕获数据的过滤标准。你可以为每个附加条件设置以下属性:data-collection:过滤器应用到的数据集合的完全限定名称。你可以为每个数据集合应用不同的过滤器。filter:指定快照要包含数据库记录所必须满足的列值,例如 "color='blue'"。快照进程会针对 filter 值评估数据集合中的记录,只捕获包含匹配值的记录。你为 filter 属性赋的具体值取决于临时快照的类型:- 对于增量快照,你需要指定一个搜索条件片段,例如 "color='blue'",快照会将其追加到查询的条件子句中。- 对于阻塞快照,你需要指定完整的 SELECT 语句,例如你可能在 snapshot.select.statement.overrides 属性中设置的语句。 |
以下示例展示了文件中典型的 execute-snapshot 消息:
{"id":"d139b9b7-7777-4547-917d-111111111111", "type":"execute-snapshot", "data":{"data-collections": ["public.MyFirstTable", "public.MySecondTable"]}}自定义信号通道
信号机制被设计为可扩展的。你可以根据需要实现通道,以最适合你环境的方式向 Debezium 发送信号。
添加信号通道包含以下几个步骤:
提供自定义信号通道
自定义信号通道是实现了 io.debezium.pipeline.signal.channels.SignalChannelReader 服务提供者接口(SPI)的 Java 类。例如:
public interface SignalChannelReader {
String name();
void init(CommonConnectorConfig connectorConfig);
List<SignalRecord> read();
void close();
}以下列表说明了上例中各行的用途:
String name();
读取器的名称。要使 Debezium 能够使用该通道,请在连接器的 signal.enabled.channels 属性中指定此名称。
void init(CommonConnectorConfig connectorConfig);
初始化该通道所需的具体配置、变量或连接。
List<SignalRecord> read();
从通道中读取信号。SignalProcessor 类会调用此方法以获取待处理的信号。
void close();
关闭所有已分配的资源。当连接器停止时,Debezium 会调用此方法。
Debezium 核心模块依赖
自定义信号通道的 Java 项目对 Debezium 核心模块具有编译依赖。你必须在项目的 pom.xml 文件中包含这些编译依赖,如以下示例所示:
<dependency>
<groupId>io.debezium</groupId>
<artifactId>debezium-connector-common</artifactId>
<version>${version.debezium}</version>
</dependency>在前面的 pom.xml 片段中,${version.debezium} 表示 Debezium 连接器的版本。
请在 META-INF/services/io.debezium.pipeline.signal.channels.SignalChannelReader 文件中声明你的实现。
部署自定义信号通道
前提条件
- 你已有一个自定义信号通道的 Java 程序。
操作步骤
要在 Debezium 连接器中使用自定义信号通道,请将 Java 项目导出为 JAR 文件,并将该文件复制到包含你要与之配合使用的每个 Debezium 连接器 JAR 文件的目录中。
例如,在典型部署中,Debezium 连接器文件存储在 Kafka Connect 目录(
/kafka/connect)的子目录中,每个连接器的 JAR 都位于各自的子目录中(/kafka/connect/debezium-connector-db2、/kafka/connect/debezium-connector-db2等)。
| 要在多个连接器中使用自定义信号通道,你必须在每个连接器的子目录中都放置一份自定义信号通道 JAR 文件的副本。 |
|---|
配置连接器以使用自定义信号通道
将自定义信号通道的名称添加到 signal.enabled.channels 配置属性中。
信号操作
你可以使用信号来触发以下操作:
要为大多数连接器使用信号来触发即席增量快照,必须先在连接器配置中启用 source 信号通道。source 通道实现了一种水印机制,用于对增量快照可能捕获、并在流式传输恢复后又被再次捕获的事件进行去重。当使用信号通道对启用了 GTIDs 的只读 MySQL 数据库触发增量快照时,无需启用 source 通道。有关更多信息,请参阅 MySQL 只读增量快照 |
|---|
有些信号并非与所有连接器都兼容。
记录信号
你可以通过创建 log 类型的信号表条目,请求连接器向日志中添加一条记录。连接器在处理该信号后,会将指定的消息打印到日志中。你还可以选择配置该信号,使生成的消息包含流式传输坐标。
表 5. 用于添加日志消息的信号记录示例
列值说明
id
924e3ff8-2245-43ca-ba77-2af9af02fa07
type
log
信号的动作类型。
data
{"message": "Signal message at offset {}"}message 参数指定要打印到日志的字符串。如果你在消息中添加占位符({}),它将被流式处理坐标所替换。
你可以通过创建信号类型为 execute-snapshot 的信号来请求连接器发起临时快照。连接器处理该信号后,会运行所请求的快照操作。
与连接器首次启动后运行的初始快照不同,临时快照发生在运行期间,即连接器已经开始从数据库流式传输变更事件之后。你可以在任何时候发起临时快照。
以下 Debezium 连接器支持临时快照:
- CockroachDB
- Db2
- MariaDB
- MongoDB
- MySQL
- Oracle
- PostgreSQL
- SQL Server
有关临时快照的更多信息,请参阅相应连接器文档中的 快照 主题。
其他资源
- CockroachDB 连接器增量快照
- Db2 连接器增量快照
- MongoDB 连接器增量快照
- MySQL 连接器增量快照
- Oracle 连接器增量快照
- PostgreSQL 连接器增量快照
- SQL Server 连接器增量快照
你可以通过创建信号类型为 stop-snapshot 的信号表条目来请求连接器停止正在进行的临时快照。连接器处理该信号后,将停止当前正在进行的快照操作。
以下 Debezium 连接器支持停止临时快照:
- CockroachDB
- Db2
- MariaDB
- MongoDB
- MySQL
- Oracle
- PostgreSQL
- SQL Server
你必须指定信号的 type。data-collections 字段是可选的。将 data-collections 字段留空,表示请求连接器停止当前快照中的所有活动。如果你想让增量快照继续进行,但又想从快照中排除特定集合,请提供一个以逗号分隔的集合名称列表或要排除的正则表达式。连接器处理该信号后,增量快照将继续进行,但会排除你所指定集合中的数据。
增量快照
增量快照是一种特定类型的临时快照。在增量快照中,连接器会捕获你所指定表的基线状态,这与初始快照类似。但与初始快照不同的是,增量快照是分块捕获表数据的,而不是一次性全部捕获。连接器使用水印方法来跟踪快照的进度。
通过分块捕获指定表的初始状态,而不是以单一的整体操作完成,增量快照相比初始快照流程具有以下优势:
- 当连接器捕获指定表的基线状态时,事务日志中近实时事件的流式传输不会中断。
- 如果增量快照过程被中断,可以从停止的位置继续执行。
- 你可以随时发起增量快照。
增量快照暂停信号
你可以通过创建信号类型为 pause-snapshot 的信号表条目,来请求连接器暂停正在进行的增量快照。连接器处理该信号后,将停止当前正在进行的快照操作。因此无法指定数据集合,因为快照处理将在处理该信号时所处的位置暂停。
你可以为以下 Debezium 连接器暂停增量快照:
- CockroachDB
- Db2
- MariaDB
- MongoDB
- MySQL
- Oracle
- PostgreSQL
- SQL Server
| 列 | 值 |
|---|---|
| id | d139b9b7-7777-4547-917d-e1775ea61d41 |
| type | pause-snapshot |
表 9. 暂停增量快照信号记录示例
你必须指定信号的 type。data 字段会被忽略。
增量快照恢复信号
你可以通过创建信号类型为 resume-snapshot 的信号表条目,来请求连接器恢复已暂停的增量快照。连接器处理该信号后,将恢复之前暂停的快照操作。
你可以为以下 Debezium 连接器恢复增量快照:
- CockroachDB
- Db2
- MariaDB
- MongoDB
- MySQL
- Oracle
- PostgreSQL
- SQL Server
| 列 | 值 |
|---|---|
| id | d139b9b7-7777-4547-917d-e1775ea61d41 |
| type | resume-snapshot |
表 10. 恢复增量快照信号记录示例
你必须指定信号的 type。data 字段会被忽略。
有关增量快照的更多信息,请参阅相应连接器文档中的快照主题。
其他资源
- CockroachDB 连接器增量快照
- Db2 连接器增量快照
- MongoDB 连接器增量快照
- MySQL 连接器增量快照
- Oracle 连接器增量快照
- PostgreSQL 连接器增量快照
- SQL Server 连接器增量快照
你可以通过创建一个信号类型为 execute-snapshot、data.type 值为 blocking 的信号,请求连接器启动一次即席阻塞快照。连接器处理该信号后,会运行所请求的快照操作。
与连接器首次启动后运行的初始快照不同,即席阻塞快照发生在运行时,即连接器停止从数据库流式传输变更事件之后。你可以在任意时刻发起即席阻塞快照。
以下 Debezium 连接器支持阻塞快照:
- CockroachDB
- Db2
- MariaDB
- MongoDB
- MySQL
- Oracle
- PostgreSQL
- SQL Server
有关阻塞快照的更多信息,请参阅相应连接器文档中的 Snapshots(快照) 主题。
其他资源
- Db2 连接器即席阻塞快照
- MongoDB 连接器即席阻塞快照
- MySQL 连接器即席阻塞快照
- Oracle 连接器即席阻塞快照
- PostgreSQL 连接器即席阻塞快照
- SQL Server 连接器即席阻塞快照
定义自定义操作
自定义操作使你能够扩展 Debezium 信号框架,以触发默认实现中未提供的操作。你可以在多个连接器中使用自定义操作。
要定义自定义信号操作,你必须定义以下接口:
@FunctionalInterface
public interface SignalAction<P extends Partition> {
/**
* @param signalPayload the content of the signal
* @return true if the signal was processed
*/
boolean arrived(SignalPayload<P> signalPayload) throws InterruptedException;
}io.debezium.pipeline.signal.actions.SignalAction 暴露了一个带有一个参数的方法,该参数表示通过信令通道发送的消息负载。
定义自定义信令操作后,可使用以下 SPI 接口让该自定义操作对信令机制可用:io.debezium.pipeline.signal.actions.SignalActionProvider。
public interface SignalActionProvider {
/**
* Create a map of signal action where the key is the name of the action.
*
* @param dispatcher the event dispatcher instance
* @param connectorConfig the connector config
* @return a concrete action
*/
<P extends Partition> Map<String, SignalAction<P>> createActions(EventDispatcher<P, ? extends DataCollectionId> dispatcher, CommonConnectorConfig connectorConfig);
}你的实现必须返回信号动作的映射(map)。将映射的键设置为动作的名称。该键用作信号的 type。
Debezium 核心模块依赖
自定义动作的 Java 项目对 Debezium 核心模块具有编译依赖。请在项目的 pom.xml 文件中添加以下编译依赖:
<dependency>
<groupId>io.debezium</groupId>
<artifactId>debezium-connector-common</artifactId>
<version>${version.debezium}</version>
</dependency>在上例中,占位符 ${version.debezium} 表示 Debezium 连接器的版本。请在 pom.xml 文件 的 <properties> 部分中为 version.debezium 属性指定一个值。例如:
<properties>
<version.debezium>3.6.3.Final</version.debezium>
</properties>在 META-INF/services/io.debezium.pipeline.signal.actions.SignalActionProvider 文件中声明你的提供者实现。
部署自定义操作
前提条件
- 你已有自定义操作的 Java 程序。
操作步骤
要在 Debezium 连接器中使用自定义操作,请将 Java 项目导出为 JAR 文件,然后把该文件复制到包含你想使用的每个 Debezium 连接器 JAR 文件的目录中。
例如,在典型部署中,Debezium 连接器文件存储在 Kafka Connect 目录(
/kafka/connect)的子目录中,每个连接器的 JAR 位于各自的子目录中(/kafka/connect/debezium-connector-db2、/kafka/connect/debezium-connector-mysql等)。
| 要在多个连接器中使用自定义操作,必须在每个连接器的子目录中各放置一份自定义信令通道 JAR 文件的副本。 |
|---|
评论
登录后参与评论
KnowForge