FETCH_CONTENT, TRY_FETCH_CONTENT

更新时间:
复制 MD 格式

下载给定URI中的文件并以Bytes Array格式返回文件内容

FETCH_CONTENT 和 TRY_FETCH_CONTENT 用于异步下载指定 URI 中的文件,并以 VARBINARY 类型返回文件内容。两者的区别在于下载失败后的处理方式:

函数

下载成功

下载失败且重试耗尽

适用场景

FETCH_CONTENT

返回文件内容

抛出异常,作业失败

文件内容是后续处理的必要输入,不允许跳过失败记录

TRY_FETCH_CONTENT

返回文件内容

返回 NULL,作业继续处理其他记录

允许跳过无效 URI、文件不存在、网络异常等失败记录

使用限制

  • 仅实时计算引擎VVR 11.7.0及以上版本支持FETCH_CONTENT函数。

  • 仅实时计算引擎VVR 11.9.0-preview.1及以上版本支持TRY_FETCH_CONTENT函数。

语法

VARBINARY FETCH_CONTENT(VARCHAR uri)
VARBINARY FETCH_CONTENT(VARCHAR uri, INTEGER concurrency)

VARBINARY TRY_FETCH_CONTENT(VARCHAR uri)
VARBINARY TRY_FETCH_CONTENT(VARCHAR uri, INTEGER concurrency)

入参

参数

数据类型

说明

uri

VARCHAR

该参数指定需要下载的文件路径。支持的URI schemes包括HTTPFlink FileSystem所支持的Schema:

  • http:// or https:// : HTTP/HTTPS链接

  • oss:// : 阿里云OSS路径

  • hdfs:// : HDFS路径

  • file:// : 本地文件路径

concurrency

INTEGER

可选参数,用于设置每个 FETCH_CONTENT 函数实例的专用 I/O 线程池大小。该线程池负责文件系统读取、HTTP 客户端的异步任务和回调等。参数省略时,默认值为max(8, 当前 JVM 可用处理器数)。例如可用处理器数为 4 时默认值为 8,为 16 时默认值为 16。

说明
  • 如果uriNULL,则返回NULL。如果非NULLuri下载失败,则抛出相关异常。

  • 如果uriOSS路径,则需要配置OSS相关鉴权信息,参考配置Bucket鉴权信息

返回结果

数据类型

说明

VARBINARY

文件内容

TRY_FETCH_CONTENT 只容忍运行时的内容获取失败。函数参数个数或类型错误、concurrency 超出取值范围等 SQL 校验错误仍会导致作业提交失败。

重试和超时配置

FETCH_CONTENT 和 TRY_FETCH_CONTENT 使用异步标量函数的统一重试和超时配置。请求并发、超时、重试由 Flink Async Scalar 算子统一管理。您可以在 SQL 作业中通过 SET 命令修改相关参数。

以下示例将重试策略设置为固定延迟重试,最多重试 2 次,两次请求之间等待 1 s,单条记录的异步调用超时时间为 30 s

SET 'table.exec.async-scalar.retry-strategy' = 'FIXED_DELAY';
SET 'table.exec.async-scalar.max-attempts' = '2';
SET 'table.exec.async-scalar.retry-delay' = '1 s';
SET 'table.exec.async-scalar.timeout' = '30 s';

配置 table.exec.async-scalar.max-attempts 为 2 时,函数最多发起 3 次请求,即首次请求加 2 次重试。如果在超时时间内没有足够时间完成全部重试,则以超时结果为准:FETCH_CONTENT 抛出异常,TRY_FETCH_CONTENT 返回 NULL

完整的参数说明如下:

参数名

默认值

说明

table.exec.async-scalar.timeout

3 min

单条输入等待异步操作完成的超时时间。

table.exec.async-scalar.max-attempts

3

最大重试次数,不包含首次调用;默认最多调用 4 次。

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

10

每个subtask 最多允许同时处于等待状态的异步操作数。该参数和concurrency 共同决定实际并发度。前者限制每个 subtask 的在途输入数量,后者限制函数实例可使用的 I/O 执行资源。实际并发通常受两者较小值限制。

table.exec.async-scalar.retry-strategy

FIXED_DELAY

重试策略,可设置为

  • FIXED_DELAY:固定间隔重试

  • NO_RETRY:不重试

table.exec.async-scalar.retry-delay

100 ms

两次调用之间的固定重试间隔。

示例1:下载文件内容

  • 测试数据

    表 1. T1

    input

    uri(VARCHAR)

    1

    http://example.com/image_url

    2

    oss://example-bucket/example.pdf

    3

    NULL

  • 测试语句

    SELECT 
        id,
        FETCH_CONTENT(uri) AS `value`
    FROM 
        T1;
  • 返回结果如下

    id (INT)

    value (VARCHAR)

    1

    x'ffd8ffe00010......'

    2

    x'aaffd8ffe000......'

    3

    NULL

示例 2:下载失败时返回 NULL

  • 测试数据

    表 1. T1

    input

    uri(VARCHAR)

    1

    http://example.com/image_url

    2

    oss://example-bucket/example.pdf

    3

    invalid://path

  • 使用 TRY_FETCH_CONTENT 下载文件。某条记录下载失败时,该记录的 content 返回 NULL,不会因为该次下载失败而中断作业。

    SELECT
      id,
      TRY_FETCH_CONTENT(uri) AS content
    FROM T2;
  • 返回结果如下

    id (INT)

    value (VARCHAR)

    1

    x'ffd8ffe00010......'

    2

    NULL

    3

    NULL

  • 您可以过滤下载失败的记录,避免将 NULL 传递给后续的多模态推理函数。

    SELECT id, content
    FROM (
      SELECT
        id,
        TRY_FETCH_CONTENT(uri) AS content
      FROM T2
    )
    WHERE content IS NOT NULL;

示例 3:指定并发数

以下示例将每个算子实例的内容获取并发数设置为 88 必须是 SQL 中的整数字面量。

SELECT
  id,
  TRY_FETCH_CONTENT(uri, 8) AS content
FROM T2;

以下调用不合法,会在 SQL 校验阶段报错。

-- concurrency 不能引用字段。SELECT TRY_FETCH_CONTENT(uri, concurrency_column) FROM T3;

-- concurrency 的有效范围为 1~1024。SELECT TRY_FETCH_CONTENT(uri, 0) FROM T3;
SELECT TRY_FETCH_CONTENT(uri, 1025) FROM T3;

相关文档