本文为您介绍如何为Flink自定义表值函数(UDTF)开发、注册和使用流程。
定义
自定义表值函数(UDTF),自定义表值函数,将0个、1个或多个标量值作为输入参数(可以是变长参数)。与自定义的标量函数类似,但与标量函数不同。表值函数可以返回任意数量的行作为输出,而不仅是1个值。返回的行可以由1个或多个列组成。调用一次函数输出多行或多列数据。详情参见User-defined Functions。
UDTF开发
Flink为您提供了UDF示例,便于您快速开发业务。Flink UDF示例中包含UDSF、UDAF和UDTF的实现,示例中已为您配置对应版本的开发环境,您无需进行环境搭建。
下载并解压ASI_UDX_Demo示例到本地。
说明ASI_UDX_Demo属于第三方搭建的网站,访问时可能会存在无法打开或访问延迟的问题。
解压完成后,会生成ASI_UDX-main文件夹。其中:
pom.xml:项目级别的配置文件,主要描述了项目的Maven坐标、依赖关系,开发者需要遵循的规则、缺陷管理系统,组织和Licenses,以及其他所有的项目相关因素。
\ASI_UDX-main\src\main\java\ASI_UDF\ASI_UDTF.java:自定义表值函数(UDTF)示例的Java代码。
在IntelliJ IDEA中,单击,打开刚才解压缩完成的ASI_UDX-main。
双击打开\ASI_UDX-main\src\main\java\ASI_UDTF后,根据您的业务,修改ASI_UDTF.java文件内容。
该示例中,ASI_UDTF.java已配置了将一行字符串按照竖线(|)分割成多列字符串的代码。
package ASI_UDTF; import org.apache.flink.api.java.tuple.Tuple2; import org.apache.flink.table.functions.TableFunction; public class ASI_UDTF extends TableFunction<Tuple2<String,String>> { public void eval(String str){ String[] split = str.split("\\|"); String name = split[0]; String place = split[1]; Tuple2<String,String> tuple2 = Tuple2.of(name,place); collect(tuple2); } }TableFunction支持的数据类型及类型推导机制请参见Flink支持的数据类型和类型推导。
说明以上两个文档链接为Flink 1.15版本对应的文档,不同Flink大版本中TableFunction支持的数据类型及推导机制可能会存在差异,请您通过VVR和Flink版本的映射关系去参考对应版本的Flink文档。Flink版本查看方法请参见空间管理与操作。
常用的复合类型Tuple和Row示例如下:
Tuple类型
TableFunction<Tuple2<String,Integer>Row类型
@FunctionHint(output = @DataTypeHint("ROW<word STRING, length INT>")) public static class SplitFunction extends TableFunction<Row> { public void eval(String str) { for (String s : str.split(" ")) { // use collect(...) to emit a row collect(Row.of(s, s.length())); } } }
双击打开\ASI_UDX-main\后,配置pom.xml。
该示例中,pom.xml文件已配置了Flink 1.11版依赖的主要JAR包信息。如果您的业务:
不依赖其他JAR包:不用配置pom.xml文件,继续下一步。
依赖其他JAR包:在pom.xml文件中添加您所需依赖的JAR包信息。
Flink 1.11版依赖的主要JAR包如下。
<dependencies> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java_2.12</artifactId> <version>1.11.0</version> <!--<scope>provided</scope>--> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-table</artifactId> <version>1.11.0</version> <type>pom</type> <!--<scope>provided</scope>--> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-core</artifactId> <version>1.11.0</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-table-common</artifactId> <version>1.11.0</version> </dependency> </dependencies>在下载文件中pom.xml所在的目录执行如下命令打包文件。
mvn package -Dcheckstyle.skip\ASI_UDX-main\target\目录下会出现ASI_UDX-1.0-SNAPSHOT.jar的JAR包,即代表完成了UDTF开发工作。
自定义异步表值函数(AsyncTableFunction)
异步表值函数(AsyncTableFunction)用于在 Flink SQL 或 Table API 中异步访问外部服务,每次调用返回零行、一行或多行结果,属于自定义表值函数(UDTF)。典型场景如根据商品 ID 查询多个标签,或根据检索条件返回多条匹配记录。异步表值函数既可以作为用户自定义函数,通过 SQL 的 LATERAL TABLE 或 Table API 的 joinLateral 调用,形成异步 Correlate 运算;也可以由维表连接器提供,作为异步 Lookup Join 的查询实现。
配置参数
通过以下参数控制异步表值函数的运行行为。
配置项 | 适用场景 | 作用 |
table.exec.async-table.max-concurrent-operations | 异步 Correlate | 控制异步表值函数的并发请求数量。 |
table.exec.async-table.timeout | 异步 Correlate | 设置调用的总等待时间,包含启用的重试。 |
table.exec.async-table.retry-strategy | 异步 Correlate | 设置函数调用的重试策略。 |
table.exec.async-lookup.buffer-capacity | 异步 Lookup Join | 控制异步查找的缓冲容量。 |
table.exec.async-lookup.timeout | 异步 Lookup Join | 设置异步维表查找的超时时间,也可由适用的 LOOKUP Hint 为单个关联指定。 |
开发示例
实现 AsyncTableFunction 时,需要声明公开、非静态的 eval 方法。第一个参数是用于提交结果的 CompletableFuture<Collection<T>>,其余参数是 SQL 调用时传入的业务参数。
函数发起异步请求后及时返回,在请求完成的回调中调用 result.complete(rows) 提交结果,或调用 result.completeExceptionally(error) 提交异常。一次调用的多行结果通过一个集合提交,无结果时提交空集合。
package com.example.udf;
import org.apache.flink.table.functions.AsyncTableFunction;
import org.apache.flink.table.annotation.DataTypeHint;
import org.apache.flink.table.annotation.FunctionHint;
import org.apache.flink.types.Row;
import org.apache.flink.table.functions.FunctionContext;
import java.io.IOException;
import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;
import java.nio.charset.StandardCharsets;
import java.time.Duration;
import java.util.Collection;
import java.util.Collections;
import java.util.stream.Collectors;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.CompletionException;
/**
* 通过异步 HTTP 接口查询商品标签。
*
* <p>接口约定:200 返回 UTF-8 文本,每行一个标签;204、404 表示无标签。
* 本示例适用于返回短文本的只读查询;其他状态码作为调用失败交给引擎处理。
*/
@FunctionHint(output = @DataTypeHint("ROW<tag STRING>"))
public class AsyncProductTags extends AsyncTableFunction<Row> {
private transient HttpClient client;
private transient URI serviceBaseUri;
@Override
public void open(FunctionContext context) {
String baseUrl = context.getJobParameter("product.service.base-url", null);
if (baseUrl == null || baseUrl.isBlank()) {
throw new IllegalArgumentException("请配置 product.service.base-url");
}
serviceBaseUri = URI.create(baseUrl);
if (!("https".equalsIgnoreCase(serviceBaseUri.getScheme())
|| "http".equalsIgnoreCase(serviceBaseUri.getScheme()))
|| serviceBaseUri.getHost() == null
|| serviceBaseUri.getRawQuery() != null
|| serviceBaseUri.getRawFragment() != null
|| !baseUrl.endsWith("/")) {
throw new IllegalArgumentException("服务地址须为以 / 结尾的 HTTP(S) 基础地址");
}
client = HttpClient.newBuilder().connectTimeout(Duration.ofSeconds(3)).build();
}
@Override
public boolean isDeterministic() {
// 商品信息可能变化,相同 ID 在不同时间查询可能得到不同结果。
return false;
}
public void eval(CompletableFuture<Collection<Row>> result, Long productId) {
if (productId == null) {
result.complete(Collections.emptyList());
return;
}
try {
HttpRequest request =
HttpRequest.newBuilder(
serviceBaseUri.resolve("products/" + productId + "/tags"))
.timeout(Duration.ofSeconds(5))
.header("Accept", "text/plain")
.GET()
.build();
client.sendAsync(request, HttpResponse.BodyHandlers.ofString(StandardCharsets.UTF_8))
.thenApply(
response -> {
if (response.statusCode() == 200) {
return response.body().lines()
.map(String::strip)
.filter(tag -> !tag.isEmpty())
.map(tag -> Row.of(tag))
.collect(Collectors.toList());
}
if (response.statusCode() == 204 || response.statusCode() == 404) {
return Collections.<Row>emptyList();
}
throw new CompletionException(
new IOException(
"商品查询失败,HTTP 状态码:" + response.statusCode()));
})
.whenComplete(
(rows, error) -> {
if (error != null) {
result.completeExceptionally(error);
} else {
result.complete(rows);
}
});
} catch (RuntimeException e) {
result.completeExceptionally(e);
}
}
}HttpClient.sendAsync 异步执行网络请求,函数无需自行创建线程池。回调只完成短文本解析和结果提交。使用其他客户端时,也应优先选择其原生异步接口。JDK 客户端的行为请参见 HttpClient API。
商品标签可能随时间变化,因此示例将 isDeterministic() 设置为 false。对于同样的输入始终返回相同结果的函数,可保留默认行为。
自定义超时处理
默认情况下,如果某次 eval 调用在超时时间之后仍未完成 CompletableFuture,且配置的重试次数也已用尽,框架会抛出 java.util.concurrent.TimeoutException 并让作业失败。支持在 AsyncTableFunction 子类中声明一个与 eval 方法相匹配的 timeout(...) 方法,由该方法返回一行兜底数据,或实现自定义的超时处理逻辑。
调用远程 LLM 集群等场景的超时概率相对较高,通常不希望超时直接导致作业失败。声明 timeout 方法后,可在超时时返回兜底结果,避免作业中断。
timeout 方法需要同时满足以下条件:
方法声明:必须是 public 实例方法,不能是 static,方法名固定为 timeout。
签名与 eval 对齐:参数列表需要与对应的 eval 严格一致,第一个参数是 CompletableFuture,其泛型类型必须与 eval 完全相同,其余参数的类型和顺序也必须与 eval 完全相同。timeout 支持重载,每一个需要兜底逻辑的 eval 重载都可以单独声明一个对应的 timeout。
同步完成(强制要求):该回调在算子的 mailbox 线程上运行,必须在方法返回前完成 future。方法返回后,框架会立即检查 future.isDone();如果尚未完成,框架会强制以 IllegalStateException 完成 future。因此:
不要在方法体内再发起一次异步调用(例如又一次重试或二次查询),并指望其回调最终完成 future。回调触发时,框架已放弃该条记录。
不要启动新线程异步完成 future,原因同上。
方法体只做一件事:给出轻量的兜底结果,例如返回一个常量行、一个 NULL 行、空集合,或使用业务自定义异常调用 completeExceptionally(...)。
抛异常是安全的:直接在方法体内同步抛出异常没有问题,框架会把异常转交给下游的 ResultFuture,效果等价于显式调用 future.completeExceptionally(thrown)。
在以上约束之外,框架还会在以下场景保证稳定的行为:
未声明 timeout 方法:沿用框架默认的 TimeoutException,不会引发代码生成阶段的失败。
签名不兼容:如果 timeout 方法可见性合法,但参数列表无法和当前调用点的实参类型对应,框架会在规划期(代码生成阶段)以 ValidationException 直接失败。错误信息中包含函数的全限定类名,并同时给出期望签名与实际签名。
空集合兜底:调用 future.complete(Collections.emptyList()) 时,对 INNER lookup join 而言该行会被丢弃;对 LEFT OUTER lookup join 而言,右表字段会被填充为 NULL。
以下示例在异步表值函数的基础上增加 timeout 回调,展示如何同步完成 future 并返回一行兜底数据:
import org.apache.flink.table.functions.AsyncTableFunction;
import java.util.Collection;
import java.util.Collections;
import java.util.concurrent.CompletableFuture;
public static class BackgroundFunctionWithTimeout extends AsyncTableFunction<Long> {
public void eval(CompletableFuture<Collection<Long>> future, Integer waitMax) {
// ...逻辑与前述开发示例一致:把异步任务派发出去,
// 然后在回调线程中完成 `future`。
}
// 参数列表必须与 eval() 对齐:CompletableFuture 的泛型相同,
// 后续参数(Integer waitMax)的类型和顺序也相同。
// 方法体必须同步完成 `future`——不要再发起异步调用,
// 也不要启动新线程异步去完成它;这个回调被触发时,
// 算子已经放弃了这条记录。
public void timeout(CompletableFuture<Collection<Long>> future, Integer waitMax) {
future.complete(Collections.singletonList(-1L));
// 也可以选择用业务自定义异常代替默认的 TimeoutException:
// future.completeExceptionally(new RuntimeException(waitMax));
}
}UDTF注册
UDTF注册过程,请参见管理自定义函数(UDF)。
UDTF使用
在注册UDTF完成后,您就可以使用UDTF,详细的操作步骤如下。
Flink SQL作业开发。详情请参见作业开发地图。
ASI_UDTF_Source表中message字段每行字符串按照竖线(|)分割成多列,代码示例如下。
CREATE TEMPORARY TABLE ASI_UDTF_Source ( `message` VARCHAR ) WITH ( 'connector'='datagen' ); CREATE TEMPORARY TABLE ASI_UDTF_Sink ( name VARCHAR, place VARCHAR ) WITH ( 'connector' = 'blackhole' ); INSERT INTO ASI_UDTF_Sink SELECT name,place FROM ASI_UDTF_Source,lateral table(ASI_UDTF(`message`)) as T(name,place);在页面,单击目标作业名称操作列的启动。
启动成功后,ASI_UDTF_Sink表会被插入ASI_UDTF_Source表中message字段按照竖线(|)分割成多列的字符。