配置

接收来自 Debezium 的通知

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

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

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 字段可包含以下值之一:

  • COMPLETED
  • ABORTED
  • SKIPPED

下表展示了报告初始快照状态的通知中可能出现的不同负载示例:

状态 负载

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 配置属性来指定通知渠道列表。默认情况下,可以使用以下通知渠道:

  • sink
  • log
  • jmx
要使用 sink 通知渠道,还必须将 notification.sink.topic.name 配置属性设置为你希望 Debezium 发送通知的主题名称。

访问 Debezium JMX 通知

要使 Debezium 能够报告通过 JMX Bean 暴露的事件,请完成以下配置步骤:

  1. 启用 JMX MBean 服务器以暴露通知 Bean。
  2. 在连接器配置的 notification.enabled.channels 属性中添加 jmx。
  3. 将你所使用的 JMX 客户端连接到 MBean 服务器。

通知通过名为 debezium.<connector-type>.management.notifications.<server> 的 Bean 的 Notifications 属性暴露。

下图展示了一条报告增量快照开始的通知:

JMX `Notifications` 属性中的字段

要丢弃某条通知,请在 Bean 上调用 reset 操作。

这些通知还会以类型为 debezium.notification 的 JMX 通知形式暴露。要使应用程序能够监听 MBean 发出的 JMX 通知,请让应用程序订阅这些通知。

自定义通知渠道

通知机制被设计为可扩展的。你可以根据需要实现渠道,以最适合你环境的方式传递通知。添加通知渠道需要以下几个步骤:

  1. 为该渠道创建一个 Java 项目以实现该渠道,并添加 Debezium Connector Core 作为依赖。
  2. 部署通知渠道。
  3. 修改连接器配置,使连接器能够使用自定义通知渠道。

配置自定义通知渠道

自定义通知渠道是实现了 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 属性中。

评论

登录后参与评论

正在加载评论…