自定义表值函数(UDTF)

更新时间:
复制 MD 格式

本文为您介绍如何为Flink自定义表值函数(UDTF)开发、注册和使用流程。

定义

自定义表值函数(UDTF),自定义表值函数,将0个、1个或多个标量值作为输入参数(可以是变长参数)。与自定义的标量函数类似,但与标量函数不同。表值函数可以返回任意数量的行作为输出,而不仅是1个值。返回的行可以由1个或多个列组成。调用一次函数输出多行或多列数据。详情参见User-defined Functions。

UDTF开发

说明

Flink为您提供了UDF示例,便于您快速开发业务。Flink UDF示例中包含UDSF、UDAF和UDTF的实现,示例中已为您配置对应版本的开发环境,您无需进行环境搭建。

  1. 下载并解压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代码。

  2. 在IntelliJ IDEA中,单击file > open,打开刚才解压缩完成的ASI_UDX-main。

  3. 双击打开\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()));
          }
        }
      }
  4. 双击打开\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>
  5. 在下载文件中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,详细的操作步骤如下。

  1. 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);
  2. 在运维中心 > 作业运维页面,单击目标作业名称操作列的启动。

    启动成功后,ASI_UDTF_Sink表会被插入ASI_UDTF_Source表中message字段按照竖线(|)分割成多列的字符。