Flink CDC自定义函数

更新时间: 2026-08-17 15:14:04

本文将向您介绍编写和使用用户自定义函数(UDF)的方法。

自定义函数(UDF)

如果内置的函数不能满足您的需求,Flink CDC数据摄入作业也支持使用Java语言编写自定义UDF函数,并像内置函数一样调用。

说明

此处类路径对应的JAR包需要在“更多配置”中作为附加依赖文件上传。

Java UDF 定义

如果您希望使用 Java 语言编写自定义函数,您需要首先创建一个 Maven 项目,并在pom.xml中引入基础依赖包:

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-cdc-common</artifactId>
    <version>${CDC 社区版本号}</version>
    <scope>provided</scope>
</dependency>

请参考下表选定使用的社区版本号:

实时计算 VVR 引擎版本

对应的社区版本号

11.7 及更高版本

3.6.0

11.3 至 11.6 版本

3.5.0

11.0 至 11.2 版本

3.4.0

8.0.11 版本

3.3.0

8.0.10 及更早版本

3.2.1

满足以下要求的Java类可以作为Flink CDC数据摄入作业UDF函数使用:

  • 实现了org.apache.flink.cdc.common.udf.UserDefinedFunction接口。

  • 拥有一个公共无参构造器。

  • 至少含有一个名为eval的公共方法。

UDF函数类可以通过@Override以下接口来实现更精细的语义控制:

  • 重写getReturnType方法来手动指定方法的返回类型。

  • 重写openclose方法来插入生命周期函数。

例如,将传入的整型参数增加1后返回的UDF函数定义如下。

public class AddOneFunctionClass implements UserDefinedFunction {
    
    public Object eval(Integer num) {
        return num + 1;
    }
    
    @Override
    public DataType getReturnType() {
        // 由于eval函数的返回类型不明确,需要
        // 使用getReturnType写明确指定类型
        return DataTypes.INT();
    }
    
    @Override
    public void open() throws Exception {
        // ...
    }

    @Override
    public void close() throws Exception {
        // ...
    }
}

类型映射关系

下表列出了UDF函数类中eval方法参数及返回值类型与getReturnType方法返回值之间的映射关系。

CDC列类型

对应的Java类

备注

BOOLEAN

java.lang.Boolean

TINYINT

java.lang.Byte

SMALLINT

java.lang.Short

INTEGER

java.lang.Integer

BIGINT

java.lang.Long

FLOAT

java.lang.Float

DOUBLE

java.lang.Double

DECIMAL

java.math.BigDecimal

DATE

java.time.LocalDate

TIME

java.time.LocalTime

TIMESTAMP

java.time.LocalDateTime

TIMESTAMP_TZ

java.time.ZonedDateTime

TIMESTAMP_LTZ

java.time.Instant

CHAR

VARCHAR

STRING

java.lang.String

BINARY

VARBINARY

BYTES

byte[]

ARRAY

java.util.List

ARRAY内部元素类型映射到List类的泛型参数。

MAP

java.util.Map

MAP键、值类型映射到Map类的泛型参数。

ROW

java.util.List

VARIANT

org.apache.flink.cdc.common.types.variant.Variant

说明

请注意和Flink SQL的Variant类路径不同。

Java UDF 注册

通过在Flink CDC数据摄入作业的pipeline块中加入如下所示的定义即可注册UDF函数:

pipeline:
  user-defined-function:
    - name: inc
      classpath: org.apache.flink.cdc.udf.examples.java.AddOneFunctionClass
    - name: format
      classpath: org.apache.flink.cdc.udf.examples.java.FormatFunctionClass
说明
  • 此处类路径对应的JAR包,需要在更多配置中作为附加依赖文件上传。

  • UDF函数名称可以在此处任意调整,无需与UDF类名一致。

Python UDF(公测中)

您可以直接在 Flink CDC 数据摄入 YAML 作业中使用 Python 语言编写内联自定义函数,无需创建 Maven 项目或额外编译 JAR 包。

说明

您需要升级到实时计算引擎 VVR 11.9(含 Preview 预览版本)及更高版本,方可使用 Python UDF 定义自定义函数的能力。

这是一项实验性功能。

语法规则

Python UDF 需要满足以下要求:

  • 必须在 python-code 中定义一个名为 eval 的顶层函数。

  • eval 函数必须声明返回类型注解。

  • eval 函数的参数数量必须与调用 UDF 时传入的参数数量一致。

  • 参数类型注解不是必需的,但建议保留,以提高代码可读性。

  • Python UDF 当前仅支持标量函数,即每次调用返回一个值。

  • 输入数据中的 SQL NULL 会被转换为 Python None

  • Python 函数可以返回 None,此时输出为 SQL NULL

例如,下面的 Python UDF 会对输入字符串执行去除首尾空格和转换为小写的操作:

def eval(value: str) -> str:
    return value.strip().lower()

您可以在 python-code 中定义辅助函数、导入模块或者声明常量,但用于执行数据转换的 eval 函数必须位于代码的顶层。例如:

import re

EMAIL_PATTERN = re.compile(r"\s+")

def normalize(value):
    return EMAIL_PATTERN.sub("", value).lower()

def eval(value: str) -> str:
    if value is None: return None
    return normalize(value)
说明

您可以直接在代码块顶层 import 内置的 Python 包,详细列表请参见预装软件包列表

类型映射

Flink CDC 根据 eval 函数的返回类型注解(Type Annotation)确定 UDF 的输出列类型。目前,只支持下列参数类型及返回值类型:

Python 返回类型注解

CDC 列类型

bool

BOOLEAN

int

BIGINT

float

DOUBLE

str

STRING

bytes

BYTES

Python UDF 注册

通过在 Flink CDC 数据摄入作业的 pipeline.user-defined-function 中配置 python-code 参数,即可注册 Python UDF。使用 YAML 的 | 多行字符串语法,可以自然地书写带缩进的 Python 代码块:

pipeline:
  user-defined-function:
    - name: normalize_email
      python-code: |
        def eval(value: str) -> str:
            if value is None:
                return None
            return value.strip().lower()

    - name: double_age
      python-code: |
        def eval(value: int) -> int:
            if value is None:
                return None
            return value * 2

其中,name 是在 Transform 表达式中使用的函数名称,无需与 Python 函数名一致。内联代码中供 Flink CDC 调用的函数名始终为 eval

Python UDF 空值及异常处理

数据摄入 Pipeline 中的 NULL 值和 Python 的 None 空值对应。若 UDF 被使用 NULL 值作为参数调用,Python eval 函数将收到 None;若 Python eval 函数返回 None,则 NULL 将用作 UDF 的返回值。

若 Python eval 函数抛出运行时异常,Transform 算子将触发 Failover。若您的 UDF 存在失败可能,请使用 Python 语言的 try + catch 语句捕捉异常,确保作业稳定运行。

Python UDF 参数配置

Python UDF 支持以下配置项。

参数

说明

是否必填

数据类型

默认值

备注

name

注册到 Transform 模块中的 UDF 名称。

STRING

同一个作业中的 UDF 名称应保持唯一。

python-code

Python UDF 的内联源代码。

STRING

必须包含一个名为 eval 的顶层函数,并为其声明受支持的返回类型注解。

python-files

Python UDF 使用的附加依赖。

LIST<STRING>

若Python UDF需要额外的附加依赖,则需要将其打包为 .zip 文件后作为附加依赖上传。

说明

附加依赖文件会被放置到 /flink/usrlib/目录中,请输入包含前缀的完整路径。

例如,如果您的附加依赖文件为 PIL.zip,则需要配置的完整路径为 /flink/usrlib/PIL.zip

UDF 使用

在完成 UDF 函数注册后,即可在 Transform 块中,像内置函数一样直接调用 UDF 函数。代码示例如下:

transform:
  - source-table: db.\.*
    projection: "*, inc(inc(inc(id))) as inc_id, format(id, 'id -> %d') as formatted_id"
    filter: inc(id) < 100

Flink SQL UDF兼容性

继承自ScalarFunction的Flink UDF函数也可以直接作为CDC YAML UDF函数注册并使用,但存在以下限制:

  • 不支持带参数的ScalarFunction

  • Flink风格的TypeInformation类型标注会被忽略。

  • openclose生命周期钩子函数不会被调用。

上一篇: Flink CDC内置函数 下一篇: 脏数据收集
阿里云首页 实时计算 Flink版 相关技术圈