Blob存储

更新时间:
复制 MD 格式

Blob 存储用于在 DLF Paimon 表中管理图片、音视频和文档等二进制数据。BLOB 数据与结构化数据按列分离存储并支持按需读取;对于无需复制的数据,可以使用 Blob View 引用上游表中的 BLOB。

概述

随着 AI 和多模态应用的发展,数据湖需要同时管理结构化数据(数值、文本)与非结构化数据(图片、视频、音频、文档)。传统方案通常将非结构化数据存放在 OSS 中,将结构化元数据存放在数据库或数据湖表中,两类数据分别管理,难以统一进行权限控制和生命周期管理。

DLF Paimon 表引入 BLOB 类型,将非结构化数据和结构化数据存储在同一张表中。

核心优势如下:

  • 列级分离存储:BLOB 数据自动写入独立的 .blob 文件,结构化列写入 Parquet、ORC 等数据文件。

  • 高效列裁剪:查询非 BLOB 列时不加载 BLOB 数据,避免无效 I/O。

  • 集合类型支持:支持 BLOBARRAY<BLOB>MAP<K, BLOB>

  • 统一管理:BLOB 文件存储在表路径下,由 DLF 统一管理元数据、权限和生命周期。

存储模式与数据类型

存储模式

DLF Paimon Append 表支持以下 BLOB 存储模式:

存储模式

说明

托管 BLOBblob-field

BLOB 原始内容写入表路径下的 .blob 文件,由 DLF 管理文件生命周期。

Blob Viewblob-view-field

下游表仅保存指向上游表 BLOB 数据的引用,读取时根据上游表、字段和 _ROW_ID 解析实际内容,无需复制 BLOB 数据或创建新的 .blob 文件。

SQL 中使用 BINARYBYTES 声明字段。托管 BLOB 通过 blob-field 表选项声明;Blob View 通过 blob-view-field 表选项声明。

Blob View 仅支持单值 BLOB。上游表必须开启 Row Tracking,以便通过稳定的 _ROW_ID 定位被引用的数据。

托管 BLOB 支持的数据类型

托管 BLOB 模式支持以下类型:

  • 单值 BLOB:BINARY/BYTES

  • BLOB 数组:ARRAY<BINARY>/ARRAY<BYTES>

  • BLOB Map:MAP<K, BINARY>/MAP<K, BYTES>

MAP<K, BLOB> 的 Key 支持整数类型、CHARVARCHAR,且不能为 NULL。

BLOB 相关表选项

选项

必填

默认值

说明

row-tracking.enabled

false

BLOB 表必须开启 Row Tracking。

data-evolution.enabled

false

BLOB 表必须开启 Data Evolution。

blob-field

指定使用托管 BLOB 模式的字段,多个字段使用逗号分隔。

blob-view-field

指定使用 Blob View 模式的字段,多个字段使用逗号分隔。

blob-view.resolve.enabled

true

是否在读取时将 Blob View 解析为上游 BLOB 内容。设置为 false 时保留引用信息。

blob-as-descriptor

false

设置为 false 时读取实际 BLOB 内容;设置为 true 时返回 BLOB Descriptor。该选项只改变读取结果,不改变存储方式。

blob.target-file-size

target-file-size

.blob 文件的滚动目标大小,不是单个 BLOB 对象的最大容量。

blob-write-null-on-missing-file

false

Flink 从文件 Descriptor 写入时,如果源文件不存在,是否写入 NULL。

blob-write-null-on-fetch-failure

false

Flink 从文件 Descriptor 写入时,如果源文件读取失败,是否写入 NULL。

使用限制

  • 必须同时开启 row-tracking.enableddata-evolution.enabled

  • BLOB 列不能作为分区列。

  • Blob View 仅支持单值 BLOB,并依赖上游表及对应行持续可用;上游数据被删除或无法读取时,Blob View 无法解析。

  • ARRAY<BLOB>MAP<K, BLOB> 仅支持托管 BLOB 存储。

  • BLOB 对象可以大于 2 GiB,但直接读取为 BINARYBYTES 或 Python bytes 时,单次物化长度不能超过 Integer.MAX_VALUE 字节。超过该大小时,需要读取 Descriptor,并使用流式 API。

使用Blob存储

通过 EMR Serverless Spark 使用

关于 EMR Serverless Spark 对接 DLF 的基础配置,请参见Serverless Spark 访问 DLF

下载并配置 JAR

使用 BLOB 功能需要使用以下 JAR:

下载附件后,将 JAR 上传至具有访问权限的 OSS 路径,并配置以下参数:

spark.emr.serverless.excludedModules    paimon
spark.emr.serverless.user.defined.jars  oss://my-bucket/jars/paimon-ali-emr-spark-3.5-1-ali-29.1.jar

建表

Spark SQL 使用 BINARY 类型声明 BLOB 列:

CREATE TABLE my_db.image_table (
    id BIGINT,
    name STRING,
    category STRING,
    image BINARY
) TBLPROPERTIES (
    'row-tracking.enabled' = 'true',
    'data-evolution.enabled' = 'true',
    'blob-field' = 'image'
);

托管 BLOB 模式同时支持集合类型。多个托管 BLOB 字段通过 blob-field 表选项统一声明:

CREATE TABLE my_db.gallery_table (
    id BIGINT,
    gallery ARRAY<BINARY>,
    renditions MAP<STRING, BINARY>
) TBLPROPERTIES (
    'row-tracking.enabled' = 'true',
    'data-evolution.enabled' = 'true',
    'blob-field' = 'gallery,renditions'
);

写入数据

通过 SQL 直接写入二进制数据:

INSERT INTO my_db.image_table VALUES
    (1, 'sample', 'photo', X'89504E470D0A1A0A');

通过 Notebook 从 OSS 读取图片并写入 BLOB 表:

image_df = (
    spark.read
         .format("binaryFile")
         .load("oss://my-bucket/path/test.jpg")
)

image_df.selectExpr(
    "CAST(1 AS BIGINT) AS id",
    "'test.jpg' AS name",
    "'photo' AS category",
    "content AS image"
).writeTo("my_db.image_table").append()

也可以使用 path_to_descriptor 读取 OSS 文件并写入托管 BLOB:

INSERT INTO my_db.image_table VALUES
    (2, 'external.jpg', 'photo',
     sys.path_to_descriptor('oss://my-bucket/images/external.jpg'));

写入过程中会读取源文件,并将内容复制到当前表管理的 .blob 文件中。

查询数据

仅读取结构化列,不加载 BLOB 数据:

SELECT id, name, category
FROM my_db.image_table
WHERE category = 'photo';

读取 BLOB 内容:

SELECT id, name, length(image) AS image_size
FROM my_db.image_table
WHERE id = 1;

通过动态参数读取 BLOB Descriptor:

SET spark.paimon.my_catalog.my_db.image_table.blob-as-descriptor=true;

SELECT id, name, sys.descriptor_to_string(image) AS image_descriptor
FROM my_db.image_table
WHERE id = 1;

RESET spark.paimon.my_catalog.my_db.image_table.blob-as-descriptor;

写入 Blob View 引用

Spark 可以通过 sys.blob_view 根据上游表、BLOB 字段和 _ROW_ID 生成 Blob View 引用,并将该引用写入通过 blob-view-field 表选项声明的下游字段。写入的是引用信息,不会复制上游 BLOB 内容。

USE my_db;

CREATE TABLE image_view_table (
    id BIGINT,
    name STRING,
    image_ref BINARY
) TBLPROPERTIES (
    'row-tracking.enabled' = 'true',
    'data-evolution.enabled' = 'true',
    'blob-view-field' = 'image_ref'
);

INSERT INTO image_view_table
SELECT
    id,
    name,
    sys.blob_view(
        'my_catalog.my_db.image_table',
        'image',
        _ROW_ID
    )
FROM `image_table$row_tracking`;

SELECT id, name, length(image_ref) AS image_size
FROM image_view_table;

sys.blob_view 的参数依次为上游表名称、上游 BLOB 字段名称和上游数据的 _ROW_ID

更新数据

Spark 支持通过 MERGE INTO 更新单值 BLOB、ARRAY<BLOB>MAP<K, BLOB>

MERGE INTO my_db.image_table AS target
USING my_db.image_update_source AS source
ON target.id = source.id
WHEN MATCHED THEN
    UPDATE SET target.image = source.image;

对于无主键的 Append 表,可以在建表时通过 upsert-key 指定业务唯一键。写入相同 upsert-key 的数据时更新已有行,否则插入新行。upsert-key 不能与主键同时配置:

CREATE TABLE my_db.image_upsert_table (
    id BIGINT,
    name STRING,
    category STRING,
    image BINARY
) TBLPROPERTIES (
    'row-tracking.enabled' = 'true',
    'data-evolution.enabled' = 'true',
    'blob-field' = 'image',
    'upsert-key' = 'id'
);

-- 首次写入 id=1,执行插入
INSERT INTO my_db.image_upsert_table VALUES
    (1, 'sample.jpg', 'photo', X'89504E470D0A1A0A');

-- 再次写入 id=1,更新已有数据
INSERT INTO my_db.image_upsert_table VALUES
    (1, 'sample-new.jpg', 'photo', X'FFD8FFE000104A46');

SELECT id, name
FROM my_db.image_upsert_table
WHERE id = 1;
-- 返回:1, sample-new.jpg

通过实时计算 Flink 使用

关于实时计算 Flink 对接 DLF 的基础配置,请参见实时计算 Flink 版访问 DLF

下载并配置 JAR

使用 BLOB 功能需要下载以下 JAR:

该包的 Catalog Factory 和 Table Factory Identifier 均为 paimon-1-ali-29-20260731

在实时计算 Flink 控制台的数据管理中上传该 JAR,并通过自定义 Catalog 使用该 JAR。

添加 Catalog

CREATE CATALOG my_catalog WITH (
    'type' = 'paimon-1-ali-29-20260731',
    'metastore' = 'rest',
    'token.provider' = 'dlf',
    'uri' = 'http://cn-hangzhou-vpc.dlf.aliyuncs.com',
    'warehouse' = 'my_catalog'
);

建表与写入

Flink SQL 使用 BYTES 类型声明 BLOB 列:

CREATE TABLE my_catalog.my_db.image_table (
    id BIGINT,
    name STRING,
    category STRING,
    image BYTES
) WITH (
    'row-tracking.enabled' = 'true',
    'data-evolution.enabled' = 'true',
    'blob-field' = 'image'
);

INSERT INTO my_catalog.my_db.image_table VALUES
    (1, 'cat.jpg', 'photo', X'89504E470D0A1A0A'),
    (2, 'dog.jpg', 'photo',
     my_catalog.sys.path_to_descriptor(
         'oss://my-bucket/images/dog.jpg'
     ));

写入 path_to_descriptor 时,Paimon 会读取源文件并将内容写入当前表管理的 .blob 文件。

查询数据

仅读取结构化列:

SELECT id, name, category
FROM my_catalog.my_db.image_table;

读取实际 BLOB 内容:

SELECT id, name
FROM my_catalog.my_db.image_table
/*+ OPTIONS('blob-as-descriptor'='false') */;

读取 BLOB Descriptor:

SELECT
    id,
    name,
    my_catalog.sys.descriptor_to_string(image) AS image_descriptor
FROM my_catalog.my_db.image_table
/*+ OPTIONS('blob-as-descriptor'='true') */;

使用 Blob View

Blob View 用于在下游表中引用上游表已有的 BLOB 数据,无需复制 BLOB 内容或创建新的 .blob 文件。下游字段通过 blob-view-field 表选项声明,仅保存引用信息,读取时解析为上游 BLOB 内容。

USE CATALOG my_catalog;
USE my_db;

CREATE TABLE image_view_table (
    id BIGINT,
    name STRING,
    image_ref BYTES
) WITH (
    'row-tracking.enabled' = 'true',
    'data-evolution.enabled' = 'true',
    'blob-view-field' = 'image_ref'
);

INSERT INTO image_view_table
SELECT
    id,
    name,
    sys.blob_view(
        'my_catalog.my_db.image_table',
        'image',
        _ROW_ID
    )
FROM `image_table$row_tracking`;

sys.blob_view 的参数依次为上游表名称、上游 BLOB 字段名称和上游数据的 _ROW_ID。查询 image_ref 时,默认返回所引用的实际 BLOB 内容。

Flink Data Evolution MERGE INTO 暂不支持直接更新托管 BLOB 字段。如需更新 BLOB 内容,请使用 Spark。

通过 PyPaimon 使用

PyPaimon 使用 PyArrow 的 large_binary() 映射 BLOB 类型。

下载并安装 PyPaimon

下载PyPaimon安装包:pypaimon-1.5.dev20260727.tar.gz

执行以下命令:

pip install pypaimon-1.5.dev20260727.tar.gz

建表与写入 BLOB 数据

import pyarrow as pa
from pypaimon import CatalogFactory, Schema

catalog = CatalogFactory.create({
    'metastore': 'rest',
    'uri': 'https://${regionId}-vpc.dlf.aliyuncs.com',
    'warehouse': 'my_catalog',
    'token.provider': 'dlf',
    'dlf.access-key-id': '<AK>',
    'dlf.access-key-secret': '<SK>',
})

pa_schema = pa.schema([
    ('id', pa.int64()),
    ('name', pa.string()),
    ('picture', pa.large_binary()),
])

schema = Schema.from_pyarrow_schema(
    pa_schema,
    options={
        'row-tracking.enabled': 'true',
        'data-evolution.enabled': 'true',
    },
)

catalog.create_table('my_db.image_table', schema, True)
table = catalog.get_table('my_db.image_table')

write_builder = table.new_batch_write_builder()
writer = write_builder.new_write()
commit = write_builder.new_commit()

data = pa.Table.from_pydict({
    'id': [1],
    'name': ['sample.png'],
    'picture': [b'\x89PNG\r\n\x1a\n'],
}, schema=pa_schema)

writer.write_arrow(data)
commit.commit(writer.prepare_commit())
writer.close()
commit.close()

创建 Blob View 下游表时,在 Schema options 中通过 blob-view-field 指定引用字段:

blob_view_pa_schema = pa.schema([
    ('id', pa.int64()),
    ('picture_ref', pa.large_binary()),
])

blob_view_schema = Schema.from_pyarrow_schema(
    blob_view_pa_schema,
    options={
        'row-tracking.enabled': 'true',
        'data-evolution.enabled': 'true',
        'blob-view-field': 'picture_ref',
    },
)

catalog.create_table(
    'my_db.image_view_table',
    blob_view_schema,
    True,
)

读取 BLOB 数据

推荐按 Batch 读取 BLOB 数据:

read_builder = table.new_read_builder()
splits = read_builder.new_scan().plan().splits()
read = read_builder.new_read()

for batch in read.to_arrow_batch_reader(
        splits, blob_parallelism=16):
    for i in range(len(batch)):
        picture = batch['picture'][i].as_py()

仅读取结构化列:

read_builder = table.new_read_builder().with_projection(['id', 'name'])
splits = read_builder.new_scan().plan().splits()
result = read_builder.new_read().to_pandas(splits)

对于大型 BLOB,可以设置 blob-as-descriptor=true,再通过流式 API 分段读取,避免一次性将完整内容加载到内存。

存储布局

定义 BLOB 列后,DLF Paimon 自动分离存储:

table/
├── bucket-0/
│   ├── data-uuid-0.parquet
│   ├── data-uuid-1.blob
│   ├── data-uuid-2.blob
│   └── ...
├── manifest/
├── schema/
└── snapshot/

结构化列写入 Parquet、ORC 等普通数据文件,BLOB 内容写入 .blob 文件。