自定义标量函数(UDSF)

更新时间:
复制 MD 格式

本文为您介绍如何开发、注册和使用Flink自定义标量函数(UDSF)。

定义

自定义标量函数(UDSF)将0个、1个或多个标量值映射到一个新的标量值。输入与输出是一对一的关系,即读入一行数据,写出一条输出值。详情参见User-defined Functions。

UDSF开发

说明

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_UDF.java:自定义标量函数(UDSF)示例的Java代码。

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

  3. 双击打开\ASI_UDX-main\src\main\java\ASI_UDF后,根据您的业务,配置ASI_UDF.java。

    该示例中,ASI_UDF.java已配置了获取每条数据中从begin~end位的字符的代码。

    package ASI_UDF;
    
    import org.apache.flink.table.functions.ScalarFunction;
    
    public class ASI_UDF extends ScalarFunction {
        public String eval(String s, Integer begin, Integer end) {
            return s.substring(begin, end);
        }
    }
  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包,即代表完成了UDSF开发工作。

自定义异步标量函数(AsyncScalarFunction)

不同于普通标量函数,自定义异步标量函数需要继承 AsyncScalarFunction,并通过 CompletableFuture 返回结果。

函数需要查询外部数据库、调用 HTTP 接口等存在 I/O 等待的操作时,可以使用异步标量函数并发处理多条记录,减少等待对作业吞吐的影响。字符串截取、数值计算等无需等待外部服务的操作,通常使用同步标量函数即可。

配置参数

通过以下参数控制异步标量函数的运行行为。参数应在提交查询或生成执行计划前设置。

参数

默认值

说明

table.exec.async-scalar.max-concurrent-operations

10

每个异步算子实例允许的最大并发异步操作数。增大该值会增加对外部服务的并发访问压力。

table.exec.async-scalar.timeout

3 min

单条记录异步处理的超时时间,包含重试耗时。

table.exec.async-scalar.retry-strategy

FIXED_DELAY

FIXED_DELAY 表示失败后按固定间隔重试;NO_RETRY 表示不重试。

table.exec.async-scalar.retry-delay

100 ms

固定间隔重试的等待时间,仅在重试策略为 FIXED_DELAY 时生效。

table.exec.async-scalar.max-attempts

3

最大重试次数,仅在重试策略为 FIXED_DELAY 时生效。

table.exec.async-scalar.output-mode

ORDERED

ORDERED 表示保持输入顺序;ALLOW_UNORDERED 表示在输入为纯追加流时允许无序输出。VVR 11.8.0 正式版及后续版本支持。

配置无序输出

默认情况下,异步函数虽然可以并发执行,但后续记录已完成的结果仍需等待前面的记录完成后再输出。当外部请求耗时差异较大,且业务允许结果乱序时,可以开启无序输出,减少慢请求造成的等待。

说明

异步调用和无序输出是两个独立的概念。异步标量函数默认并发调用、按输入顺序输出;配置为允许无序输出后,系统会根据输入数据的变更类型决定是否启用无序输出。

在执行查询前增加以下配置:

SET 'table.exec.async-scalar.output-mode' = 'ALLOW_UNORDERED';

自定义异步标量函数示例

以下示例通过商品服务的 HTTP 接口查询商品名称。函数在 open 中创建并复用一个 HttpClient,在 eval 中调用 sendAsync() 发起请求,通过完成回调返回结果。sendAsync() 返回 CompletableFuture,适合将 HTTP 调用与异步标量函数衔接,详情请参见 JDK HttpClient 文档。

示例使用 3 秒连接超时、5 秒请求超时,适用于响应体较小的商品名称查询。请根据实际接口的响应时间调整。

package com.example.udf;

import org.apache.flink.table.functions.AsyncScalarFunction;
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.concurrent.CompletableFuture;
import java.util.concurrent.CompletionException;

/**
 * 通过异步 HTTP 接口查询商品名称。
 *
 * <p>接口约定:200 返回 UTF-8 商品名称,404 表示商品不存在。
 * 本示例适用于返回短文本的只读查询;其他状态码作为调用失败交给引擎处理。
 */
public class AsyncProductName extends AsyncScalarFunction {

    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<String> result, Long productId) {
        if (productId == null) {
            result.complete(null);
            return;
        }

        try {
            HttpRequest request =
                    HttpRequest.newBuilder(
                                    serviceBaseUri.resolve("products/" + productId + "/name"))
                            .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();
                                }
                                if (response.statusCode() == 404) {
                                    return null;
                                }
                                throw new CompletionException(
                                        new IOException(
                                                "商品查询失败,HTTP 状态码:" + response.statusCode()));
                            })
                    .whenComplete(
                            (name, error) -> {
                                if (error != null) {
                                    result.completeExceptionally(error);
                                } else {
                                    result.complete(name);
                                }
                            });
        } catch (RuntimeException e) {
            result.completeExceptionally(e);
        }
    }
}

开发时请注意:

  • 调用 result.complete(value) 返回结果;需要返回 SQL NULL 时,调用 result.complete(null)。

  • 调用失败时,通过 result.completeExceptionally(exception) 返回异常,交由引擎按重试策略处理。

  • eval 在发起请求后直接返回。使用 sendAsync() 的完成回调传递结果,避免在 eval 或回调中调用阻塞式 send()、get()、join(),或执行耗时计算。

  • 每个函数实例在 open 中初始化并复用客户端,避免每条记录创建客户端。JDK 11/17 的 HttpClient 没有公开的 close() 方法,其内部资源由 JDK 管理。使用提供关闭接口的其他异步客户端时,应在函数的 close 中释放客户端。

  • 商品信息会随时间变化,因此示例将 isDeterministic() 设置为 false。涉及更新流时,还需结合查询语义处理非确定性结果。

UDSF注册

UDSF注册过程,请参见管理自定义函数(UDF)。

UDSF使用

在注册UDSF完成后,您就可以使用UDSF,详细的操作步骤如下。

  1. Flink SQL作业开发。详情请参见作业开发地图。

    获取ASI_UDSF_Source表中a字段中每行字符串中第2~4位的字符,代码示例如下。

    CREATE TEMPORARY TABLE ASI_UDSF_Source (
      a VARCHAR,
      b INT,
      c INT
    ) WITH (
      'connector' = 'datagen'
    );
    
    CREATE TEMPORARY TABLE ASI_UDSF_Sink (
      a VARCHAR
    ) WITH (
      'connector' = 'blackhole'
    );
    
    INSERT INTO ASI_UDSF_Sink
    SELECT ASI_UDSF(a,2,4)
    FROM ASI_UDSF_Source;
  2. 在运维中心 > 作业运维页面,单击目标作业名称操作列的启动。

    启动成功后,ASI_UDSF_Sink表每行会被插入ASI_UDSF_Source表中a字段每行字符串的第2~4位字符。