接收来自 Debezium 的通知
Debezium 通知
概述
Debezium 通知提供了一种获取连接器状态信息的机制。通知可以通过以下通道发送:
SinkNotificationChannel
通过 Connect API 将通知发送到已配置的主题。
LogNotificationChannel
将通知追加到日志中。
JmxNotificationChannel
以属性的形式将通知暴露在 JMX Bean 中。
自定义
将通知发送到你自行实现的自定义通道。
Debezium 通知格式
通知消息包含以下信息:
属性说明
id
分配给该通知的唯一标识符。对于增量快照通知,id 与发送 execute-snapshot 信号时使用的 id 相同。
由于代码的限制,在同时进行多个增量快照时,id 与原始信号并非一一对应。在这种情况下,所有通知都会使用最后发送的信号的 id。 |
|---|
aggregate_type
通知所关联的聚合根的数据类型。在领域驱动设计中,导出的事件始终应当引用某个聚合。
type
提供与 aggregate_type 字段中指定的事件相关的状态信息。
additional_data
一个 Map<String,String>,包含关于该通知的详细信息。示例请参见 Debezium 关于增量快照进度的通知。
timestamp
通知创建的时间。该值表示自 UNIX 纪元以来经过的毫秒数。
可用的通知
Debezium 通知会传递有关初始快照或增量快照进度的信息。有关特定连接器如何使用通知,请参阅该连接器的文档。
Debezium 关于初始快照状态的通知
以下示例展示了一条典型的通知,它提供了初始快照的状态:
{
"id": "5563ae14-49f8-4579-9641-c1bbc2d76f99",
"aggregate_type": "Initial Snapshot",
"type": "COMPLETED",
"additional_data" : {
"connector_name": "myConnector"
},
"timestamp": "1695817046353"
}以下列表描述了前面初始快照通知示例中的部分字段:
type
type 字段可包含以下值之一:
COMPLETEDABORTEDSKIPPED
下表展示了报告初始快照状态的通知中可能出现的不同负载示例:
状态 负载
STARTED
{
"id":"ff81ba59-15ea-42ae-b5d0-4d74f1f4038f",
"aggregate_type":"Initial Snapshot",
"type":"STARTED",
"additional_data":{
"connector_name":"my-connector"
},
"timestamp": "1695817046353"
}IN_PROGRESS
{
"id":"6d82a3ec-ba86-4b36-9168-7423b0dd5c1d",
"aggregate_type":"Initial Snapshot",
"type":"IN_PROGRESS",
"additional_data":{
"connector_name":"my-connector",
"data_collections":"table1, table2",
"current_collection_in_progress":"table1"
},
"timestamp": "1695817046353"
}字段 data_collection 不受 MongoDB 连接器支持。
TABLE_SCAN_COMPLETED
{
"id":"6d82a3ec-ba86-4b36-9168-7423b0dd5c1d",
"aggregate_type":"Initial Snapshot",
"type":"TABLE_SCAN_COMPLETED",
"additional_data":{
"connector_name":"my-connector",
"data_collection":"table1, table2",
"scanned_collection":"table1",
"total_rows_scanned":"100",
"status":"SUCCEEDED"
},
"timestamp": "1695817046353"
}在前面的示例中,additional_data.status 字段可以取以下值之一:
SQL_EXCEPTION
执行快照时发生了 SQL 异常。
SUCCEEDED
快照成功完成。
目前 MongoDB 连接器尚不支持 total_rows_scanned 和 data_collection 字段。
COMPLETED
{
"id":"ff81ba59-15ea-42ae-b5d0-4d74f1f4038f",
"aggregate_type":"Initial Snapshot",
"type":"COMPLETED",
"additional_data":{
"connector_name":"my-connector"
},
"timestamp": "1695817046353"
}已中止
{
"id":"ff81ba59-15ea-42ae-b5d0-4d74f1f4038f",
"aggregate_type":"Initial Snapshot",
"type":"ABORTED",
"additional_data":{
"connector_name":"my-connector"
},
"timestamp": "1695817046353"
}已跳过
{
"id":"ff81ba59-15ea-42ae-b5d0-4d74f1f4038f",
"aggregate_type":"Initial Snapshot",
"type":"SKIPPED",
"additional_data":{
"connector_name":"my-connector"
},
"timestamp": "1695817046353"
}关于增量快照进度的 Debezium 通知
下表列出了报告增量快照状态的通知中可能出现的各种负载示例:
状态 负载
开始
{
"id":"ff81ba59-15ea-42ae-b5d0-4d74f1f4038f",
"aggregate_type":"Incremental Snapshot",
"type":"STARTED",
"additional_data":{
"connector_name":"my-connector",
"data_collections":"table1, table2"
},
"timestamp": "1695817046353"
}已暂停
{
"id":"068d07a5-d16b-4c4a-b95f-8ad061a69d51",
"aggregate_type":"Incremental Snapshot",
"type":"PAUSED",
"additional_data":{
"connector_name":"my-connector",
"data_collections":"table1, table2"
},
"timestamp": "1695817046353"
}已恢复
{
"id":"a9468204-769d-430f-96d2-b0933d4839f3",
"aggregate_type":"Incremental Snapshot",
"type":"RESUMED",
"additional_data":{
"connector_name":"my-connector",
"data_collections":"table1, table2"
},
"timestamp": "1695817046353"
}从 Debezium 接收通知
已停止
{
"id":"83fb3d6c-190b-4e40-96eb-f8f427bf482c",
"aggregate_type":"Incremental Snapshot",
"type":"ABORTED",
"additional_data":{
"connector_name":"my-connector"
},
"timestamp": "1695817046353"
}处理分块
{
"id":"d02047d6-377f-4a21-a4e9-cb6e817cf744",
"aggregate_type":"Incremental Snapshot",
"type":"IN_PROGRESS",
"additional_data":{
"connector_name":"my-connector",
"data_collections":"table1, table2",
"current_collection_in_progress":"table1",
"maximum_key":"100",
"last_processed_key":"50"
},
"timestamp": "1695817046353"
}表快照已完成
{
"id":"6d82a3ec-ba86-4b36-9168-7423b0dd5c1d",
"aggregate_type":"Incremental Snapshot",
"type":"TABLE_SCAN_COMPLETED",
"additional_data":{
"connector_name":"my-connector",
"data_collection":"table1, table2",
"scanned_collection":"table1",
"total_rows_scanned":"100",
"status":"SUCCEEDED"
},
"timestamp": "1695817046353"
}在上例中,additional_data.status 字段可能包含以下值之一:
EMPTY
该表中没有任何值。
NO_PRIMARY_KEY
无法完成快照;该表没有主键。
SKIPPED
无法为该类型的表执行快照。详情请参阅日志。
SQL_EXCEPTION
执行快照时发生 SQL 异常。
SUCCEEDED
快照已成功完成。
UNKNOWN_SCHEMA
找不到该表的结构(schema)。请查看日志以获取已知表的列表。
完成
{
"id":"6d82a3ec-ba86-4b36-9168-7423b0dd5c1d",
"aggregate_type":"Incremental Snapshot",
"type":"COMPLETED",
"additional_data":{
"connector_name":"my-connector"
},
"timestamp": "1695817046353"
}启用 Debezium 通知
要使 Debezium 能够发出通知,请通过设置 notification.enabled.channels 配置属性来指定通知渠道列表。默认情况下,可以使用以下通知渠道:
sinklogjmx
要使用 sink 通知渠道,还必须将 notification.sink.topic.name 配置属性设置为你希望 Debezium 发送通知的主题名称。 |
|---|
访问 Debezium JMX 通知
要使 Debezium 能够报告通过 JMX Bean 暴露的事件,请完成以下配置步骤:
- 启用 JMX MBean 服务器以暴露通知 Bean。
- 在连接器配置的
notification.enabled.channels属性中添加jmx。 - 将你所使用的 JMX 客户端连接到 MBean 服务器。
通知通过名为 debezium.<connector-type>.management.notifications.<server> 的 Bean 的 Notifications 属性暴露。
下图展示了一条报告增量快照开始的通知:

要丢弃某条通知,请在 Bean 上调用 reset 操作。
这些通知还会以类型为 debezium.notification 的 JMX 通知形式暴露。要使应用程序能够监听 MBean 发出的 JMX 通知,请让应用程序订阅这些通知。
自定义通知渠道
通知机制被设计为可扩展的。你可以根据需要实现渠道,以最适合你环境的方式传递通知。添加通知渠道需要以下几个步骤:
配置自定义通知渠道
自定义通知渠道是实现了 io.debezium.pipeline.notification.channels.NotificationChannel 服务提供者接口(SPI)的 Java 类。例如:
public interface NotificationChannel {
String name();
void init(CommonConnectorConfig config);
void send(Notification notification);
void close();
}以下列表说明了前面 NotificationChannel 接口示例中各行的作用:
String name()
通道的名称。要使 Debezium 能够使用该通道,请在连接器的 notification.enabled.channels 属性中指定此名称。
void init(CommonConnectorConfig config)
初始化该通道所需的特定配置、变量或连接。
void send(Notification notification)
通过该通道发送通知。Debezium 调用此方法来报告其状态。
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.notification.channels.NotificationChannel 文件中声明你的实现。
部署自定义通知渠道
前提条件
- 你已有一个自定义通知渠道的 Java 程序。
操作步骤
要在 Debezium 连接器中使用通知渠道,请将 Java 项目导出为 JAR 文件,并将该文件复制到你想使用的每个 Debezium 连接器的 JAR 文件所在目录中。
例如,在典型部署中,Debezium 连接器文件存储在 Kafka Connect 目录(
/kafka/connect)的子目录中,每个连接器的 JAR 位于各自的子目录(/kafka/connect/debezium-connector-db2、/kafka/connect/debezium-connector-mysql等)。要在连接器中使用信号渠道,请将转换器 JAR 文件添加到该连接器的子目录中。
| 要在多个连接器中使用自定义通知渠道,你必须在每个连接器子目录中都放置一份通知渠道 JAR 文件的副本。 |
|---|
配置连接器以使用自定义通知渠道
在连接器配置中,将自定义通知渠道的名称添加到 notification.enabled.channels 属性中。
评论
登录后参与评论
KnowForge