应用程序接口与服务提供者接口

自定义转换器

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

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

自定义转换器

数据类型转换

默认情况下,Debezium 连接器使用内置的类型映射将源数据类型映射为 Kafka Connect 模式类型。当下游应用需要不同的数据格式时,你可以通过开发并部署自定义转换器来覆盖这些默认设置。

Debezium 变更事件记录中的每个字段都对应源表或源数据集合中的一个字段或列。当连接器向 Kafka 发出变更事件记录时,会将源中每个字段的数据类型转换为 Kafka Connect 模式类型。列值同样会被转换,以匹配目标字段的模式类型。每个连接器都有一套默认映射,用于指定该连接器如何转换每种数据类型。这些默认映射在各连接器的数据类型文档中均有说明。

尽管默认映射通常已足够使用,但在某些应用中你可能需要采用其他映射方式。例如,如果默认映射以 UNIX 纪元以来的毫秒数格式导出某一列,而你的下游应用只能以格式化字符串的形式消费该列的值,那么你就需要自定义映射。自定义数据类型映射的方式是开发并部署一个自定义转换器。你可以配置自定义转换器使其作用于某一类型的所有列,也可以缩小其作用范围,使其仅适用于特定的表列。转换函数会拦截所有符合指定条件的列的数据类型转换请求,并执行指定的转换。不符合指定条件的列则会被转换器忽略。

自定义转换器是实现了 Debezium 服务提供者接口(SPI)的 Java 类。通过在连接器配置中设置 converters 属性来启用并配置自定义转换器。converters 属性指定连接器可用的转换器,并可包含用于进一步调整转换行为的子属性。

启动连接器后,连接器配置中已启用的转换器会被实例化并添加到注册表中。注册表会将每个转换器与其需要处理的列或字段关联起来。每当 Debezium 处理新的变更事件时,就会调用已配置的转换器,来转换与其注册相关联的列或字段。

以下说明仅适用于 Debezium 关系型数据库源连接器。你不能使用这些信息为 Debezium MongoDB 连接器或 Debezium JDBC 接收器连接器创建自定义转换器。

实现自定义转换器

下面的示例展示了一个 Java 类的转换器实现,该类实现了接口 io.debezium.spi.converter.CustomConverter:

public interface CustomConverter<S, F extends ConvertedField> {

    @FunctionalInterface
    interface Converter {
        Object convert(Object input);
    }

    public interface ConverterRegistration<S> {
        void register(S fieldSchema, Converter converter);
    }

    void configure(Properties props);

    void converterFor(F field, ConverterRegistration<S> registration);
}

以下列表说明了前面的 Java 转换器类示例中各选择元素的用途:

Converter

将数据从一种类型转换为另一种类型。

ConverterRegistration

用于注册转换器的回调。

register()

为当前字段注册给定的模式(schema)和转换器。对于同一个字段,不应多次调用此方法。

converterFor()

注册自定义的值转换器和模式转换器,以便用于特定字段。

自定义转换器方法

CustomConverter 接口的实现必须包含以下方法:

configure()

将连接器配置中指定的属性传递给转换器实例。configure 方法在连接器初始化时运行。你可以在多个连接器中使用同一个转换器,并根据连接器的属性设置来调整其行为。configure 方法接受以下参数:

props

包含要传递给转换器实例的属性。每个属性都指定了针对特定类型列的值进行转换的格式。

converterFor()

注册转换器,以处理数据源中的特定列或字段。Debezium 调用 converterFor() 方法,提示转换器调用 registration 来执行转换。converterFor 方法对每一列运行一次。该方法接受以下参数:

field

一个对象,用于传递有关所处理字段或列的元数据。列元数据可以包括列或字段的名称、表或集合的名称、数据类型、大小等等。

registration

类型为 io.debezium.spi.converter.CustomConverter.ConverterRegistration 的对象,它提供目标模式定义以及用于转换列数据的代码。当源列的类型与转换器应处理的类型匹配时,转换器会调用 registration 参数,并调用 register 方法来为模式中的每一列定义转换器。模式使用 Kafka Connect 的 SchemaBuilder API 表示。将来会新增一个独立的模式定义 API。

Debezium 自定义转换器示例

以下示例展示了如何实现一个转换器,将类型为 isbn 的源列映射为 Kafka Connect 的 STRING 模式值。你可以以此示例为模板,编写用于处理源数据库中自定义或非标准列类型的转换器。

示例中的转换器执行以下操作:

  • 运行 configure 方法,根据连接器配置中指定的 schema.name 属性值来配置转换器。转换器的配置是针对每个实例的。

  • 运行 converterFor 方法,注册该转换器,以处理数据类型设置为 isbn 的源列中的值。

  • 根据为 schema.name 属性指定的值,识别目标 STRING 模式。

    • 将源列中的 ISBN 数据转换为 String 值。

示例 1. 一个简单的自定义转换器

public static class IsbnConverter implements CustomConverter<SchemaBuilder, RelationalColumn> {

    private String isbnSchemaName;

    @Override
    public void configure(Properties props) {
        isbnSchemaName = props.getProperty("schema.name");
    }

    @Override
    public void converterFor(RelationalColumn column,
            ConverterRegistration<SchemaBuilder> registration) {

        if ("isbn".equals(column.typeName())) {
            registration.register(SchemaBuilder.string().name(isbnSchemaName), x -> x.toString());
        }
    }
}

Debezium 与 Kafka Connect API 模块依赖

自定义转换器的 Java 项目在编译时依赖 Debezium API 和 Kafka Connect API 库模块。这些编译依赖必须包含在项目的 pom.xml 中,示例如下:

<dependency>
    <groupId>io.debezium</groupId>
    <artifactId>debezium-api</artifactId>
    <version>${version.debezium}</version>
</dependency>
<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>connect-api</artifactId>
    <version>${version.kafka}</version>
</dependency>

以下列表说明了上文 pom.xml 示例中编译依赖版本的用途:

${version.debezium}

表示 Debezium 连接器的版本。

${version.kafka}

表示你环境中 Apache Kafka 的版本。

配置和使用转换器

自定义转换器针对源表中的特定列或列类型进行处理,用于指定如何将源中的数据类型转换为 Kafka Connect 模式类型。要将自定义转换器与连接器配合使用,需要将转换器 JAR 文件与连接器文件一起部署,然后配置连接器以使用该转换器。

自定义转换器旨在修改 Debezium 关系型数据库源连接器发出的消息。你无法为 Debezium MongoDB 连接器或 Debezium JDBC 接收器连接器配置自定义转换器。

部署自定义转换器

要将自定义转换器与 Debezium 连接器配合使用,请将转换器导出为 JAR 文件,并将其放入连接器的插件目录中。

前提条件

  • 你已经编写了自定义转换器 Java 程序。

操作步骤

  • 要将自定义转换器与 Debezium 连接器配合使用,请将 Java 项目导出为 JAR 文件,然后将该文件复制到包含你想要配合使用的每个 Debezium 连接器 JAR 文件的目录中。

    例如,在典型的部署中,Debezium 连接器文件存储在 Kafka Connect 目录(/kafka/connect)的子目录中,每个连接器的 JAR 文件都位于各自的子目录中(/kafka/connect/debezium-connector-db2、/kafka/connect/debezium-connector-mysql 等)。要将转换器与某个连接器配合使用,请将转换器 JAR 文件添加到该连接器的子目录中。

要将转换器与多个连接器配合使用,你必须在每个连接器子目录中都放置一份转换器 JAR 文件的副本。

配置连接器以使用自定义转换器

要使连接器能够使用自定义转换器,请在 Debezium 源连接器的配置中添加属性,以指定转换器的名称和类。你无法为 Debezium JDBC 接收器连接器配置自定义转换器。如果转换器需要更多信息来定制特定数据类型的格式,你可以定义其他配置选项来提供这些信息。

前提条件

操作步骤

  • 通过向连接器配置中添加以下必需属性,为连接器实例启用转换器:

    converters: <converterSymbolicName>
    <converterSymbolicName>.type: <fullyQualifiedConverterClassName>

以下列表说明了前一个示例中的属性:

converters

必需的 converters 属性用于枚举一个逗号分隔的符号名列表,指定要与连接器一起使用的转换器实例。该属性列出的值将作为你为转换器指定的其他属性名称的前缀。

<converterSymbolicName>.type

必需的 <converterSymbolicName>.type 属性用于指定实现该转换器的类的名称。

例如,针对前文的自定义转换器示例,你需要向连接器配置中添加以下属性:

converters: isbn
isbn.type: io.debezium.test.IsbnConverter
  • 要将其他属性与自定义转换器关联起来,请在属性名前加上转换器的符号名和一个点号(.)。该符号名是你作为 converters 属性的值所指定的标签。例如,要为前文的 isbn 转换器添加一个属性,用于指定传递给转换器代码中 configure 方法的 schema.name,请添加以下属性:

    isbn.schema.name: io.debezium.postgresql.type.Isbn

评论

登录后参与评论

正在加载评论…