FETCH_CONTENT

更新时间:
复制 MD 格式

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

使用限制

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

语法

VARBINARY FETCH_CONTENT(VARCHAR uri)
VARBINARY 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。

说明

FETCH_CONTENT 使用异步 I/O 拉取 URI 指向的内容。请求并发、超时、重试由 Flink Async Scalar 算子统一管理。相关参数如下:

参数名

默认值

说明

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

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

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

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

出参

数据类型

说明

VARBINARY

文件内容

示例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

相关文档