非结构化的文件路径及其他信息存入Doris引擎
更新时间:2026-07-31
概览
说明:下文所述Doris引擎,为平台中计算资源功能下 创建的 分析与AI搜索实例
百度胜算平台支持对图片等非结构化数据信息写入Doris引擎。平台提供两种方式:
- 工作流模版写入Doris:通过可视化工作流对图片等非结构化数据进行特征嵌入生成向量,将图像数据与向量写入Doris表,快速完成向量数据存储,开箱即用。
- PySpark任务写入Doris:通过PySpark脚本,将Iceberg表中的数据同步到Doris表,适合已有Iceberg存量数据,需要批量迁移、持续同步至Doris的业务场景。
注意:只有写入Doris的数据包含向量字段,且目标Doris表预先创建ANN向量索引,入库后才可使用Doris交互查询开展相似度检索。 写入数据完成后,可使用Notebook交互查询能力,对已写入的数据开展查询。
将数据写入Doris
方式一:通过工作流模板写入Doris
本方案以图片为例,同样适用于其他类型的非结构化数据。 通过百度胜算工作流预置模版的能力,依托Ray算子完成图像解析、多模态特征向量提取,将图像信息与向量数据写入Doris,支撑图像相似度检索。
前置条件
- 已创建项目;
- 已创建Ray实例(用于工作流算子任务);
- 已创建Doris实例(作为标签数据的存储引擎);
- 准备好一张.jpg图片(用于工作流处理)。
步骤一:创建数据卷
- 登录百度胜算控制台,进入已经创建好的工作空间。在侧边导航栏选择元数据。
- 单击新建按钮,选择创建数据目录,配置数据目录名称为demo_test。

- 选择已创建好的数据目录demo_test,依次单击default>立即创建>创建数据卷。
- 配置数据卷名称为image_data,数据卷类型选择内部卷,然后单击确认。

- 重复上述步骤,再创建1个数据卷,命名为output_image。各数据卷的作用说明如下:
| 卷名 | 作用 |
|---|---|
| image_data | 用来存储要处理的数据的路径信息。 |
| output_image | 用来存储算子处理后的数据。 |
- 在image_data数据卷右上角单击上传文件按钮,选择一张事前准备好的.jpg图片上传即可。
步骤二:创建Doris数据表
- 在元数据中找到Doris数据目录,依次单击default>立即创建>创建数据表。
- 在右侧DDL语句中输入以下信息,然后单击确定即可。
SQL
1CREATE TABLE `test_only_image` (
2 `images` varchar(2560) NULL DEFAULT "" COMMENT "图片路径",
3 `embedding` array<float> NOT NULL COMMENT "图片向量",
4 INDEX index_emb (`embedding`) USING ANN PROPERTIES("distance" = "cosine", "hnsw_m" = "16", "dimension" = "4096", "algorithm" = "hnsw", "hnsw_efConstruction" = "200")
5) ENGINE=OLAP
6UNIQUE KEY(`images`)
7DISTRIBUTED BY HASH(`images`) BUCKETS 10;
步骤三:创建图片入库模版工作流
- 在左侧导航栏单击数据处理>工作流,单击预置模版页签。
- 在搜索框输入关键词图片入库进行搜索,选择图片入库模版,单击使用模版。

- 在使用模版创建工作流对话框输入名称为picture_demo,所属位置选择已经创建好的项目,然后单击确定。

- 展开图片入库任务节点,单击文件元数据加载器算子,配置算子参数data_path为待处理图片所在的文件夹路径。

- 单击多模态数据入库算子,table_name默认为test_only_image;检查是否配置doris_compute_id,若没有配置,则需前往分析与AI搜索实例页面复制Doris实例ID。
说明:在创建Doris实例后,会默认填入。

- 最后单击数据输出器算子,配置算子参数export_path为输出的数据文件Volume路径。

- 选中图片入库任务节点,单击右侧执行资源按钮,核对算力配置无误,确保资源规格满足图片入库需求。

- 确认全部参数配置无误后,单击保存按钮。
步骤四:运行工作流并查询结果
运行步骤一配置好的工作流,将图像数据写入Doris,然后查询验证结果。
- 单击立即运行按钮,然后单击运行记录,查看是否运行成功。

- 运行成功后,前往步骤二创建的Doris表中,单击数据预览查看结果。

方式二:通过PySpark任务写入Doris
适用于已有Iceberg数据表,需要将数据同步到Doris进行标签化查询的场景。
前置条件
- 已创建通用独占资源队列。
- 已创建Spark实例。
步骤一:在元数据模块创建表
在元数据中创建Iceberg表和Doris数据表,确保两边表结构(schema)保持一致。
- 在方法一创建好的demo_test数据目录中,依次单击default>立即创建>创建数据表;
- 在右侧DDl语句中输入以下信息,单击确定即可。
SQL
1CREATE TABLE iceberg_demo (
2 id INT,
3 data STRING
4)
5USING ICEBERG
- 然后前往Doris数据目录,依次单击default>立即创建>创建数据表。
- 在右侧DDL语句中输入以下信息,然后单击确定即可。
SQL
1CREATE TABLE `doris_demo` (
2 `id` INT NULL,
3 `data` STRING NULL
4)
5UNIQUE KEY (`id`)
6DISTRIBUTED BY HASH (`id`) BUCKETS 10
步骤二:在项目文件里上传pyspark脚本文件
- 将以下代码保存为python文件。
Python
1#!/usr/bin/env python3
2"""
3通用 Iceberg -> Doris 同步脚本
4
5用途:
6- 直接读取任意 schema 的 Iceberg 表,并原样写入 Doris 表。
7- 不对 Iceberg 表的字段做任何假设或解析,读到什么列就写什么列。
8- 默认约定:Iceberg 表和 Doris 表的 schema 一致(列名、顺序、类型可对应)。
9 如需只写部分列,用 --columns 指定。
10
11平台配置边界:
12- Spark Catalog、Warehouse、BOS/HDFS、Iceberg Runtime 和 Doris Connector 由平台或提交命令配置。
13- 本脚本不覆盖平台已经配置的 Catalog 属性,避免连接到错误的 Iceberg 存储位置。
14"""
15
16import argparse
17import os
18import re
19from datetime import date
20
21from pyspark.sql import SparkSession, functions as F
22
23IDENTIFIER = re.compile(r"^[A-Za-z_][A-Za-z0-9_]*$")
24
25
26def env_or_default(name, default=None):
27 """读取环境变量,未设置或为空时使用默认值。"""
28 value = os.getenv(name)
29 return value if value not in (None, "") else default
30
31
32def parse_args():
33 """读取任务参数。
34
35 参数来源优先级:命令行参数 > 环境变量 > 脚本默认值。
36 Spark Catalog、Warehouse、凭证和 Iceberg 实现不在脚本中重设,由平台统一注入。
37 """
38 parser = argparse.ArgumentParser(
39 description="通用 Iceberg -> Doris 同步(表名可配置,不强制 schema)"
40 )
41
42 # ---- Iceberg 源表(可配置)----
43 parser.add_argument(
44 "--catalog",
45 default=env_or_default("ICEBERG_CATALOG", "deepway_catalog"),
46 help="平台已注册的 Iceberg Catalog,默认读取 ICEBERG_CATALOG",
47 )
48 parser.add_argument(
49 "--iceberg-namespace",
50 default=env_or_default("ICEBERG_NAMESPACE"),
51 help="Iceberg namespace(库名),默认读取 ICEBERG_NAMESPACE",
52 )
53 parser.add_argument(
54 "--iceberg-table",
55 default=env_or_default("ICEBERG_TABLE"),
56 help="Iceberg 表名,默认读取 ICEBERG_TABLE",
57 )
58
59 # ---- Doris 目标表(可配置)----
60 parser.add_argument(
61 "--doris-fenodes",
62 default=env_or_default("DORIS_FENODES"),
63 help="Doris FE 地址,例如 fe-host:8030",
64 )
65 parser.add_argument("--doris-user", default=env_or_default("DORIS_USER"))
66 parser.add_argument(
67 "--doris-password", default=env_or_default("DORIS_PASSWORD")
68 )
69 parser.add_argument(
70 "--doris-catalog",
71 default=env_or_default("DORIS_CATALOG", "b0a956f100cc2e3f"),
72 help="Doris catalog.id,默认读取 DORIS_CATALOG",
73 )
74 parser.add_argument(
75 "--doris-db",
76 default=env_or_default("DORIS_DB"),
77 help="Doris 库名,默认读取 DORIS_DB",
78 )
79 parser.add_argument(
80 "--doris-table",
81 default=env_or_default("DORIS_TABLE"),
82 help="Doris 表名,默认读取 DORIS_TABLE",
83 )
84
85 # ---- 可选:列裁剪与日期过滤 ----
86 parser.add_argument(
87 "--columns",
88 default=env_or_default("SYNC_COLUMNS"),
89 help="可选,逗号分隔的列名,只同步这些列;不填则同步 Iceberg 表全部列",
90 )
91 parser.add_argument(
92 "--biz-date",
93 default=env_or_default("BIZ_DATE"),
94 help="可选,业务日期 yyyy-MM-dd;配合 --date-column 使用时只同步当天数据",
95 )
96 parser.add_argument(
97 "--date-column",
98 default=env_or_default("DATE_COLUMN"),
99 help="可选,用于按 --biz-date 过滤的日期/时间列名",
100 )
101
102 args = parser.parse_args()
103
104 # 必填校验
105 required_values = {
106 "--iceberg-namespace": args.iceberg_namespace,
107 "--iceberg-table": args.iceberg_table,
108 "--doris-fenodes": args.doris_fenodes,
109 "--doris-user": args.doris_user,
110 "--doris-password": args.doris_password,
111 "--doris-catalog": args.doris_catalog,
112 "--doris-db": args.doris_db,
113 "--doris-table": args.doris_table,
114 }
115 missing = [name for name, value in required_values.items() if not value]
116 if missing:
117 parser.error(
118 "缺少必填配置:" + ", ".join(missing) + ";请通过命令行参数或环境变量提供"
119 )
120
121 # 标识符校验,防止拼接 SQL 时被注入
122 identifiers = [
123 args.catalog,
124 args.iceberg_namespace,
125 args.iceberg_table,
126 args.doris_db,
127 args.doris_table,
128 ]
129 if args.date_column:
130 identifiers.append(args.date_column)
131 for value in identifiers:
132 if not IDENTIFIER.fullmatch(value):
133 parser.error(f"非法标识符:{value}")
134
135 # 日期过滤参数一致性校验
136 if args.biz_date:
137 try:
138 date.fromisoformat(args.biz_date)
139 except ValueError as exc:
140 parser.error(f"--biz-date 必须为 yyyy-MM-dd:{exc}")
141 if not args.date_column:
142 parser.error("提供了 --biz-date 时必须同时提供 --date-column")
143
144 # 解析并校验 --columns
145 args.column_list = None
146 if args.columns:
147 cols = [c.strip() for c in args.columns.split(",") if c.strip()]
148 for c in cols:
149 if not IDENTIFIER.fullmatch(c):
150 parser.error(f"非法列名:{c}")
151 args.column_list = cols
152
153 return args
154
155
156def build_spark(args):
157 """创建 SparkSession,但不覆盖平台已经注入的 Catalog 配置。"""
158 return (
159 SparkSession.builder.appName(
160 f"iceberg-to-doris-{args.iceberg_table}"
161 )
162 .config("spark.sql.session.timeZone", "UTC")
163 .config("spark.sql.storeAssignmentPolicy", "ANSI")
164 .config("spark.sql.adaptive.enabled", "true")
165 .getOrCreate()
166 )
167
168
169def read_iceberg(spark, args):
170 """读取 Iceberg 源表,可选按日期过滤和列裁剪。
171
172 不解析、不改写 schema:读到的 DataFrame 结构即为 Iceberg 表结构。
173 """
174 table = f"`{args.catalog}`.`{args.iceberg_namespace}`.`{args.iceberg_table}`"
175 df = spark.table(table)
176
177 if args.biz_date and args.date_column:
178 df = df.where(
179 F.to_date(F.col(args.date_column)) == F.lit(args.biz_date)
180 )
181
182 if args.column_list:
183 df = df.select(*args.column_list)
184
185 return df
186
187
188def write_doris(dataframe, args):
189 """通过 Spark Doris Connector 将 DataFrame 原样写入 Doris 表。
190
191 默认约定两边 schema 一致,因此直接整表写入,不做列映射。
192 使用 JSON 按行写入,兼容结构化字段和数组类型。
193 """
194 (
195 dataframe.write.format("doris")
196 .mode("append")
197 .option("doris.fenodes", args.doris_fenodes)
198 .option("doris.query.port", "9030")
199 .option("doris.table.identifier", f"{args.doris_db}.{args.doris_table}")
200 .option("doris.catalog.id", args.doris_catalog)
201 .option("user", args.doris_user)
202 .option("password", args.doris_password)
203 .option("doris.sink.properties.format", "json")
204 .option("doris.sink.properties.read_json_by_line", "true")
205 .option("doris.sink.batch.size", "10000")
206 .option("doris.sink.max-retries", "3")
207 .save()
208 )
209
210
211def main():
212 """读取 Iceberg 表并原样写入 Doris 表。"""
213 args = parse_args()
214 spark = build_spark(args)
215 try:
216 src = f"{args.catalog}.{args.iceberg_namespace}.{args.iceberg_table}"
217 dst = f"{args.doris_db}.{args.doris_table}"
218 print(f"[1/2] 读取 Iceberg 表 `{src}`")
219 df = read_iceberg(spark, args)
220 print(f" 源表字段:{df.columns}")
221
222 print(f"[2/2] 写入 Doris 表 `{dst}`")
223 write_doris(df, args)
224 print("完成")
225 finally:
226 spark.stop()
227
228
229if __name__ == "__main__":
230 main()
- 进入工作台的目标项目中,依次单击创建>导入文件,上传脚本文件,单击确定即可。
步骤三:创建工作流
- 在左侧导航栏单击数据处理>工作流>创建,在创建工作流面板配置名称、所属位置选择已经目标项目即可。
- 单击创建空白工作流,拖拽PySpark任务至画布中间,进行配置;
- 配置程序文件为步骤二上传的pyspark脚本文件。
- 主类参数配置:["--catalog","deepway_catalog","--iceberg-namespace","xiangong_demo","--iceberg-table","iceberg_demo","--doris-fenodes","192.168.18.30:8030","--doris-user","root","--doris-password","xxx","--doris-catalog","b0a956f100cc2e3f","--doris-db","xiangong_demo","--doris-table","doris_demo"]
说明:这里的"--doris-fenodes","192.168.18.30:8030","--doris-user","root","--doris-password","xxx"可以在创建Doris实例后联系内部对接同学获取,--doris-catalog在Doris数据目录的存储路径最后一段字符串。
- 单击右侧执行资源按钮,新增两组Spark配置:
| 参数名 | 参数值 |
|---|---|
| spark.driver.extraJavaOptions | add-opens=java.base/java.nio=org.apache.arrow.memory.core,ALL-UNNAMED |
| spark.executor.extraJavaOptions | add-opens=java.base/java.nio=org.apache.arrow.memory.core,ALL-UNNAMED |
- 确认全部参数配置无误后,单击保存按钮。
步骤四:运行工作流并查询结果
运行PySpark工作流,将Iceberg数据同步到Doris,然后查询验证同步结果。
- 单击立即运行按钮,然后单击运行记录,查看是否运行成功。

- 运行成功后,前往步骤二创建的Doris表中,单击数据预览查看结果。

数据查询(Notebook)
- 在目标项目中依次单击创建>创建Notebook。
- 连接右上角Doris实例,输入SQL语句查询,然后单击运行即可。

评价此篇文章

