Flink CDC自定义函数
本文将向您介绍编写和使用用户自定义函数(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方法来手动指定方法的返回类型。重写
open和close方法来插入生命周期函数。
例如,将传入的整型参数增加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会被转换为 PythonNone。Python 函数可以返回
None,此时输出为 SQLNULL。
例如,下面的 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 列类型 |
|
|
|
|
|
|
|
|
|
|
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 支持以下配置项。
参数 | 说明 | 是否必填 | 数据类型 | 默认值 | 备注 |
| 注册到 Transform 模块中的 UDF 名称。 | 是 | STRING | 无 | 同一个作业中的 UDF 名称应保持唯一。 |
| Python UDF 的内联源代码。 | 是 | STRING | 无 | 必须包含一个名为 |
| Python UDF 使用的附加依赖。 | 否 | LIST<STRING> | 无 | 若Python UDF需要额外的附加依赖,则需要将其打包为 .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) < 100Flink SQL UDF兼容性
继承自ScalarFunction的Flink UDF函数也可以直接作为CDC YAML UDF函数注册并使用,但存在以下限制:
不支持带参数的
ScalarFunction。Flink风格的TypeInformation类型标注会被忽略。
open和close生命周期钩子函数不会被调用。