Blob 存储用于在 DLF Paimon 表中管理图片、音视频和文档等二进制数据。BLOB 数据与结构化数据按列分离存储并支持按需读取;对于无需复制的数据,可以使用 Blob View 引用上游表中的 BLOB。
概述
随着 AI 和多模态应用的发展,数据湖需要同时管理结构化数据(数值、文本)与非结构化数据(图片、视频、音频、文档)。传统方案通常将非结构化数据存放在 OSS 中,将结构化元数据存放在数据库或数据湖表中,两类数据分别管理,难以统一进行权限控制和生命周期管理。
DLF Paimon 表引入 BLOB 类型,将非结构化数据和结构化数据存储在同一张表中。
核心优势如下:
-
列级分离存储:BLOB 数据自动写入独立的
.blob文件,结构化列写入 Parquet、ORC 等数据文件。 -
高效列裁剪:查询非 BLOB 列时不加载 BLOB 数据,避免无效 I/O。
-
集合类型支持:支持
BLOB、ARRAY<BLOB>和MAP<K, BLOB>。 -
统一管理:BLOB 文件存储在表路径下,由 DLF 统一管理元数据、权限和生命周期。
存储模式与数据类型
存储模式
DLF Paimon Append 表支持以下 BLOB 存储模式:
存储模式 |
说明 |
托管 BLOB( |
BLOB 原始内容写入表路径下的 |
Blob View( |
下游表仅保存指向上游表 BLOB 数据的引用,读取时根据上游表、字段和 |
SQL 中使用 BINARY 或 BYTES 声明字段。托管 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 支持整数类型、CHAR 和 VARCHAR,且不能为 NULL。
BLOB 相关表选项
|
选项 |
必填 |
默认值 |
说明 |
|
|
是 |
|
BLOB 表必须开启 Row Tracking。 |
|
|
是 |
|
BLOB 表必须开启 Data Evolution。 |
|
|
否 |
无 |
指定使用托管 BLOB 模式的字段,多个字段使用逗号分隔。 |
|
|
否 |
无 |
指定使用 Blob View 模式的字段,多个字段使用逗号分隔。 |
|
|
否 |
|
是否在读取时将 Blob View 解析为上游 BLOB 内容。设置为 |
|
|
否 |
|
设置为 |
|
|
否 |
|
|
|
|
否 |
|
Flink 从文件 Descriptor 写入时,如果源文件不存在,是否写入 NULL。 |
|
|
否 |
|
Flink 从文件 Descriptor 写入时,如果源文件读取失败,是否写入 NULL。 |
使用限制
-
必须同时开启
row-tracking.enabled和data-evolution.enabled。 -
BLOB 列不能作为分区列。
-
Blob View 仅支持单值 BLOB,并依赖上游表及对应行持续可用;上游数据被删除或无法读取时,Blob View 无法解析。
-
ARRAY<BLOB>和MAP<K, BLOB>仅支持托管 BLOB 存储。 -
BLOB 对象可以大于 2 GiB,但直接读取为
BINARY、BYTES或 Pythonbytes时,单次物化长度不能超过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 文件。