源连接器

YashanDB(崖山数据库)

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

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

YashanDB 的 Debezium 连接器

概述

YashanDB 是由深圳计算科学研究院自主研发的新一代数据库管理系统。它以经典数据库理论为基础,融合了有界计算、近似计算、并行可扩展、跨模态融合计算等原创计算范式,从而使 YashanDB 能够满足金融、政务、能源等核心行业对高性能、高并发、高安全性的严苛要求。

Debezium YashanDB 连接器捕获并记录 YashanDB 服务器上数据库中发生的行级变更,包括连接器运行期间新增的表。你可以配置连接器,使其针对特定的模式和表子集发出变更事件,并将这些变更事件同步到 Kafka。

有关与此连接器兼容的 YashanDB 版本信息,请参阅 Debezium 发行版概述。

Debezium 可以通过原生的 YStream 数据库程序包从 YashanDB 摄取变更事件。有关 YashanDB 和 YStream 的更多信息,请参阅 YashanDB 官网。

Debezium Yashan连接器的工作原理

为了最佳地配置和运行 Debezium YashanDB 连接器,了解连接器的工作方式非常有帮助。

YStream 机制

Debezium YashanDB 连接器使用 YashanDB 的 YStream 接口从数据库事务日志中捕获变更。YStream 是 YashanDB 中原生的变更数据捕获(CDC)引擎。

YashanDB 连接器的工作方式如下:

  1. 连接器通过 JDBC 连接到 YashanDB 数据库。
  2. 连接器建立与指定 YStream 服务的连接,并使用 YStream 客户端 API 实时获取已提交的事务数据。
  3. 连接器首先执行初始快照,捕获指定表的当前状态。
  4. 快照完成后,连接器切换到流式模式,通过 YStream 持续捕获增量数据变更。

快照

YashanDB 服务器上的重做日志通常配置为不保留数据库的完整历史记录。因此,Debezium YashanDB 连接器无法从日志中检索数据库的全部历史。为了使连接器能够建立数据库当前状态的基线,连接器首次启动时会对数据库执行初始一致性快照。

表 1. snapshot.mode 连接器配置属性的设置

设置 说明
always 每次连接器启动时都执行快照。快照完成后,连接器开始流式传输后续数据库变更的事件记录。
initial 连接器执行数据库快照。快照完成后,连接器开始流式传输后续数据库变更的事件记录。
initial_only

连接器执行数据库快照,并在开始流式传输任何变更事件记录之前停止。连接器不会捕获快照之后发生的任何变更事件。

no_data

连接器捕获所有相关表的结构,执行默认快照工作流中描述的所有步骤,但不会创建 READ 事件来表示连接器启动时的数据集。

recovery

设置此选项可恢复丢失或损坏的数据库模式历史主题。重新启动后,连接器会运行一次快照,从源表重建该主题。

when_needed

连接器启动后,仅在检测到以下情况之一时才执行快照:

  • 无法检测到任何主题偏移量。
  • 之前记录的偏移量指定的日志位置在服务器上不可用。

自定义快照器 SPI

要在标准快照模式提供的功能之外自定义快照行为,您可以实现一个或多个 Debezium 快照器 SPI 接口。这些接口可以控制是否运行快照、如何查询数据以及是否锁定表。

io.debezium.snapshot.spi.Snapshotter

控制连接器是否执行快照。

io.debezium.snapshot.spi.SnapshotQuery

控制快照期间如何查询数据。

io.debezium.snapshot.spi.SnapshotLock

控制连接器在执行快照时是否锁定表。

io.debezium.snapshot.spi.Snapshotter 接口。所有内置的快照模式都实现了此接口。

/**
 * {@link Snapshotter} is used to determine the following details about the snapshot process:
 * <p>
 * - Whether a snapshot occurs. <br>
 * - Whether streaming continues during the snapshot. <br>
 * - Whether the snapshot includes schema (if supported). <br>
 * - Whether to snapshot data or schema following an error.
 * <p>
 * Although Debezium provides many default snapshot modes,
 * to provide more advanced functionality, such as partial snapshots,
 * you can customize implementation of the interface.
 * For more information, see the documentation.
 *
 *
 *
 */
@Incubating
public interface Snapshotter extends Configurable {

    /**
     * @return the name of the snapshotter.
     *
     *
     */
    String name();

    /**
     * @param offsetExists is {@code true} when the connector has an offset context (i.e. restarted)
     * @param snapshotInProgress is {@code true} when the connector is started, but a snapshot is already in progress
     *
     * @return {@code true} if the snapshotter should take a data snapshot
     */
    boolean shouldSnapshotData(boolean offsetExists, boolean snapshotInProgress);

    /**
     * @param offsetExists is {@code true} when the connector has an offset context (i.e. restarted)
     * @param snapshotInProgress is {@code true} when the connector is started, but a snapshot is already in progress
     *
     * @return {@code true} if the snapshotter should take a schema snapshot
     */
    boolean shouldSnapshotSchema(boolean offsetExists, boolean snapshotInProgress);

    /**
     * @return {@code true} if the snapshotter should stream after taking a snapshot
     */
    boolean shouldStream();

    /**
     * @return {@code true} whether the schema can be recovered if database schema history is corrupted.
     */
    boolean shouldSnapshotOnSchemaError();

    /**
     * @return {@code true} whether the snapshot should be re-executed when there is a gap in data stream.
     */
    boolean shouldSnapshotOnDataError();

    /**
     *
     * @return {@code true} if streaming should resume from the start of the snapshot
     * transaction, or {@code false} for when a connector resumes and takes a snapshot,
     * streaming should resume from where streaming previously left off.
     */
    default boolean shouldStreamEventsStartingFromSnapshot() {
        return true;
    }

    /**
     * Lifecycle hook called after the snapshot phase is successful.
     */
    default void snapshotCompleted() {
        // no operation
    }

    /**
     * Lifecycle hook called after the snapshot phase is aborted.
     */
    default void snapshotAborted() {
        // no operation
    }
}

io.debezium.snapshot.spi.SnapshotQuery 接口。所有内置的快照查询模式都实现了该接口。

/**
 * {@link SnapshotQuery} is used to determine the query used during a data snapshot
 *
 *
 */
public interface SnapshotQuery extends Configurable, Service {

    /**
     * @return the name of the snapshot lock.
     *
     *
     */
    String name();

    /**
     * Generate a valid query string for the specified table, or an empty {@link Optional}
     * to skip snapshotting this table (but that table will still be streamed from)
     *
     * @param tableId the table to generate a query for
     * @param snapshotSelectColumns the columns to be used in the snapshot select based on the column
     *                              include/exclude filters
     * @return a valid query string, or none to skip snapshotting this table
     */
    Optional<String> snapshotQuery(String tableId, List<String> snapshotSelectColumns);

}

io.debezium.snapshot.spi.SnapshotLock 接口。所有内置的快照锁模式都实现了该接口。

/**
 * {@link SnapshotLock} is used to determine the table lock mode used during schema snapshot
 *
 *
 */
public interface SnapshotLock extends Configurable, Service {

    /**
     * @return the name of the snapshot lock.
     *
     *
     */
    String name();

    /**
     * Returns a SQL statement for locking the given table during snapshotting, if required by the specific snapshotter
     * implementation.
     */
    Optional<String> tableLockingStatement(Duration lockTimeout, String tableId);

}

有关更多信息,请参阅连接器配置属性表中的 snapshot.mode。

主题名称

默认情况下,YashanDB 连接器会将表中发生的所有 INSERT、UPDATE 和 DELETE 操作的变更事件写入该表专用的单个 Apache Kafka 主题。连接器按照以下约定命名变更事件主题:

topicPrefix.schemaName.tableName

以下列表定义了默认名称中各组成部分的含义:

topicPrefix

由 topic.prefix 连接器配置属性指定的主题前缀。

schemaName

发生该操作所在 schema 的名称。

tableName

发生该操作所在表的名称。

例如,如果服务器名称为 fulfillment,schema 名称为 inventory,且数据库中包含名为 orders、customers 和 products 的表,那么 Debezium YashanDB 连接器会将事件发送到以下 Kafka 主题,数据库中的每个表对应一个主题:

fulfillment.inventory.orders
fulfillment.inventory.customers
fulfillment.inventory.products

连接器采用类似的命名约定来标记其内部的数据库模式历史主题、模式变更主题和事务元数据主题。

如果默认的主题名称不能满足你的需求,可以配置自定义主题名称。要配置自定义主题名称,请在逻辑主题路由 SMT 中指定正则表达式。有关使用逻辑主题路由 SMT 自定义主题命名的更多信息,请参阅主题路由。

模式历史主题

当数据库客户端查询数据库时,客户端使用数据库当前的模式。然而,数据库模式可能随时发生变化,这意味着连接器必须能够识别在每次插入、更新或删除操作被记录时的模式是什么。此外,连接器不一定能将当前模式应用于每个事件。如果某个事件相对较旧,它可能是在应用当前模式之前被记录的。

为确保在模式变更之后发生的事件能够被正确处理,YashanDB 在日志中不仅包含影响数据的行级变更,还包含应用于数据库的 DDL 语句。当连接器在日志中遇到这些 DDL 语句时,会对其进行解析,并更新内存中每个表的模式表示。连接器使用该模式表示来识别每次插入、更新或删除操作发生时表的结构,并生成相应的变更事件。在一个独立的数据库模式历史 Kafka 主题中,连接器会记录所有的 DDL 语句,以及每条 DDL 语句在日志中出现的位置。

当连接器在崩溃或正常停止之后重新启动时,它会从日志中的一个特定位置,即一个特定的时间点开始读取。连接器通过读取数据库模式历史 Kafka 主题并解析所有 DDL 语句(直到连接器开始读取的日志位置),来重建该时间点上存在的表结构。

这个数据库模式历史主题仅供连接器内部使用。除此之外,连接器还可以将模式变更事件发送到另一个面向消费方应用的主题。

模式变更主题

你可以配置 Debezium YashanDB 连接器,使其产生描述应用于数据库中表的结构变更的模式变更事件。连接器将模式变更事件写入名为 <serverName> 的 Kafka 主题,其中 serverName 是在 topic.prefix 配置属性中指定的命名空间。

Debezium 在从新表流式传输数据时,或在表结构发生变更时,会向 schema change 主题发送新消息。

连接器发送到 schema change 主题的消息包含有效负载(payload),并且还可选地包含变更事件消息的 schema。

schema change 事件的 schema 包含以下元素:

name

schema change 事件消息的名称。

type

变更事件消息的类型。

version

schema 的版本。版本是一个整数,每次 schema 发生更改时递增。

fields

变更事件消息中包含的字段。

示例:YashanDB 连接器 schema change 主题的 schema

以下示例展示了 JSON 格式的典型 schema。

{
  "schema": {
    "type": "struct",
    "fields": [
      {
        "type": "string",
        "optional": false,
        "field": "databaseName"
      }
    ],
    "optional": false,
    "name": "io.debezium.connector.yashandb.SchemaChangeKey",
    "version": 1
  },
  "payload": {
    "databaseName": "inventory"
  }
}

架构变更事件消息的负载包含以下元素:

ddl

提供导致架构变更的 SQL CREATE、ALTER 或 DROP 语句。

databaseName

应用这些语句的数据库名称。databaseName 的值用作消息键。

schemaName

应用这些语句的架构(schema)名称。

tableChanges

架构变更后整张表结构的结构化表示。tableChanges 字段包含一个数组,其中列出了该表的每一列。由于结构化表示以 JSON 或 Avro 格式呈现数据,消费者无需先通过 DDL 解析器处理即可轻松读取消息。

当连接器被配置为捕获某张表时,它不仅会将该表的架构变更历史保存到架构变更主题中,还会保存到内部的数据库架构历史主题中。内部数据库架构历史主题仅供连接器使用,不面向消费应用程序直接使用。请确保需要接收架构变更通知的应用程序仅从架构变更主题中消费该信息。

切勿对数据库架构历史主题进行分区。为使数据库架构历史主题正常工作,它必须保持连接器向其发送的事件记录在全局范围内的一致顺序。

为确保该主题不会被分散到多个分区中,请通过以下任一方法设置该主题的分区数:

  • 如果你手动创建数据库架构历史主题,请将分区数指定为 1。
  • 如果你使用 Apache Kafka 代理自动创建数据库架构历史主题,请将 Kafka num.partitions 配置选项的值设置为 1。

示例:发送到 YashanDB 连接器架构变更主题的消息

以下示例展示了一个典型的 JSON 格式架构变更消息。该消息包含表结构的逻辑表示。

{
  "schema": {
  ...
  },
  "payload": {
    "source": {
      "version": "3.6.3.Final",
      "connector": "yashandb",
      "name": "server1",
      "ts_ms": 1780017300692,
      "snapshot": "false",
      "db": "",
      "sequence": null,
      "ts_us": 1780017300692853,
      "ts_ns": 1780017300692853000,
      "schema": "DEBEZIUM",
      "table": "CUSTOMERS",
      "txId": "131072047",
      "scn": "828249295637925888",
      "batch_row_id": 0,
      "position_scn": 828249295637925888,
      "group_lsn": 3692924,
      "group_offset": 220,
      "instance_id": "0"
    },
    "ts_ms": 1780045807728, (1)
    "databaseName": "inventory", (2)
    "schemaName": "DEBEZIUM", (3)
    "ddl": "CREATE TABLE \"DEBEZIUM\".\"CUSTOMERS\" \n   (    \"ID\" NUMBER(9,0) NOT NULL ENABLE, \n    \"NAME\" VARCHAR2(255) \n   )", (4)
    "tableChanges": [ (5)
      {
        "type": "CREATE", (6)
        "id": "\"DEBEZIUM\".\"CUSTOMERS\"", (7)
        "table": { (8)
          "defaultCharsetName": null,
          "primaryKeyColumnNames": [ (9)
            "ID"
          ],
          "columns": [ (10)
            {
              "name": "ID",
              "jdbcType": 2,
              "nativeType": null,
              "typeName": "NUMBER",
              "typeExpression": "NUMBER",
              "charsetName": null,
              "length": 9,
              "scale": 0,
              "position": 1,
              "optional": false,
              "autoIncremented": false,
              "generated": false
            },
            {
              "name": "NAME",
              "jdbcType": 12,
              "nativeType": null,
              "typeName": "VARCHAR2",
              "typeExpression": "VARCHAR2",
              "charsetName": null,
              "length": 255,
              "scale": null,
              "position": 2,
              "optional": true,
              "autoIncremented": false,
              "generated": false
            }
          ],
          "attributes": [ (11)
            {
              "customAttribute": "attributeValue"
            }
          ]
        }
      }
    ]
  }
}

表 2. 发送到模式更改主题的消息中各字段的说明

编号 字段名 说明
1 ts_ms 可选字段,用于显示连接器处理事件的时间。该时间基于运行 Kafka Connect 任务的 JVM 中的系统时钟。
2 databaseName 标识包含该变更的数据库。
3 schemaName 标识包含该变更的模式。
4 ddl 此字段包含导致模式更改的 DDL。
5 tableChanges 一个数组,包含一项或多项由 DDL 命令生成的模式更改。
6 type 描述变更的类型。type 可设置为以下值之一:CREATE(表已创建)、ALTER(表已修改)、DROP(表已删除)。
7 id 被创建、修改或删除的表的完整标识符。对于表重命名的情况,该标识符是 <旧表名>,<新表名> 的拼接。
8 table 表示应用变更之后的表元数据。
9 primaryKeyColumnNames 构成该表主键的列的列表。
10 columns 变更表中每一列的元数据。
11 attributes 每次表变更的自定义属性元数据。

在连接器发送到模式更改主题的消息中,消息键是包含该模式更改的数据库名称。在下面的示例中,payload 字段包含 databaseName 键:

{
  "schema": {
    "type": "struct",
    "fields": [
      {
        "type": "string",
        "optional": false,
        "field": "databaseName"
      }
    ],
    "optional": false,
    "name": "io.debezium.connector.yashandb.SchemaChangeKey",
    "version": 1
  },
  "payload": {
    "databaseName": "inventory"
  }
}

事务元数据

Debezium 可以生成表示事务元数据边界的事件,并丰富数据变更事件消息。

Debezium 接收事务元数据的限制

Debezium 仅会注册并接收连接器部署之后发生的事务的元数据。连接器部署之前发生的事务的元数据不可用。

数据库事务由一个位于 BEGIN 和 END 关键字之间的语句块来表示。Debezium 会为每个事务中的 BEGIN 和 END 分隔符生成事务边界事件。事务边界事件包含以下字段:

status

BEGIN 或 END。

id

唯一事务标识符的字符串表示形式。

ts_ms

事务边界事件(BEGIN 或 END 事件)在数据源处发生的时间。

event_count(针对 END 事件)

该事务发出的事件总数。

data_collections(针对 END 事件)

一个由 data_collection 和 event_count 元素成对组成的数组,用于表示连接器为源自某个数据集合的变更所发出的事件数量。

以下示例展示了一个典型的事务边界消息:

示例:YashanDB 连接器事务边界事件

{
  "status": "BEGIN",
  "id": "5.6.641",
  "ts_ms": 1486500577125,
  "event_count": null,
  "data_collections": null
}

{
  "status": "END",
  "id": "5.6.641",
  "ts_ms": 1486500577691,
  "event_count": 2,
  "data_collections": [
    {
      "data_collection": "inventory.DEBEZIUM.CUSTOMERS",
      "event_count": 1
    },
    {
      "data_collection": "inventory.DEBEZIUM.ORDERS",
      "event_count": 1
    }
  ]
}

数据变更事件

YashanDB 连接器发出的每个数据变更事件都包含一个键和一个值。键和值的结构取决于变更事件所源自的表。有关 Debezium 如何构造主题名称的信息,请参阅主题名称。

Debezium YashanDB 连接器确保所有 Kafka Connect schema 名称都是合法的 Avro schema 名称。要成为合法的 Avro schema 名称,逻辑服务器名称必须以字母字符或下划线([a-z,A-Z,_])开头。逻辑服务器名称中的其余字符,以及 schema 和表名称中的所有字符,都必须是字母数字字符或下划线([a-z,A-Z,0-9,\_])。连接器会自动将非法字符替换为下划线字符。

如果多个逻辑服务器名称、schema 名称或表名称之间的唯一区别字符是非法字符,并且这些字符被替换为下划线,就可能导致意外的命名冲突。

Debezium 和 Kafka Connect 是围绕事件消息的持续流设计的。但是,这些事件的结构可能会随时间变化,这对主题消费者来说可能难以处理。为了便于处理可变的事件结构,Kafka Connect 中的每个事件都是自包含的。每条消息的键和值都由两部分组成:schema 和 payload。schema 描述 payload 的结构,而 payload 则包含实际数据。

变更事件键

对于每个发生变更的表,变更事件键的结构是:在事件创建时,为表的主键(或唯一键约束)中的每个列都存在一个字段。

例如,考虑以下针对在 inventory 数据库 schema 中定义的 customers 表的 SQL:

CREATE TABLE customers (
  ID INT NOT NULL PRIMARY KEY,
  NAME VARCHAR(255)
);

如果 <topic.prefix>.transaction 配置属性的值设置为 server1,那么数据库中 customers 表发生的每个变更事件,其 JSON 表示都具有如下的键结构:

{
    "schema": {
        "type": "struct",
        "fields": [
            {
                "type": "int32",
                "optional": false,
                "field": "ID"
            }
        ],
        "optional": false,
        "name": "server1.inventory.customers.Key"
    },
    "payload": {
        "ID": 1001
    }
}

键的 schema 部分包含一个 Kafka Connect schema,用于描述键部分的内容。在前面的示例中,payload 的值不是可选的,其结构由名为 server1.inventory.customers.Key 的 schema 定义,并且包含一个必需的 ID 字段,类型为 int32。键的 payload 字段的值表明它确实是一个结构(在 JSON 中就是一个对象),其中只有一个 ID 字段,其值为 1001。

因此,你可以将该键理解为描述 inventory.customers 表中的某一行(由名为 server1 的连接器输出),该行的 ID 主键列的值为 1001。

更改事件值

更改事件消息中值的结构与消息中更改事件的消息键的结构相同,同时包含一个 schema 部分和一个 payload 部分。

更改事件值的 payload

更改事件值的 payload 部分中的 envelope 结构包含以下字段:

op

一个必填字段,包含描述操作类型的字符串值。YashanDB 连接器更改事件值的 payload 中的 op 字段包含以下值之一:c(创建或插入)、u(更新)、d(删除)或 r(读取,表示快照)。

before

一个可选字段,如果存在,则描述事件发生之前行的状态。其结构由 Kafka Connect schema server1.inventory.customers.Value 描述,server1 连接器将该 schema 用于 inventory.customers 表中的所有行。

after

一个可选字段,如果存在,则包含变更发生之后行的状态。其结构由与 before 字段所使用的相同的 server1.inventory.customers.Value Kafka Connect schema 描述。

source

一个必填字段,包含描述事件源元数据的结构。对于 YashanDB 连接器,该结构包含以下字段:

  • Debezium 版本。
  • 连接器类型和名称。
  • 时间戳(ts_ms、ts_us、ts_ns),表示连接器处理事件的时间,基于运行 Kafka Connect 任务的 JVM 的系统时钟。对于快照,时间戳表示快照发生的时间。
  • 该事件是否属于正在进行的快照。
  • 数据库和 schema 名称。
  • 表名。
  • 事务 ID(txId),快照中不包含。
  • 批量操作(如批量插入)中的行序号(batch_row_id)。
  • 该逻辑日志条目所属事务的提交 SCN(position_scn)。
  • 该逻辑日志条目所属日志组的 SCN(group_lsn)。
  • 逻辑日志条目在其日志组中的物理偏移量(group_offset)。
  • 该事务所属实例的实例标识(instance_id)。

ts_ms

提供以毫秒为单位的时间戳。

ts_us

提供以微秒为单位的时间戳。

ts_ns

提供以纳秒为单位的时间戳。

变更事件值的 Schema

事件消息 值 的 schema 部分包含一个模式,用于描述负载的信封(envelope)结构及其嵌套字段。

create 事件

以下示例展示了 变更事件键 示例中所描述的 customers 表的 create 事件的值:

{
    "schema": {
        "type": "struct",
        "fields": [
            {
                "type": "struct",
                "fields": [
                    {
                        "type": "int32",
                        "optional": false,
                        "field": "ID"
                    },
                    {
                        "type": "string",
                        "optional": true,
                        "field": "NAME"
                    }
                ],
                "optional": true,
                "name": "server1.inventory.customers.Value",
                "field": "before"
            },
            {
                "type": "struct",
                "fields": [
                    {
                        "type": "int32",
                        "optional": false,
                        "field": "ID"
                    },
                    {
                        "type": "string",
                        "optional": true,
                        "field": "NAME"
                    }
                ],
                "optional": true,
                "name": "server1.inventory.customers.Value",
                "field": "after"
            },
            {
                "type": "struct",
                "fields": [
                    {
                        "type": "string",
                        "optional": true,
                        "field": "version"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "connector"
                    },
                    {
                        "type": "string",
                        "optional": false,
                        "field": "name"
                    },
                    {
                        "type": "int64",
                        "optional": true,
                        "field": "ts_ms"
                    },
                    {
                        "type": "string",
                        "optional": true,
                        "field": "snapshot"
                    },
                    {
                        "type": "string",
                        "optional": true,
                        "field": "db"
                    },
                    {
                        "type": "string",
                        "optional": true,
                        "field": "sequence"
                    },
                    {
                        "type": "int64",
                        "optional": true,
                        "field": "ts_us"
                    },
                    {
                        "type": "int64",
                        "optional": true,
                        "field": "ts_ns"
                    },
                    {
                        "type": "string",
                        "optional": true,
                        "field": "schema"
                    },
                    {
                        "type": "string",
                        "optional": true,
                        "field": "table"
                    },
                    {
                        "type": "string",
                        "optional": true,
                        "field": "txId"
                    },
                    {
                        "type": "int64",
                        "optional": true,
                        "field": "batch_row_id"
                    },
                    {
                        "type": "int64",
                        "optional": true,
                        "field": "position_scn"
                    },
                    {
                        "type": "int64",
                        "optional": true,
                        "field": "group_lsn"
                    },
                    {
                        "type": "int64",
                        "optional": true,
                        "field": "group_offset"
                    },
                    {
                        "type": "string",
                        "optional": true,
                        "field": "instance_id"
                    }
                ],
                "optional": false,
                "name": "io.debezium.connector.yashandb.Source",
                "field": "source"
            },
            {
                "type": "string",
                "optional": false,
                "field": "op"
            },
            {
                "type": "int64",
                "optional": true,
                "field": "ts_ms"
            },
            {
                "type": "int64",
                "optional": true,
                "field": "ts_us"
            },
            {
                "type": "int64",
                "optional": true,
                "field": "ts_ns"
            }
        ],
        "optional": false,
        "name": "server1.inventory.customers.Envelope"
    },
    "payload": {
        "before": null,
        "after": {
            "ID": 1001,
            "NAME": "Test Record"
        },
        "source": {
            "version": "3.6.3.Final",
            "connector": "yashandb",
            "name": "server1",
            "ts_ms": 1780017300692,
            "snapshot": "false",
            "db": "",
            "sequence": null,
            "ts_us": 1780017300692853,
            "ts_ns": 1780017300692853000,
            "schema": "DEBEZIUM",
            "table": "CUSTOMERS",
            "txId": "131072047",
            "batch_row_id": 0,
            "position_scn": 828249295637925888,
            "group_lsn": 3692924,
            "group_offset": 220,
            "instance_id": "0"
        },
        "op": "c",
        "ts_ms": 1688252618953,
        "ts_us": 1688252618953000,
        "ts_ns": 1688252618953000000
    }
}

以下列表描述了前述 create 事件消息值部分中的部分字段:

schema

指定事件值的架构。该架构描述了值的负载(payload)结构。只要表的架构保持不变,Debezium 为某张表发出的每个变更事件都使用相同的值架构。

name

schema 部分可以包含多个 name 字段。每个 name 字段指定事件值负载中某个字段所对应的架构。

server1.inventory.customers.Value 是负载中 before 和 after 字段的架构。该架构特定于 customers 表。

before 和 after 字段的架构名称采用 logicalName.schemaName.tableName.Value 的形式。这种格式可确保架构名称在数据库内唯一。在使用 Avro 转换器的环境中,唯一的架构名称可确保每个逻辑源中每张表的 Avro 架构都拥有各自的发展演进历史。

"name": "io.debezium.connector.yashandb.Source"

io.debezium.connector.yashandb.Source 是负载中 source 字段的架构。该架构特定于 YashanDB 连接器,连接器对其生成的所有事件都使用此架构。

"name": "server1.inventory.customers.Envelope"

指定负载整体结构所对应架构的名称。该架构名称由以下几部分组成:

server1

指定生成此事件的连接器的名称。

inventory

指定包含被更改表的数据库。

customers

指定被更改的表。

payload

指定被更改行的实际数据。

由于事件的 JSON 表示同时包含消息的架构部分和负载部分,它通常比其所描述的行本身更大。为了减小连接器流式传输到 Kafka 主题的消息体积,可以使用 Avro 转换器。

op

指定导致连接器生成此事件的操作类型。在此示例中,c 表示执行了 create 操作,产生了一条新行。该字段可以包含以下值之一:

c创建一行。
u更新一行。
d删除一行。
r读取一行(仅适用于快照)。

ts_ms、ts_us、ts_ns

分别以毫秒、微秒和纳秒显示时间戳,表示连接器处理该事件的时间。该时间基于运行 Kafka Connect 任务的 JVM 中的系统时钟。

before

一个可选字段,指定事件发生前的行状态。由于前例中的 op 字段为 c(create,创建),该变更事件描述的是新数据,因此 before 字段的值为 null。

after

一个可选字段,指定事件发生后的行状态。在此示例中,after 字段包含新行的 ID 和 NAME 列的值。

source

一个必填字段,描述事件的源元数据。此字段包含的信息可用于将该事件与其他事件进行比较,包括事件的来源、事件发生的顺序,以及事件是否属于同一事务。源元数据提供以下信息:

version

Debezium 版本。

connector

连接器类型。

name

连接器名称。

ts_ms、ts_us、ts_ns

以毫秒、微秒和纳秒为单位显示的时间戳,表示该变更在数据库中发生的时间。

snapshot

指定该事件是否由快照操作产生。

db

包含新行的数据库名称。

schema

包含新行的 Schema 名称。

table

包含新行的表名。

txId

事务标识符。

batch_row_id

在批处理操作(如批量插入)中的行序号。

position_scn

该逻辑日志条目所属事务的提交 SCN。

group_lsn

该逻辑日志条目所属日志组的 SCN。

group_offset

该逻辑日志条目在其日志组内的物理偏移量。

instance_id

该事务所属实例的实例标识符。

update(更新)事件

以下示例展示了一个 update 变更事件,它是连接器从与前例 create(创建)事件相同的表中捕获的。

{
    "schema": { ... },
    "payload": {
        "before": {
            "ID": 1001,
            "NAME": "Test Record"
        },
        "after": {
            "ID": 1001,
            "NAME": "Updated Record"
        },
        "source": {
            "version": "3.6.3.Final",
            "connector": "yashandb",
            "name": "server1",
            "ts_ms": 1780017300692,
            "snapshot": "false",
            "db": "",
            "sequence": null,
            "ts_us": 1780017300692853,
            "ts_ns": 1780017300692853000,
            "schema": "DEBEZIUM",
            "table": "CUSTOMERS",
            "txId": "131072047",
            "batch_row_id": 0,
            "position_scn": 828249295637925888,
            "group_lsn": 3692924,
            "group_offset": 220,
            "instance_id": "0"
        },
        "op": "u",
        "ts_ms": 1688252619000,
        "ts_us": 1688252619000000,
        "ts_ns": 1688252619000000000
    }
}

该负载的结构与 create(插入)事件的负载相同,但以下值有所不同:

  • op 字段的值为 u,表示该行因更新而发生变化。
  • before 字段显示行的先前状态,其中包含数据库提交更新之前的值。
  • after 字段显示更新后的行状态,其中 NAME 的值已被设置为 Updated Record。
  • source 字段的结构与之前相同,但值不同,因为连接器是从日志中的不同位置捕获该事件的。
  • ts_ms 字段显示表示 Debezium 处理该事件时间的时间戳。

payload 部分还揭示了其他一些有用的信息。例如,通过比较 before 和 after 结构,我们可以确定某行是如何因一次提交而发生变化的。source 结构提供了 YashanDB 对此次变更的记录信息,从而提供了可追溯性。它还能让我们了解该事件相对于本主题及其他主题中其他事件的发生时间关系。它是在另一个事件之前、之后,还是与另一个事件同属一次提交?

delete 事件

以下示例展示了前述 create 和 update 事件示例中所涉及表的 delete 事件。delete 事件的 schema 部分与这些事件的 schema 部分完全相同。

{
    "schema": { ... },
    "payload": {
        "before": {
            "ID": 1001,
            "NAME": "Updated Record"
        },
        "after": null,
        "source": {
            "version": "3.6.3.Final",
            "connector": "yashandb",
            "name": "server1",
            "ts_ms": 1780017300692,
            "snapshot": "false",
            "db": "",
            "sequence": null,
            "ts_us": 1780017300692853,
            "ts_ns": 1780017300692853000,
            "schema": "DEBEZIUM",
            "table": "CUSTOMERS",
            "txId": "131072047",
            "batch_row_id": 0,
            "position_scn": 828249295637925888,
            "group_lsn": 3692924,
            "group_offset": 220,
            "instance_id": "0"
        },
        "op": "d",
        "ts_ms": 1688252620000,
        "ts_us": 1688252620000000,
        "ts_ns": 1688252620000000000
    }
}

该负载显示的行与前面的 create 和 update 事件具有相同的键,但包含不同的值:

  • op 字段为 d,表示该行已被删除。
  • before 字段包含该行在随数据库提交被删除之前的值。
  • after 字段为 null,表示该行已不复存在。
  • source 字段与 create 事件中的 source 字段结构相同,但部分字段的值有所不同。
  • ts_ms 显示的时间戳表示 Debezium 处理此事件的时间。

delete 事件为消费者提供了处理该行删除所需的信息。

YashanDB 连接器的事件设计为与 Kafka 日志压缩配合使用,该功能允许在至少保留每个键的最新一条消息的前提下,删除一些较旧的消息。这样 Kafka 就能在确保主题包含完整数据集的同时回收存储空间,从而使该数据集可用于重新加载基于键的状态。

当某行被删除时,前面示例中显示的 delete 事件值仍然可用于日志压缩,因为 Kafka 能够删除所有使用相同键的较早消息。必须将消息值设置为 null,以指示 Kafka 删除所有共享相同键的消息。为了实现这一点,默认情况下,Debezium YashanDB 连接器总会在 delete 事件之后紧跟一个特殊的 tombstone(墓碑)事件,该事件具有相同的键,但值为 null。

truncate 事件

truncate 变更事件表示某个表已被截断。此时消息键为 null,消息值如下所示:

{
    "schema": { ... },
    "payload": {
        "before": null,
        "after": null,
        "source": { (1)
            "version": "3.6.3.Final",
            "connector": "yashandb",
            "name": "my_topic",
            "ts_ms": 1780017300692,
            "snapshot": "false",
            "db": "",
            "sequence": null,
            "ts_us": 1780017300692853,
            "ts_ns": 1780017300692853000,
            "schema": "DEBEZIUM",
            "table": "CUSTOMERS",
            "txId": "131072047",
            "batch_row_id": 0,
            "position_scn": 828249295637925888,
            "group_lsn": 3692924,
            "group_offset": 220,
            "instance_id": "0"
        },
        "op": "t", (2)
        "ts_ms": 1688252630000, (3)
        "ts_us": 1688252630000000, (3)
        "ts_ns": 1688252630000000000 (3)
    }
}

表 3. truncate 事件值字段说明

序号 字段名 说明
1 source 必填字段,描述事件的源元数据。在 truncate 事件值中,source 字段的结构与同一表的 create、update 和 delete 事件相同,提供以下元数据:

- Debezium 版本
- 连接器类型和名称
- 变更在数据库中发生的时间戳(ts_ms、ts_us、ts_ns)
- 该事件是否属于快照的一部分(truncate 事件始终为 false)
- 包含该表的数据库和 schema
- 表名
- 事务标识符(txId)
- 批量操作(如批量插入)中的行序号(batch_row_id)
- 该逻辑日志条目所属事务的提交 SCN(position_scn)
- 该逻辑日志条目所属的日志组的 SCN(group_lsn)
- 该逻辑日志条目在其日志组内的物理偏移量(group_offset)
- 该事务所属实例的实例标识符(instance_id)
2 op 必填字符串,描述操作类型。op 字段的值为 t,表示该表已被截断。
3 ts_ms、ts_us、ts_ns 可选字段,显示连接器处理该事件的时间。该时间基于运行 Kafka Connect 任务的 JVM 中的系统时钟。

在 source 对象中,ts_ms 表示变更在数据库中发生的时间。通过比较 payload.source.ts_ms 的值与 payload.ts_ms 的值,可以确定源数据库更新与 Debezium 之间的延迟。

如果单个 TRUNCATE 操作影响多个表,连接器会为每个被截断的表发出一条 truncate 变更事件记录。

truncate 事件表示对整张表所做的变更,且它没有消息键。因此,对于具有多个分区的主题,与某张表相关的变更事件(create、update 等)或 truncate 事件没有顺序保证。例如,如果消费者从多个分区读取某张表的事件,它可能在从一个分区收到删除该表所有数据的 truncate 事件之后,才从另一个分区收到该表的 update 事件。只有使用单个分区的主题才能保证顺序。

如果你不希望连接器捕获 truncate(截断)事件,可以将 skipped.operations 选项设置为过滤掉这些事件。

数据类型映射

当 Debezium YashanDB 连接器检测到表中某行的值发生更改时,会发出一个表示该更改的更改事件。每条更改事件记录的结构与原始表相同,事件记录中的每个列值都对应一个字段。表列的数据类型决定了连接器在更改事件字段中如何表示该列的值。

对于表中的每一列,Debezium 会将源数据类型映射为一种字面量类型。

字面量类型描述了 Debezium 如何使用以下 Kafka Connect 模式类型之一来字面量地表示值:INT8、INT16、INT32、INT64、FLOAT32、FLOAT64、BOOLEAN、STRING、BYTES、ARRAY、MAP、STRUCT。

如果某一列未被映射,连接器将忽略该列的更改。连接器会映射表中的所有其他列,并为这些列生成更改事件。如果某一列未被映射,它将不会包含在更改事件中。

字符类型

下表描述了连接器如何映射字符数据类型。

YashanDB 数据类型字面量类型(模式类型)语义类型(模式名称)及说明
CHAR[(M)]STRING不适用
NCHAR[(M)]STRING不适用
VARCHAR[(M)]STRING不适用
NVARCHAR[(M)]STRING不适用

表 4. YashanDB 字符类型的映射

二进制和字符 LOB 类型

下表描述了连接器如何映射二进制和字符大对象(LOB)数据类型。

YashanDB 数据类型字面量类型(schema 类型)语义类型(schema 名称)及说明
BLOBBYTES根据连接器配置中 lob.enabled 属性的设置,连接器会将该类型的 LOB 值映射为原始字节或 base64 编码的字符串。
CLOBSTRING不适用
NCLOBSTRING不适用
RAWBYTES根据连接器配置中 lob.enabled 属性的设置,连接器会将该类型的 LOB 值映射为原始字节或 base64 编码的字符串。

表 5. YashanDB 二进制和字符 LOB 类型的映射

数值类型

下表描述了连接器如何映射数值数据类型。

表 6. YashanDB 数值类型的映射

YashanDB 数据类型 字面量类型(schema 类型) 语义类型(schema 名称)及说明

TINYINT

INT8

不适用

SMALLINT

INT16

不适用

INT

INT32

不适用

BIGINT

INT64

不适用

FLOAT

FLOAT32

不适用

DOUBLE

FLOAT64

不适用

NUMBER

BYTES / INT8 / INT16 / INT32 / INT64

org.apache.kafka.connect.data.Decimal

根据 decimal.handling.mode 属性的取值,连接器会将 NUMBER 映射为 BYTES 表示形式,或映射为某种 INT 类型。对于标度为负数的 NUMBER,请使用 decimal.handling.mode=string 以避免序列化问题。

BIT(1)

BOOLEAN

不适用

BIT(n)

BYTES

不适用

BOOLEAN

BOOLEAN

不适用

时间类型

下表描述了连接器如何映射时间数据类型。连接器转换时间类型的方式取决于 time.precision.mode 配置属性。

表 7. 当 time.precision.mode 为 connect 时 YashanDB 时间类型的映射

YashanDB 数据类型 字面量类型(schema 类型) 语义类型(schema 名称)及说明

DATE

INT64

io.debezium.time.Timestamp

表示自 UNIX 纪元以来的毫秒数。

TIME

INT64

io.debezium.time.MicroTime

表示自午夜起经过的微秒数。

TIMESTAMP

INT64

io.debezium.time.MicroTimestamp

表示自 UNIX 纪元以来的微秒数。

INTERVAL YEAR TO MONTH

FLOAT64

io.debezium.time.MicroDuration

表示该区间内的月数,以带小数精度的年数表示。

INTERVAL DAY TO SECOND

FLOAT64

io.debezium.time.MicroDuration

表示该区间内的微秒数。

其他类型

下表描述了连接器如何映射其他数据类型。

YashanDB 数据类型字面量类型(schema 类型)语义类型(schema 名称)及说明
ROWIDSTRING不适用

表 8. YashanDB 其他类型的映射

LOB 处理

YashanDB 仅在 SQL 语句中显式设置或修改时,才为 CLOB、NCLOB 和 BLOB 数据类型,以及大小超过 32000 的 VARCHAR、NVARCHAR 和 RAW 数据类型提供列值。对于未被修改的 LOB 列,连接器不会将这些列的值包含在变更事件中。

要捕获 LOB 值并将其序列化到变更事件中,请将 lob.enabled 选项设置为 true。启用 LOB 处理后,连接器在发出 LOB 数据时会产生一定的性能开销。

Decimal 处理

您可以通过修改连接器的 decimal.handling.mode 配置属性的值,来更改连接器映射 NUMBER 数据类型的方式。

当该属性设置为默认值 precise 时,连接器会将这些数据类型映射为 Kafka Connect 的 org.apache.kafka.connect.data.Decimal 逻辑类型。

当该属性值设置为 double 或 string 时,连接器将使用替代映射方式。

YashanDB 数据类型precisedoublestring
NUMBERorg.apache.kafka.connect.data.DecimalFLOAT64STRING

表 9. 数值数据类型的映射

DATE、TIME 和 TIMESTAMP 类型默认映射为 INT64(时间戳形式)。如果您希望将它们映射为固定格式的字符串(例如 yyyy-MM-dd HH:mm:ss.SSSSSS),请参阅数据类型转换章节。

自定义转换器

默认情况下,Debezium YashanDB 连接器提供了多个针对 YashanDB 数据类型的 CustomConverter 实现。这些自定义转换器根据连接器配置为特定数据类型提供替代映射。要向连接器添加 CustomConverter,请按照自定义转换器文档中的说明操作。

Debezium YashanDB 连接器提供以下自定义转换器:

TimestampToStringConverter

TimestampToStringConverter 将 TIMESTAMP 类型的数据转换为自定义格式的字符串。

属性说明
yashandb_timestamp_formatter.typeio.debezium.connector.yashandb.converters.TimestampToStringConverter
yashandb_timestamp_formatter.format.datetime日期时间格式,例如:yyyy-MM-dd HH:mm:ss.SSSSSS

表 10. 转换器配置

DateToStringConverter

DateToStringConverter 将 DATE 类型的数据转换为自定义格式的字符串。

属性说明
yashandb_date_formatter.typeio.debezium.connector.yashandb.converters.DateToStringConverter
yashandb_date_formatter.format.date日期格式,例如:yyyy-MM-dd

表 11. 转换器配置

TimeToStringConverter

TimeToStringConverter 将 TIME 类型的数据转换为自定义格式的字符串。

属性说明
yashandb_time_formatter.typeio.debezium.connector.yashandb.converters.TimeToStringConverter
yashandb_time_formatter.format.time时间格式,例如:HH:mm:ss.SSSSSS

表 12. 转换器配置

使用示例

在配置中指定以下内容:

# Name two converters: yashandb_timestamp_formatter for TIMESTAMP, yashandb_date_formatter for DATE
"converters": "yashandb_timestamp_formatter,yashandb_date_formatter"
# Bind yashandb_timestamp_formatter to TimestampToStringConverter class
"yashandb_timestamp_formatter.type": "io.debezium.connector.yashandb.converters.TimestampToStringConverter"
# Format TIMESTAMP data as yyyy-MM-dd HH:mm:ss.SSSSSS
"yashandb_timestamp_formatter.format.datetime": "yyyy-MM-dd HH:mm:ss.SSSSSS"
# Bind yashandb_date_formatter to DateToStringConverter class
"yashandb_date_formatter.type": "io.debezium.connector.yashandb.converters.DateToStringConverter"
# Format DATE data as yyyy-MM-dd
"yashandb_date_formatter.format.date": "yyyy-MM-dd"

配置 YashanDB

在部署和运行 Debezium YashanDB 连接器之前,请先调整 YashanDB 数据库配置,以确保兼容性。

排除捕获的模式

当 Debezium YashanDB 连接器捕获表时,会自动排除以下模式中的表:

  • SYS
  • MDSYS
  • XA_SYS

要使连接器能够捕获某张表的更改,该表必须使用未在上述列表中列出的模式。

配置 YStream 内存池

增量数据依赖 YStream 从 YashanDB 实时获取已提交的数据。当您将 YashanDB 用作包含增量同步的任务的数据源时,必须首先在 YashanDB 上为 YStream 分配一个内存池:

ALTER SYSTEM SET STREAM_POOL_SIZE = '512M';
未能配置此参数可能导致任务失败。此参数为全局参数。更多信息请参阅 YashanDB 文档。

启用补充日志

读取增量数据变更需要启用补充日志。

数据库级补充日志

要允许连接器监视数据库中的所有对象(包括新对象),需在 YashanDB 中启用数据库级补充日志:

ALTER DATABASE ADD SUPPLEMENTAL LOG TABLE TYPE (HEAP);
ALTER DATABASE ADD SUPPLEMENTAL LOG DATA (ALL) COLUMNS;

表级补充日志

如果希望连接器仅监控特定的表,请在 YashanDB 中启用表级补充日志:

ALTER TABLE tablename ADD SUPPLEMENTAL LOG DATA (ALL) COLUMNS;
需要注意的是,未能启用补充日志或错误地启用补充日志,可能会导致数据丢失或任务失败。

为连接器创建用户账户

为了让 Debezium YashanDB 连接器能够捕获变更事件,它必须以具有特定权限的 YashanDB 用户身份运行。

YashanDB 连接器用户正常运行需要以下权限:

-- Create user
CREATE USER username IDENTIFIED BY password;

-- Session permissions
GRANT CREATE SESSION TO username;
GRANT ALTER SESSION TO username;

-- Debezium specific permissions
GRANT SELECT ANY TABLE TO username;
GRANT LOCK ANY TABLE TO username;
GRANT FLASHBACK ANY TABLE TO username;
GRANT SELECT ON V_$DATABASE TO username;
GRANT SELECT ON V_$TRANSACTION TO username;
GRANT SELECT ON V_$YSTREAM_SERVER TO username;

-- YStream permissions
GRANT YSTREAM_CAPTURE TO username;

配置 YStream 服务

YashanDB 连接器需要 YStream 服务来获取增量数据。以下是如何配置 YStream 服务的步骤。

创建 YStream 服务

使用 DBMS_YSTREAM_ADM.CREATE 存储过程创建 YStream 服务:

DBMS_YSTREAM_ADM.CREATE(
    server_name  IN VARCHAR(64),
    connect_user IN VARCHAR(64) DEFAULT NULL,
    start_scn    IN BIGINT DEFAULT NULL);

start_scn 通过查询 SELECT CURRENT_SCN FROM V$DATABASE 获得。

EXEC DBMS_YSTREAM_ADM.CREATE('serverName', 'connect_user', start_scn);

向 YStream 服务添加表

使用 DBMS_YSTREAM_ADM.ADD_TABLES 存储过程向 YStream 服务添加需要捕获的表:

DBMS_YSTREAM_ADM.ADD_TABLES(
    server_name IN VARCHAR(64),
    table_names IN VARCHAR(4096),
    schemas     IN VARCHAR(4096));
添加到服务中的表名和模式必须与连接器配置为捕获的表名和模式一致。

设置 YStream 服务参数

前提条件

  • YStream 服务可供配置。查询 V$YSTREAM_SERVER 视图以获取服务状态。

使用 DBMS_YSTREAM_ADM.SET_PARAMETER 过程 为现有服务设置参数:

DBMS_YSTREAM_ADM.SET_PARAMETER(
    server_name IN VARCHAR(64),
    parameter   IN VARCHAR(64),
    value       IN VARCHAR(64));

启动 YStream 服务

使用 DBMS_YSTREAM_ADM.START 存储过程启动 YStream 服务:

DBMS_YSTREAM_ADM.START(server_name IN VARCHAR(64));
如果需要使用备用数据库连接到 YStream,必须在备用数据库节点上启动 YStream 服务。

部署 YashanDB 连接器

要部署 Debezium YashanDB 连接器,需要安装 Debezium YashanDB 连接器归档文件、配置连接器,并通过将其配置添加到 Kafka Connect 来启动该连接器。

前提条件

操作步骤

  1. 下载 Debezium YashanDB 连接器插件归档文件。
  2. 将所有文件解压到 Kafka Connect 环境中。
  3. 将包含 JAR 文件的目录添加到 Kafka Connect 的 plugin.path 中。
  4. 重启 Kafka Connect 进程,以加载新的 JAR 文件。

后续步骤

Debezium YashanDB 连接器配置

通常,通过提交一个指定连接器配置属性的 JSON 请求来注册 Debezium YashanDB 连接器。以下示例展示了在端口 1688 上注册逻辑名称为 server1 的 Debezium YashanDB 连接器实例的 JSON 请求:

示例:Debezium YashanDB 连接器配置

{
    "name": "inventory-connector",  (1)
    "config": {
        "connector.class" : "io.debezium.connector.yashandb.YashanDbConnector",  (2)
        "database.hostname" : "<YASHANDB_IP_ADDRESS>",  (3)
        "database.port" : "1688",  (4)
        "database.user" : "username",  (5)
        "database.password" : "password",   (6)
        "database.dbname" : "dbname",  (7)
        "topic.prefix" : "server1",  (8)
        "tasks.max" : "1",  (9)
        "database.ystream.server.name" : "server1",  (10)
        "schema.history.internal.kafka.bootstrap.servers" : "kafka:9092", (11)
        "schema.history.internal.kafka.topic": "schema-changes.inventory"  (12)
    }
}
1在向 Kafka Connect 服务注册连接器时为该连接器指定的名称。
2此 YashanDB 连接器类的名称。
3YashanDB 实例的地址。
4YashanDB 实例的端口号。
5YashanDB 用户的名称,参见为连接器创建用户账户。
6YashanDB 用户的密码,参见为连接器创建用户账户。
7要捕获其更改的数据库名称。
8用于标识连接器捕获更改所来自的 YashanDB 数据库服务器并为其提供命名空间的主题前缀。
9为该连接器创建的最大任务数。
10连接器用于捕获更改的 YStream 服务名称。
11该连接器用于将 DDL 语句写入和恢复到数据库模式历史主题的 Kafka 代理列表。
12连接器写入和恢复 DDL 语句所使用的数据库模式历史主题名称。该主题仅供内部使用,不应由使用者使用。

有关可为 Debezium YashanDB 连接器设置的配置属性的完整列表,请参阅连接器配置属性。

你可以通过 POST 命令将此配置发送到正在运行的 Kafka Connect 服务。该服务会记录该配置并启动一个连接器任务,执行以下操作:

  • 连接到 YashanDB 数据库。
  • 读取数据库的重做日志。
  • 为你指定的表中发生的每项操作发出变更事件。
  • 将变更事件记录流式传输到 Kafka 主题。

添加连接器配置

要开始运行 Debezium YashanDB 连接器,请创建连接器配置,并将该配置添加到 Kafka Connect 集群。

前提条件

操作步骤

  1. 为 YashanDB 连接器创建配置。
  2. 使用 Kafka Connect REST API 将该连接器配置添加到 Kafka Connect 集群。

结果

连接器启动后,会对其所配置的 YashanDB 数据库执行一致性快照。随后,连接器开始生成行级操作的数据变更事件,并将变更事件记录流式传输到 Kafka 主题。

连接器属性

Debezium YashanDB 连接器提供了众多配置属性,你可以利用这些属性调整连接器行为以满足应用需求。许多属性都有默认值。属性相关信息组织如下:

Debezium YashanDB 连接器必需的配置属性

Debezium YashanDB 连接器使用众多配置属性来创建连接器实例。这些属性的说明按如下方式组织:

必需的配置属性

属性

默认值

说明

name

无默认值

连接器的唯一名称。指定的名称只能用于注册一次连接器。如果尝试重复使用连接器名称,注册将失败。(所有 Kafka Connect 连接器都需要此属性。)

connector.class

无默认值

连接器的 Java 类名。对于 YashanDB 连接器,请始终使用以下值:

io.debezium.connector.yashandb.YashanDbConnector

tasks.max

1

为该连接器创建的最大任务数。YashanDB 连接器始终使用单个任务,因此不会使用此值,接受默认值即可。

database.hostname

无默认值

YashanDB 数据库服务器的 IP 地址或主机名。

database.port

无默认值

YashanDB 数据库服务器的整数端口号。

database.user

无默认值

连接器用于连接 YashanDB 数据库服务器的 YashanDB 用户账户名。

database.password

无默认值

连接 YashanDB 数据库服务器时使用的密码。

database.dbname

无默认值

YashanDB 数据库的名称。

database.url

无默认值

YashanDB 数据库的 JDBC URL。格式:jdbc:yasdb://<host>:1688/<dbname>。

database.ystream.server.name

无默认值

YashanDB 数据库上 YStream 服务的名称。请指定在“配置 YStream 服务”步骤中创建的 YStream 服务名。

topic.prefix

无默认值

为连接器捕获更改的 YashanDB 数据库服务器提供命名空间的主题前缀。你设置的值将用作连接器发出的所有 Kafka 主题名称的前缀。请指定在整个 Debezium 环境中所有连接器唯一的话题前缀。可用字符包括字母、数字、连字符、点和下划线。

schema.history.internal.kafka.bootstrap.servers

无默认值

连接器用于向数据库架构历史主题写入和恢复 DDL 语句的 Kafka 代理列表。

schema.history.internal.kafka.topic

无默认值

连接器写入和恢复 DDL 语句的数据库架构历史主题名称。此主题仅供内部使用,不应由消费者使用。

可选配置属性

属性

默认值

描述

ystream.blocking.queue.size

128

YashanDB 客户端内置阻塞队列的大小。增量逻辑日志直接从此队列获取。

ystream.poll.timeout

10

从阻塞队列中获取下一个结果的超时时间(秒)。

ystream.client.response.timeout

60

YStream 服务器等待 YStream 客户端响应的最长秒数。

schema.include.list

无默认值

一个可选的、以逗号分隔的正则表达式列表,用于匹配你想要捕获变更的架构(schema)名称。对于名称未包含在 schema.exclude.list 中的任何架构,连接器都会捕获其变更,系统架构除外。默认情况下,将捕获所有非系统架构的变更。为匹配架构名称,Debezium 会将你指定的正则表达式作为锚定正则表达式应用。

schema.exclude.list

无默认值

一个可选的、以逗号分隔的正则表达式列表,用于匹配你不希望捕获变更的架构名称。对于名称未包含在 schema.exclude.list 中的任何架构,连接器都会捕获其变更,系统架构除外。

table.include.list

无默认值

一个可选的、以逗号分隔的正则表达式列表,用于匹配待捕获表的完全限定标识符。设置此属性后,连接器仅从指定的表中捕获变更。每个表标识符使用以下格式:<schema_name>.<table_name>。默认情况下,连接器会监控每个被捕获数据库架构中的所有非系统表。

为匹配表名,Debezium 会将你指定的正则表达式作为锚定正则表达式应用。

table.exclude.list

无默认值

一个可选的、以逗号分隔的正则表达式列表,用于匹配需要从监控中排除的表的完全限定标识符。连接器会从排除列表中未指定的任何表捕获变更事件。请使用以下格式指定每个表标识符:<schema_name>.<table_name>。

为匹配表名,Debezium 会将你指定的正则表达式作为锚定正则表达式应用。

max.batch.size

2048

一个正整数值,用于指定该连接器每次迭代时要处理的每批事件的最大大小。

max.queue.size

8192

一个正整数值,用于指定阻塞队列可容纳的最大记录数。当 Debezium 从数据库读取事件流时,会先将事件放入阻塞队列,然后再将事件写入 Kafka。当连接器接收消息的速度快于写入 Kafka 的速度,或者 Kafka 不可用时,阻塞队列可防止数据丢失。

max.queue.size.in.bytes

0(已禁用)

一个长整型值,用于指定阻塞队列以字节为单位的最大容量。默认情况下,阻塞队列未指定容量限制。若要指定队列可占用的字节数,请将此属性设置为一个正的长整型值。如果同时设置了 max.queue.size,则当队列大小达到任一属性所指定的限制时,对队列的写入将被阻塞。

poll.interval.ms

500(0.5 秒)

一个正整数值,用于指定连接器在每次迭代中等待新变更事件出现的毫秒数。

skipped.operations

t

一个以逗号分隔的运算类型列表,用于指定连接器在流式传输期间要跳过的操作。你可以配置连接器跳过以下类型的运算:c(插入/创建)、u(更新)、d(删除)、t(截断)。默认情况下,仅跳过截断运算。

snapshot.mode

initial

指定连接器对被捕获表执行快照时所使用的模式。你可以设置以下值:

always

每次连接器启动时都执行快照。

initial

在连接器启动时执行初始快照。

initial_only

仅执行初始快照,不流式传输后续的变更。

no_data

仅捕获表结构,不捕获数据。

recovery

恢复丢失或损坏的 schema history 主题。

when_needed

仅在需要时执行快照。

snapshot.fetch.size

10000

指定快照期间每次从每张表中读取的最大行数。连接器会按指定大小分多个批次读取表内容。

snapshot.max.threads

1

指定连接器执行初始快照时使用的线程数。若要启用并行初始快照,请将该属性设置为大于 1 的值。在并行初始快照中,连接器会同时处理多张表。

此功能处于孵化阶段,可能会有所变更。

legacy.snapshot.max.threads

false

指定并行初始快照是否使用旧版的每线程一张表算法。设置为 false(默认值)时,连接器会将源表的内容分块分配给所有并行初始快照线程处理,以获得最佳性能。设置为 true 时,连接器会为每个线程分配一张表进行处理。

snapshot.locking.mode

无默认值

指定连接器在同步快照数据之前是否使用锁定模式来防止 DDL 变更。可设置为以下值之一:

none不获取任何锁。
shared获取共享锁。

lob.enabled

false

指定大对象(CLOB 或 BLOB 等)列值是否在变更事件中输出。默认情况下,变更事件中会包含大对象列,但这些列不包含值。处理和管理大对象列类型及负载会带来一定的开销。若要捕获大对象值并在变更事件中对其进行序列化,请将此选项设置为 true。

decimal.handling.mode

precise

指定连接器应如何处理 NUMBER 列的浮点值。可设置为以下选项之一:

precise

(默认)使用 java.math.BigDecimal 精确表示值,在变更事件中以二进制形式表示。

double

使用 FLOAT64(双精度浮点)表示值。

string

使用 STRING 表示值。

unavailable.value.placeholder

__debezium_unavailable_value

指定连接器提供的常量,用于表示原始值未发生变更且数据库未提供该值。例如,如果 LOB 检索失败,则会使用此占位符代替。

signal.data.collection

无默认值

用于向连接器发送信号的数据集合的完全限定名称。请使用以下格式指定集合名称:<databaseName>.<schemaName>.<tableName>。

signal.enabled.channels

source

为连接器启用的信号通道名称列表。默认情况下,可用通道包括:source、kafka、file、jmx。

»notification.enabled.channels

无默认值

为连接器启用的通知渠道名称列表。默认提供以下渠道:sink、log、jmx。

»incremental.snapshot.chunk.size

1024

在增量快照分块期间,连接器获取并读入内存的最大行数。增大分块大小可以提高效率,因为快照执行的查询次数更少,但每次查询的数据量更大。不过,较大的分块大小也需要更多内存来缓冲快照数据。请根据环境调整分块大小,以获得最佳性能。

»topic.naming.strategy

io.debezium.schema.SchemaTopicNamingStrategy

用于确定数据变更、架构变更、事务、心跳事件等主题名称的类名。默认为 SchemaTopicNamingStrategy。

»topic.delimiter

.

指定主题名称的分隔符。默认为点号(.)。

»converters

无默认值

列出连接器可以使用的自定义转换器实例的符号名称,以逗号分隔。要让连接器使用自定义转换器,必须设置此属性。

对于为连接器配置的每个转换器,还必须添加一个 .type 属性,用于指定实现该转换器接口的类的完全限定名。

»<converter_name>.type

无默认值

为 Debezium 配置自定义转换器的类名。

»<converter_name>.<param_name>

无默认值

自定义转换器的配置,根据转换器的使用方式进行设置。

»query.fetch.size

10000

JDBC 查询的抓取大小。

»ddl.parse.fail.retry.read.table

false

增量 DDL 解析失败后,在处理 DML 事件时完全读取源表结构以进行分析。当 schema.history.internal.skip.unparseable.ddl 和 ddl.parse.fail.retry.read.table 均设置为 true 时,此属性生效。

schema.history.internal

无默认值

处理连接器架构变更的类的名称。对于 YashanDB,使用 io.debezium.relational.history.KafkaSchemaHistory。

schema.history.internal.store.only.captured.tables.ddl

false

指定连接器是记录数据库中所有表的 DDL 语句,还是只记录被捕获表的 DDL 语句。设置为 true 时,仅存储被捕获表的 DDL。

Debezium YashanDB 连接器数据库架构历史配置属性

Debezium 提供了一组 schema.history.internal.* 属性,用于控制连接器与架构历史主题的交互方式。

下表描述了用于配置 Debezium 连接器的 schema.history.internal 属性。

表 13. 连接器数据库架构历史配置属性

属性 默认值 说明

schema.history.internal.kafka.topic

无默认值

连接器存储数据库架构历史的 Kafka 主题的完整名称。

schema.history.internal.kafka.bootstrap.servers

无默认值

连接器用于与 Kafka 集群建立初始连接的主机/端口对列表。此连接用于检索连接器之前存储的数据库架构历史,以及写入从源数据库读取的每条 DDL 语句。每个主机/端口对都应指向 Kafka Connect 进程所使用的同一个 Kafka 集群。

schema.history.internal.kafka.recovery.poll.interval.ms

100

一个整数值,指定连接器在启动/恢复期间轮询持久化数据时应等待的最大毫秒数。默认值为 100 毫秒。

schema.history.internal.kafka.query.timeout.ms

3000

一个整数值,指定连接器使用 Kafka 管理客户端获取集群信息时应等待的最大毫秒数。

schema.history.internal.kafka.create.timeout.ms

30000

一个整数值,指定连接器使用 Kafka 管理客户端创建 Kafka 架构历史主题时应等待的最大毫秒数。

schema.history.internal.kafka.recovery.attempts

100

连接器在恢复失败并报错之前,尝试读取已持久化的历史数据的最大次数。在未接收到数据后最长等待时间为 recovery.attempts × recovery.poll.interval.ms。

schema.history.internal.skip.unparseable.ddl

false

一个布尔值,用于指定连接器是忽略格式错误或未知的数据库语句,还是停止处理以便由人工修复该问题。安全的默认值为 false。跳过解析时应格外谨慎,因为在处理 binlog 时可能导致数据丢失或数据损坏。

schema.history.internal.store.only.captured.tables.ddl

false

一个布尔值,用于指定连接器是记录某个模式或数据库中所有表的模式结构,还是仅记录被指定为捕获目标的表的模式结构。可指定以下值之一:

false(默认值)

在进行数据库快照时,连接器会记录数据库中所有非系统表的模式数据,包括未被指定为捕获目标的表。最好保留默认设置。如果以后决定捕获最初未指定为捕获目标的表的变更,连接器可以轻松开始从这些表中捕获数据,因为它们的模式结构已经存储在模式历史主题中。Debezium 需要表的模式历史,以便识别变更事件发生时该表所处的结构。

true

在进行数据库快照时,连接器仅记录 Debezium 捕获变更事件所来自的表的模式结构。如果你更改了默认值,并且之后配置连接器从数据库中的其他表捕获数据,连接器将缺少从这些表捕获变更事件所需的模式信息。

schema.history.internal.store.only.captured.databases.ddl

false

一个布尔值,用于指定连接器是否记录数据库实例中所有逻辑数据库的模式结构。可指定以下值之一:

true

连接器仅记录 Debezium 捕获变更事件所来自的逻辑数据库和模式中的表的模式结构。

false

连接器记录所有逻辑数据库的模式结构。

schema.history.internal.memory.optimization

off

控制 Debezium 如何通过驻留器(interner)在内存中对相同的模式对象(表、列、属性)进行去重。可指定以下值之一:

off

不执行去重(默认值)。

on

每个连接器使用各自独立的驻留池。在不干扰其他连接器的前提下,减少单个连接器内的堆内存占用。

shared

所有配置为 shared 的连接器共享同一个全局驻留池。当有大量连接器跟踪结构相似的表时,可最大化去重效果。

YashanDB 连接器透传配置属性

连接器支持透传属性,使 Debezium 能够指定自定义配置选项,以便对 Apache Kafka 生产者和消费者的行为进行微调。有关 Kafka 生产者和消费者全部配置属性的信息,请参阅 Kafka 文档。

用于配置生产者和消费者客户端与模式历史主题交互方式的透传属性

Debezium 依赖 Apache Kafka 生产者将模式更改写入数据库模式历史主题。同样,它依赖 Kafka 消费者在连接器启动时从数据库模式历史主题中读取数据。你可以通过为以 schema.history.internal.producer.* 和 schema.history.internal.consumer.* 前缀开头的一组透传配置属性赋值,来定义 Kafka 生产者和消费者客户端的配置。这些透传的生产者和消费者数据库模式历史属性可控制一系列行为,例如这些客户端如何与 Kafka 代理建立安全连接,如以下示例所示:

schema.history.internal.producer.security.protocol=SSL
schema.history.internal.producer.ssl.keystore.location=/var/private/ssl/kafka.server.keystore.jks
schema.history.internal.producer.ssl.keystore.password=test1234
schema.history.internal.producer.ssl.truststore.location=/var/private/ssl/kafka.server.truststore.jks
schema.history.internal.producer.ssl.truststore.password=test1234
schema.history.internal.producer.ssl.key.password=test1234

schema.history.internal.consumer.security.protocol=SSL
schema.history.internal.consumer.ssl.keystore.location=/var/private/ssl/kafka.server.keystore.jks
schema.history.internal.consumer.ssl.keystore.password=test1234
schema.history.internal.consumer.ssl.truststore.location=/var/private/ssl/kafka.server.truststore.jks
schema.history.internal.consumer.ssl.truststore.password=test1234
schema.history.internal.consumer.ssl.key.password=test1234

Debezium 在将属性传递给 Kafka 客户端之前,会先从属性名中移除前缀。

有关 Kafka 生产者配置属性和 Kafka 消费者配置属性的更多信息,请参阅 Apache Kafka 文档。

用于配置 YashanDB 连接器与 Kafka 信号主题交互方式的透传属性

Debezium 提供了一组 signal.* 属性,用于控制连接器与 Kafka 信号主题的交互方式。

下表描述了 Kafka signal 属性。

表 14. Kafka signals 配置属性

属性 默认值 说明

signal.kafka.topic

<topic.prefix>-signal

连接器用于监视临时信号的 Kafka 主题的名称。

如果自动创建主题已禁用,你必须手动创建所需的信号主题。为保证信号顺序,必须使用信号主题。信号主题必须只有一个分区。

signal.kafka.groupId

kafka-signal

Kafka 消费者所使用的组 ID 的名称。

signal.kafka.bootstrap.servers

无默认值

连接器用于与 Kafka 集群建立初始连接的主机与端口对列表。每个对都指向 Debezium Kafka Connect 进程所使用的 Kafka 集群。

signal.kafka.poll.timeout.ms

100

一个整数值,用于指定连接器在轮询信号时等待的最长时间(毫秒)。

用于配置信号通道 Kafka 消费者客户端的透传属性

Debezium 连接器支持对信号 Kafka 消费者进行透传配置。透传的信号属性以 signal.consumer.* 前缀开头。例如,连接器会将 signal.consumer.security.protocol=SSL 这样的属性传递给 Kafka 消费者。

Debezium 在将属性传递给 Kafka 信号消费者之前,会先移除这些属性的前缀。

用于配置 YashanDB 连接器 sink 通知渠道的透传属性

下表描述了可用于配置 Debezium sink notification 渠道的属性。

属性默认值描述
notification.sink.topic.name无默认值接收来自 Debezium 通知的主题名称。当您将 notification.enabled.channels 属性配置为包含 sink 作为启用的通知渠道之一时,此属性为必填项。

表 15. Sink 通知配置属性

Debezium 连接器透传数据库驱动配置属性

Debezium 连接器支持对数据库驱动进行透传配置。透传数据库属性以 driver.* 前缀开头。例如,连接器会将 driver.foobar=false 之类的属性传递给 JDBC URL。

Debezium 在将这些属性传递给数据库驱动之前,会先去除属性中的前缀。

监控

除了 Apache Kafka 和 Kafka Connect 内置的 JMX 指标支持之外,Debezium YashanDB 连接器还提供了三种指标类型。

有关如何使用 JMX 暴露指标的详细信息,请参阅 Debezium 监控文档。

自定义 MBean 名称

Debezium 连接器通过连接器的 MBean 名称暴露指标。这些指标特定于每个连接器实例,提供有关连接器快照、流式传输和模式历史处理过程的行为数据。

默认情况下,当你部署一个配置正确的连接器时,Debezium 会为连接器的每一种不同指标生成一个唯一的 MBean 名称。要查看某个连接器进程的指标,你需要配置可观测性栈来监控其 MBean。但这些默认的 MBean 名称取决于连接器配置;配置的更改可能导致 MBean 名称发生变化。MBean 名称一旦改变,就会破坏连接器实例与 MBean 之间的关联,从而干扰监控活动。在这种情况下,如果你希望恢复监控,就必须重新配置可观测性栈,以使用新的 MBean 名称。

为防止 MBean 名称变更导致监控中断,你可以配置自定义指标标签。通过在连接器配置中添加 custom.metric.tags 属性即可配置自定义指标。该属性接受键值对,其中每个键代表 MBean 对象名称的一个标签,对应的值代表该标签的值。例如:k1=v1,k2=v2。Debezium 会将指定的标签附加到连接器的 MBean 名称之后。

为连接器配置了 custom.metric.tags 属性后,你就可以配置可观测性栈来检索与指定标签关联的指标。随后,可观测性栈将使用指定的标签(而非易变的 MBean 名称)来唯一标识连接器。此后,即使 Debezium 重新定义了构造 MBean 名称的方式,或者连接器配置中的 topic.prefix 发生变化,指标采集也不会中断,因为指标抓取任务使用指定的标签模式来标识连接器。

使用自定义标签的另一个好处是,你可以使用能够反映数据管道架构的标签,从而以适合自身运维需求的方式来组织指标。例如,你可以指定标签,其值用于声明连接器活动的类型、应用上下文或数据源,例如 db1-streaming-for-application-abc。如果你指定了多个键值对,所有指定的键值对都会被附加到连接器的 MBean 名称中。

下面的示例说明了标签如何修改默认的 MBean 名称。

示例 1. 自定义标签如何修改连接器的 MBean 名称

默认情况下,YashanDB 连接器对流式传输指标使用以下 MBean 名称:

debezium.yashandb:type=connector-metrics,context=streaming,server=<topic.prefix>

如果你将 custom.metric.tags 的值设置为 database=salesdb-streaming,table=inventory,Debezium 会生成如下的自定义 MBean 名称:

debezium.yashandb:type=connector-metrics,context=streaming,server=<topic.prefix>,database=salesdb-streaming,table=inventory

快照指标

MBean 为 debezium.yashandb:type=connector-metrics,context=snapshot,server=<topic.prefix>。

下表列出了可用于监控 Debezium 快照操作的 JMX 指标,包括行计数、表进度、耗时和队列容量。除非有快照操作正在进行,或者自上次连接器启动以来已发生过快照,否则不会暴露快照指标。

属性类型描述
LastEventstring连接器读取到的最后一个快照事件。
MilliSecondsSinceLastEventlong自连接器读取并处理最近一个事件以来经过的毫秒数。
NumberOfErroneousEventslong记录连接器在快照操作期间识别为错误的变更事件数量。每当连接器在初始快照、增量快照或临时快照过程中遇到无法处理的事件时,该指标值就会递增。事件处理失败的原因可能包括:格式错误、与模式不兼容,或在转换过程中出现故障。该指标值在连接器任务的整个生命周期内保持不变。如果快照被中断且连接器任务重新启动,该指标计数将重置为 0。
TotalNumberOfEventsSeenlong自上次启动或重置以来,此连接器看到的事件总数。
NumberOfEventsFilteredlong被连接器上配置的包含/排除列表过滤规则所过滤的事件数量。
CapturedTablesstring[]连接器捕获的表列表。
QueueTotalCapacityint用于在快照程序和主 Kafka Connect 循环之间传递事件的队列长度。
QueueRemainingCapacityint用于在快照程序和主 Kafka Connect 循环之间传递事件的队列剩余容量。
TotalTableCountint纳入快照的表的总数。
RemainingTableCountint快照尚未复制的表的数量。
SnapshotRunningboolean快照是否已启动。
SnapshotPausedboolean快照是否已暂停。
SnapshotAbortedboolean快照是否已中止。
SnapshotCompletedboolean快照是否已完成。
SnapshotSkippedboolean快照是否已被跳过。
SnapshotDurationInSecondslong快照迄今为止所花费的总秒数,即使快照尚未完成。该值也包含快照处于暂停状态的时间。
SnapshotPausedDurationInSecondslong快照被暂停的总秒数。如果快照多次暂停,暂停时间将累计相加。
RowsScannedMap<String, Long>包含快照中每张表已扫描行数的映射。处理过程中会逐步向该映射中添加表。每扫描 10,000 行以及完成一张表时更新一次。
TableChunkCountsMap<String, Long>使用基于分块的多线程快照时,包含快照中每张表的分块数量的映射。
TableChunksCompletedCountsMap<String, Long>使用基于分块的多线程快照时,包含快照中每张表已完成分块数量的映射。
MaxQueueSizeInByteslong队列缓冲区的最大字节数。当 max.queue.size.in.bytes 设置为正的 long 值时,此指标可用。
CurrentQueueSizeInByteslong队列中记录的当前字节大小。

下表列出了连接器执行增量快照时可用的附加 JMX 指标,其中包括可用于跟踪快照进度的分块标识和表边界标识。

属性类型描述
ChunkIdstring当前快照分块的标识符。
ChunkFromstring定义当前分块的主键集合的下界。
ChunkTostring定义当前分块的主键集合的上界。
TableFromstring当前正在快照的表的主键集合的下界。
TableTostring当前正在快照的表的主键集合的上界。

流处理指标

MBean 为 debezium.yashandb:type=connector-metrics,context=streaming,server=<topic.prefix>。

通用流处理指标

下表列出了可用于监控 Debezium 流处理操作的 JMX 指标,包括按类型统计的事件数量、相对源头的延迟、队列容量以及连接状态。

属性类型描述
LastEventstring连接器读取到的最后一个流式事件。
MilliSecondsSinceLastEventlong距离连接器读取并处理最近一个事件所经过的毫秒数。
NumberOfErroneousEventslong记录连接器在流式处理过程中识别为错误的变更事件数量。在流式会话的整个生命周期内,连接器每遇到一个无法处理的事件,该指标就会递增。事件可能因格式错误、与模式不兼容,或在转换过程中失败而导致处理失败。该指标值在连接器任务的生命周期内持续保留。连接器重启后,该指标计数会重置为 0。
TotalNumberOfEventsSeenlong自上次启动连接器或重置指标以来,源数据库报告的数据变更事件总数。代表 Debezium 需要处理的数据变更工作负载。
TotalNumberOfCreateEventsSeenlong自上次启动或重置指标以来,连接器处理的创建事件总数。
TotalNumberOfUpdateEventsSeenlong自上次启动或重置指标以来,连接器处理的更新事件总数。
TotalNumberOfDeleteEventsSeenlong自上次启动或重置指标以来,连接器处理的删除事件总数。
NumberOfEventsFilteredlong被连接器上配置的包含/排除列表过滤规则过滤掉的事件数量。
NumberOfUnchangedEventsSkippedlong自上次启动连接器或重置指标以来,因被监视列未发生变化而跳过的更新事件数量。如果 skip.messages.without.change 为 false,默认值为 -1。对于 YashanDB 以及其他不支持跳过未变更事件的连接器,即使该属性设置为 true,该值也可能始终为 0。
CapturedTablesstring[]连接器捕获的表列表。
QueueTotalCapacityint用于在流式处理组件与 Kafka Connect 主循环之间传递事件的队列长度。
QueueRemainingCapacityint用于在流式处理组件与 Kafka Connect 主循环之间传递事件的队列剩余容量。
Connectedboolean指示连接器当前是否已连接到数据库服务器的标志。
MilliSecondsBehindSourcelong最后一个变更事件的时间戳与连接器处理它之间相差的毫秒数。该值会包含数据库服务器和连接器所在机器之间时钟的任何差异。
MilliSecondsBehindSourceMinValuelong连接器运行期间观察到的落后于源的最小延迟(毫秒)。
MilliSecondsBehindSourceMaxValuelong连接器运行期间观察到的落后于源的最大延迟(毫秒)。
MilliSecondsBehindSourceAverageValuedouble连接器运行期间所有观测值计算得出的落后于源的平均延迟(毫秒)。
MilliSecondsBehindSourceP50double落后于源延迟的第 50 百分位数(中位数)。与平均值相比,该指标受离群值影响较小,能更稳健地衡量典型延迟。当 statistics.metrics.enabled 设置为 true(默认)时可用。
MilliSecondsBehindSourceP95double落后于源延迟的第 95 百分位数。该指标表明 95% 的延迟测量值低于此值,可用于识别尾部延迟和设定 SLA 阈值。当 statistics.metrics.enabled 设置为 true(默认)时可用。
MilliSecondsBehindSourceP99double落后于源延迟的第 99 百分位数。该指标表明 99% 的延迟测量值低于此值,可用于了解最坏情况下的性能表现。当 statistics.metrics.enabled 设置为 true(默认)时可用。

| [NumberOfCommittedTransactions](

YStream 流式指标

Debezium YashanDB 连接器提供以下 YStream 适配器特有的附加流式指标:

属性类型描述
ErrorCountlong连接器在流式传输阶段检测到的错误数量。
WarningCountlong连接器在流式传输阶段检测到的警告数量。

表 16. YStream 特有流式指标说明

架构历史指标

MBean 为 debezium.yashandb:type=connector-metrics,context=schema-history,server=<topic.prefix>。

下表列出了可用于监控连接器架构历史过程的 JMX 指标,包括恢复状态、已应用的架构变更数量以及最近一次变更的时间戳。

属性类型描述
Statusstring数据库模式历史的状态,取值为 STOPPED、RECOVERING(正在从存储中恢复历史记录)或 RUNNING。
RecoveryStartTimelong恢复开始的时间,以纪元秒(epoch seconds)表示。
ChangesRecoveredlong恢复阶段读取到的变更数量。
ChangesAppliedlong恢复和运行期间应用的模式变更总数。
MilliSecondsSinceLast​RecoveredChangelong自从上一次从历史存储中恢复变更以来经过的毫秒数。
MilliSecondsSinceLast​AppliedChangelong自从上一次应用变更以来经过的毫秒数。
LastRecoveredChangestring从历史存储中恢复的最新一条变更的字符串表示。
LastAppliedChangestring最近一次应用的变更的字符串表示。

常见问题

错误:YashanDB 尚无 YStream 服务器 'serverxx',请检查参数是否填写正确(参数为 'database.ystream.server.name')。

database.ystream.server.name 参数所对应的 YStream 服务器在 YashanDB 数据库中不存在。请按照配置 YStream 服务中的说明创建相关的 YStream 服务器。

错误:YashanDB YStream 服务器状态为 xxx。请执行 'DBMS_YSTREAM_ADM.START(…​)' 以启动 YStream 服务器。

database.ystream.server.name 参数所对应的 YStream 服务器未处于运行状态。请在数据库中执行以下命令以启动 YStream 服务:

EXEC DBMS_YSTREAM_ADM.START('server1');

Decimal 值同步到 Kafka 后,为什么序列化数据不正确?

Debezium 对负精度的 Decimal 进行了特殊处理。若要规避此问题,请使用参数 decimal.handling.mode=string。

DATE/TIME/TIMESTAMP 值同步到 Kafka 后,为什么是时间戳形式,而不是 'yyyy-MM-dd HHss.SSSSSS' 形式?

默认情况下,Debezium 将时间类型映射为 INT64。有关如何将时间类型映射为固定格式字符串,请参阅数据类型转换章节,配置自定义转换器。

YashanDB 连接器是否支持从检查点恢复?任务停止或失败后,能否从上次提交的位置继续捕获增量数据?

支持。YashanDB 连接器支持从检查点恢复。基于 Kafka 的两阶段提交机制,任务停止或失败后,最后一次成功提交的日志位置会记录在 Kafka 的元数据中。任务恢复时,连接器会获取最后一次提交的日志位置,并从该位置开始捕获数据,从而确保向 Kafka 主题的数据同步实现精确一次(exactly-once)语义。

YashanDB 连接器会在任务日志中捕获数据库中所有表的元数据结构。能否让连接器只捕获已配置表(schema.include.list 和 table.include.list)的元数据结构?

默认情况下,Debezium 会捕获数据库中所有表的结构。您可以在任务配置中设置 schema.history.internal.store.only.captured.tables.ddl=true,以仅捕获已配置表的结构。

删除并重建任务后,重启时未捕获表的元数据结构(例如日志中没有出现 "Capturing structure of table"),这是什么原因造成的?

连接器将表的元数据存储在 schema.history.internal.kafka.topic 属性指定的 Kafka 主题中。当您删除并以相同名称重建任务时,只要 schema history 主题未被删除,连接器就会从 schema history 主题中读取已有的元数据,而不是重新捕获表结构的快照。若要再次捕获表的元数据,可以删除 schema history 主题,或更改任务名称。

文档中说明不支持自定义数据类型、XMLTYPE 和 JSON 数据类型,但 XMLTYPE 和 JSON 数据仍然可以同步到 Kafka 主题,这是为什么?

YashanDB 连接器的类型支持取决于 YashanDB YStream 的支持范围。XMLTYPE 和 JSON 数据可以正常同步,但数据的正确性无法保证,具体取决于 YStream 的支持情况。

YashanDB 源端的 DDL 导致任务在解析该 DDL 时失败,如何跳过这个导致失败的 DDL?

Debezium 提供了跳过 DDL 解析失败的功能。你可以通过修改任务配置,将 schema.history.internal.skip.unparseable.ddl=true 设置为 true 来启用该功能。

YashanDB 源端执行了类似 "CREATE TABLE …​ AS SELECT" 的 DDL,随后该表的 DML 数据解析失败,应如何处理?

这类 DDL 无法被解析以获取表的元数据,从而导致后续的 DML 数据解析失败。Debezium Oracle 连接器也存在同样的问题。如果你遇到此问题,请重新启动一个新的同步任务。

通过存储过程包(如 DBMS_STATS.CREATE_STAT_TABLE)创建的增量 DDL 无法被识别为具体的 DDL 语句,原因是什么?应如何处理?

连接器无法将通过存储过程包创建的 DDL 语句识别为具体的 DDL 语句。因此,在数据同步过程中可能会出现元数据不匹配,进而导致同步错误。为避免此问题,请直接执行 DDL 语句,而不要通过存储过程包来执行。

评论

登录后参与评论

正在加载评论…