配置

向 Debezium 发送信号

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

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

向 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 信号通道已启用。

对于每个要使用信号功能的连接器,你必须显式配置信号功能。

操作步骤

  1. 在源数据库中创建一个信号数据集合表,用于向连接器发送信号。有关信号数据集合所需结构的信息,请参阅信号数据集合的结构。

  2. 对于实现了原生变更数据捕获(CDC)机制的源数据库(如 Db2 或 SQL Server),请为信号表启用 CDC。

  3. 将信号数据集合的名称添加到 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 消费者,用于消费发送到所配置信号主题的信号。

消费者可用的其他配置

消息格式

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 信号

操作步骤

  1. 使用你喜欢的 JMX 客户端(例如 JConsole 或 JDK Mission Control)连接到 MBean 服务器。

  2. 查找 MBean debezium.<connector-type>.management.signals.<server>。该 MBean 暴露了 signal 操作,接受以下输入参数:

    p0

    信号的 ID。

    p1

    信号的类型,例如 execute-snapshot。

    p2

    一个 JSON 数据字段,其中包含与指定信号类型相关的附加信息。

  3. 通过为输入参数提供值来发送 execute-snapshot 信号。在 JSON 数据字段中,包含下表所列的信息:

    表 3. 执行快照的数据字段

    字段 默认值

    type

    incremental

    要运行的快照类型。目前 Debezium 支持 incremental 和 blocking 两种类型。

    data-collections

    无

    一个由逗号分隔的正则表达式数组,用于匹配要包含在快照中的表的全限定名。

    additional-conditions

    无

    一个可选数组,用于指定连接器求值的一组附加条件,以确定要包含在快照中的记录子集。每个附加条件都是一个对象,用于指定即席快照所捕获数据的过滤标准。每个附加条件可以设置以下属性:

    data-collection

    过滤器所应用的数据集合的全限定名。你可以为每个数据集合应用不同的过滤器。

    filter

    指定数据库记录中必须存在的列值,快照才会包含该记录,例如 "color='blue'"。快照进程会使用 filter 的值对数据集合中的记录进行求值,仅捕获包含匹配值的记录。

    你为 filter 属性指定的具体值取决于即席快照的类型:

    • 对于增量快照,你需要指定一个搜索条件片段,例如 "color='blue'",快照会将其追加到查询的条件子句中。
    • 对于阻塞快照,你需要指定完整的 SELECT 语句,例如你可能会设置在 snapshot.select.statement.overrides 属性中的语句。

下图展示了使用 JConsole 发送信号的示例:

使用 JConsole 发送 `execute-snapshot` 信号

启用文件信号通道

你可以通过在连接器配置的 signal.enabled.channels 属性中添加 file 来启用 File 信号通道。启用信号通道后,你必须配置连接器以从文件中读取信号。默认情况下,信号文件创建在连接器 classpath 的根目录下,文件名为 file-signals.txt。如果你想使用其他文件,请在连接器配置中设置 signal.file 属性,并指定文件名和路径。文件路径必须在连接器运行环境中可访问。

消息格式

信号文件中的信号以 JSON 对象的形式表示,由 id、type 和 data 字段组成。

id 字段是信号的唯一标识符,通常是一个 UUID 字符串。

当信号类型设置为 execute-snapshot 时,data 字段必须包含下表中列出的字段:

表 4. 执行快照的数据字段

字段默认值说明
typeincremental要执行的快照类型。目前 Debezium 支持 incremental 和 blocking 两种类型。
data-collectionsN/A一个正则表达式数组(以逗号分隔),用于匹配要包含在快照中的数据集合的完全限定名称。命名格式取决于数据库。
additional-conditionsN/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 发送信号。

添加信号通道包含以下几个步骤:

  1. 创建用于该通道的 Java 项目来实现通道,并将 Debezium Core 添加为依赖项。
  2. 部署自定义信号通道。
  3. 修改连接器配置,使连接器能够使用自定义信号通道。

提供自定义信号通道

自定义信号通道是实现了 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

有关临时快照的更多信息,请参阅相应连接器文档中的 快照 主题。

其他资源

你可以通过创建信号类型为 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
列值
idd139b9b7-7777-4547-917d-e1775ea61d41
typepause-snapshot

表 9. 暂停增量快照信号记录示例

你必须指定信号的 type。data 字段会被忽略。

增量快照恢复信号

你可以通过创建信号类型为 resume-snapshot 的信号表条目,来请求连接器恢复已暂停的增量快照。连接器处理该信号后,将恢复之前暂停的快照操作。

你可以为以下 Debezium 连接器恢复增量快照:

  • CockroachDB
  • Db2
  • MariaDB
  • MongoDB
  • MySQL
  • Oracle
  • PostgreSQL
  • SQL Server
列值
idd139b9b7-7777-4547-917d-e1775ea61d41
typeresume-snapshot

表 10. 恢复增量快照信号记录示例

你必须指定信号的 type。data 字段会被忽略。

有关增量快照的更多信息,请参阅相应连接器文档中的快照主题。

其他资源

你可以通过创建一个信号类型为 execute-snapshot、data.type 值为 blocking 的信号,请求连接器启动一次即席阻塞快照。连接器处理该信号后,会运行所请求的快照操作。

与连接器首次启动后运行的初始快照不同,即席阻塞快照发生在运行时,即连接器停止从数据库流式传输变更事件之后。你可以在任意时刻发起即席阻塞快照。

以下 Debezium 连接器支持阻塞快照:

  • CockroachDB
  • Db2
  • MariaDB
  • MongoDB
  • MySQL
  • Oracle
  • PostgreSQL
  • SQL Server

有关阻塞快照的更多信息,请参阅相应连接器文档中的 Snapshots(快照) 主题。

其他资源

定义自定义操作

自定义操作使你能够扩展 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 文件的副本。

评论

登录后参与评论

正在加载评论…