从 MinIO 加载数据到 StarRocksINSERTFILES() 与 Broker Load 完整实战指南【免费下载链接】starrocksThe worlds fastest open query engine for sub-second analytics both on and off the data lakehouse. With the flexibility to support nearly any scenario, StarRocks provides best-in-class performance for multi-dimensional analytics, real-time analytics, and ad-hoc queries. A Linux Foundation project.项目地址: https://gitcode.com/GitHub_Trending/st/starrocks本文基于 StarRocks 官方文档 Load data from MinIO 编写系统讲解如何将存储在 MinIO 对象存储中的数据加载进 StarRocks 的两种主流方式同步的INSERTFILES()与异步的Broker Load。文章覆盖两种方式的使用场景、配置参数、完整可复现的 SQL 示例、任务进度查询方法并结合仓库源码剖析aws.s3.*系列参数在 BE 端的真实解析逻辑帮助你根据文件格式、数据规模与业务场景选择最合适的加载方案。选型概览两种加载方式的定位StarRocks 为 MinIO 数据加载提供了两条路径各自适用于不同的业务诉求加载方式同步/异步支持的格式推荐场景INSERT FILES()同步Parquet、ORC、CSVCSV 自 v3.3.0 起支持大多数常规场景易用性最高Broker Load异步Parquet、ORC、CSV、JSONJSON 自 v3.2.3 起支持长时间运行的大批量任务、需要 JSON 格式、或加载过程中需要执行 DELETE 等数据变更INSERT FILES()从 v3.1 起可用语法简洁、即写即用是官方推荐的首选方案。但FILES()目前仅支持 Parquet、ORC 与 CSV 三种格式若需加载 JSON 等其他格式或在加载过程中对 Primary Key 表执行 DELETE 等变更则应改用 Broker Load。两种方式互补优先用 INSERTFILES()遇到其不支持的格式或需求时切换到 Broker Load。开始之前三项准备工作无论选择哪种加载方式都需要先完成以下三步。1. 准备源数据确保待加载的源数据已妥善存放在 MinIO bucket 中。官方建议关注 bucket 与 StarRocks 集群所在区域的一致性——当两者位于同一区域时跨区域数据传输成本会显著降低。本文使用官方提供的示例数据集1000 万行用户行为数据进行演示可先用curl下载curl -O https://starrocks-examples.s3.amazonaws.com/user_behavior_ten_million_rows.parquet将下载的 Parquet 文件放入 MinIO并记下 bucket 名称。本文示例统一使用 bucket 名/starrocks对象路径为s3://starrocks/user_behavior_ten_million_rows.parquet。2. 检查权限执行数据加载的用户需要具备目标数据库/表的INSERT 权限Broker Load 同样需要以及对应数据源的读取权限。MinIO 侧则通过 Access Key 体系鉴权请确保用于加载的 Access Key 具备对应 bucket 的读取权限。3. 收集连接信息使用 MinIO Access Key 认证时需要准备以下信息存储数据的bucket名称访问 bucket 中特定对象时的object key对象名MinIO 的endpoint如http://minio:9000用于访问凭据的access key 与 secret key在 MinIO 控制台创建 Access Key 的入口如下图所示使用 INSERTFILES()同步加载FILES()是 StarRocks 提供的表函数能够根据你指定的路径相关属性直接读取云存储中的文件自动推断文件内数据的表结构并以数据行的形式返回文件内容。基于它你可以实现三类操作直接使用 SELECT 查询 MinIO 中的数据使用 CREATE TABLE AS SELECTCTAS 建表并加载数据使用 INSERT 将数据加载进已存在的表。核心参数说明INSERTFILES()通过一组aws.s3.*属性描述 MinIO 的连接方式。下表汇总了示例中出现的全部参数及其含义参数示例值说明aws.s3.endpointhttp://minio:9000MinIO 服务地址需与你的 MinIO 部署保持一致paths3://starrocks/xxx.parquet待读取对象的完整路径格式为s3://bucket/objectaws.s3.enable_sslfalse是否启用 SSL 连接MinIO 使用 HTTPS 时需设为trueaws.s3.access_keyAAAAAAAAAAAAAAAAAAAAMinIO Access Key示例值需替换aws.s3.secret_keyBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBMinIO Secret Key示例值需替换formatparquet文件格式取值parquet、orc、csvcsv 自 v3.3.0 起aws.s3.use_aws_sdk_default_behaviorfalse是否使用 AWS SDK 默认凭据链aws.s3.use_instance_profilefalse是否使用实例配置文件Instance Profile凭据aws.s3.enable_path_style_accesstrue是否启用 S3 path-style 访问MinIO 等自建对象存储通常必须为true这些参数并非只是文档约定而是由 BE 端代码真实解析的。在 cloud_configuration_factory.h 中定义了完整的参数常量清单包括aws.s3.access_key、aws.s3.secret_key、aws.s3.session_token、aws.s3.iam_role_arn、aws.s3.region等在 cloud_configuration_factory.cpp 的CloudConfigurationFactory::create_aws()中这些属性被逐一读取并填入AWSCloudConfiguration结构体其中enable_path_style_access默认值为false而 MinIO 等 S3 兼容存储不提供虚拟主机风格的 DNS 解析因此示例中显式置为trueenable_ssl默认值为trueMinIO 走 HTTP 时需显式置为falseuse_aws_sdk_default_behavior与use_instance_profile默认均为false即默认走显式 Access Key 认证路径。典型示例一用 SELECT 直接查询 MinIO 数据在正式建表之前先用 SELECTFILES()直接预览数据集内容可以不落地存储数据即可查看数据集全貌查询字段的 min/max 值辅助决定目标表的数据类型检查是否存在NULL值。SELECT * FROM FILES ( aws.s3.endpoint http://minio:9000, path s3://starrocks/user_behavior_ten_million_rows.parquet, aws.s3.enable_ssl false, aws.s3.access_key AAAAAAAAAAAAAAAAAAAA, aws.s3.secret_key BBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBB, format parquet, aws.s3.use_aws_sdk_default_behavior false, aws.s3.use_instance_profile false, aws.s3.enable_path_style_access true ) LIMIT 3;系统返回如下查询结果注意列名由 Parquet 文件本身提供---------------------------------------------------------------- | UserID | ItemID | CategoryID | BehaviorType | Timestamp | ---------------------------------------------------------------- | 543711 | 829192 | 2355072 | pv | 2017-11-27 08:22:37 | | 543711 | 2056618 | 3645362 | pv | 2017-11-27 10:16:46 | | 543711 | 1165492 | 3645362 | pv | 2017-11-27 10:17:00 | ---------------------------------------------------------------- 3 rows in set (0.41 sec)典型示例二用 CTAS 建表并加载数据将上一个查询用 CREATE TABLE AS SELECTCTAS包装起来即可让 StarRocks 自动推断表结构、创建表并完成数据加载。由于 Parquet 格式自带列名使用FILES()时无需手动指定列名与类型。:::note CTAS 配合 schema inference 建表时无法直接设置副本数需在建表前通过 FE 配置指定。单副本系统的示例如下ADMIN SET FRONTEND CONFIG (default_replication_num 1);:::创建数据库并切换CREATE DATABASE IF NOT EXISTS mydatabase; USE mydatabase;执行 CTASCREATE TABLE user_behavior_inferred AS SELECT * FROM FILES ( aws.s3.endpoint http://minio:9000, path s3://starrocks/user_behavior_ten_million_rows.parquet, aws.s3.enable_ssl false, aws.s3.access_key AAAAAAAAAAAAAAAAAAAA, aws.s3.secret_key BBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBB, format parquet, aws.s3.use_aws_sdk_default_behavior false, aws.s3.use_instance_profile false, aws.s3.enable_path_style_access true );加载成功后返回如下信息1000 万行数据耗时约 3 秒Query OK, 10000000 rows affected (3.17 sec) {label:insert_a5da3ff5-9ee4-11ee-90b0-02420a060004, status:VISIBLE, txnId:17}用 DESCRIBE 查看自动推断出的表结构DESCRIBE user_behavior_inferred;------------------------------------------------------------- | Field | Type | Null | Key | Default | Extra | ------------------------------------------------------------- | UserID | bigint | YES | true | NULL | | | ItemID | bigint | YES | true | NULL | | | CategoryID | bigint | YES | true | NULL | | | BehaviorType | varchar(1048576) | YES | false | NULL | | | Timestamp | varchar(1048576) | YES | false | NULL | | -------------------------------------------------------------验证数据已成功加载SELECT * from user_behavior_inferred LIMIT 3;--------------------------------------------------------------- | UserID | ItemID | CategoryID | BehaviorType | Timestamp | --------------------------------------------------------------- | 58 | 158350 | 2355072 | pv | 2017-11-27 13:06:51 | | 58 | 158590 | 3194735 | pv | 2017-11-27 02:21:04 | | 58 | 215073 | 3002561 | pv | 2017-11-30 10:55:42 | ---------------------------------------------------------------典型示例三用 INSERT 加载进已存在表当你希望对目标表做更多定制时例如指定列数据类型、NULL 约束、默认值、Key 类型与列、数据分区分桶策略建议采用「手动建表 INSERT INTO SELECT FROM FILES()」的方式。关于如何设计高效的表结构可参考 StarRocks 表设计。基于对 Parquet 数据内容的了解可以做出如下设计决策通过查询确认Timestamp列的数据与datetime类型匹配故 DDL 中将其声明为datetime数据集中不存在NULL值故所有列均声明为NOT NULL结合预期查询模式将排序列与分桶列设为UserID你的场景也可以改用或追加ItemID。创建数据库并建表建议目标表结构与 Parquet 文件保持一致CREATE DATABASE IF NOT EXISTS mydatabase; USE mydatabase;CREATE TABLE user_behavior_declared ( UserID int(11) NOT NULL, ItemID int(11) NOT NULL, CategoryID int(11) NOT NULL, BehaviorType varchar(65533) NOT NULL, Timestamp datetime NOT NULL ) ENGINE OLAP DUPLICATE KEY(UserID) DISTRIBUTED BY HASH(UserID) PROPERTIES ( replication_num 1 );对比手动声明的表结构与前面FILES()推断出的结构可以直观看到差异DESCRIBE user_behavior_declared;----------------------------------------------------------- | Field | Type | Null | Key | Default | Extra | ----------------------------------------------------------- | UserID | int | NO | true | NULL | | | ItemID | int | NO | false | NULL | | | CategoryID | int | NO | false | NULL | | | BehaviorType | varchar(65533) | NO | false | NULL | | | Timestamp | datetime | NO | false | NULL | | ----------------------------------------------------------- 5 rows in set (0.00 sec)对比要点集中在三处数据类型、是否可空nullable、Key 字段。生产环境中建议手工指定表结构以更好地控制目标表 schema 并提升查询性能——例如将时间戳字段用datetime而非varchar存储查询与存储效率都更高。执行 INSERT INTO SELECT FROM FILES() 完成加载INSERT INTO user_behavior_declared SELECT * FROM FILES ( aws.s3.endpoint http://minio:9000, path s3://starrocks/user_behavior_ten_million_rows.parquet, aws.s3.enable_ssl false, aws.s3.access_key AAAAAAAAAAAAAAAAAAAA, aws.s3.secret_key BBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBB, format parquet, aws.s3.use_aws_sdk_default_behavior false, aws.s3.use_instance_profile false, aws.s3.enable_path_style_access true );查询验证加载结果SELECT * from user_behavior_declared LIMIT 3;---------------------------------------------------------------- | UserID | ItemID | CategoryID | BehaviorType | Timestamp | ---------------------------------------------------------------- | 58 | 4309692 | 1165503 | pv | 2017-11-25 14:06:52 | | 58 | 181489 | 1165503 | pv | 2017-11-25 14:07:22 | | 58 | 3722956 | 1165503 | pv | 2017-11-25 14:09:28 | ----------------------------------------------------------------查询 INSERT 加载进度从 v3.1 起可以通过 Information Schema 中的 loads 视图查询 INSERT 任务的进度SELECT * FROM information_schema.loads ORDER BY JOB_ID DESC;提交了多个加载任务时可按LABEL过滤定位特定任务SELECT * FROM information_schema.loads WHERE LABEL insert_e3b882f5-7eb3-11ee-ae77-00163e267b60 \G*************************** 1. row *************************** JOB_ID: 10243 LABEL: insert_e3b882f5-7eb3-11ee-ae77-00163e267b60 DATABASE_NAME: mydatabase STATE: FINISHED PROGRESS: ETL:100%; LOAD:100% TYPE: INSERT PRIORITY: NORMAL SCAN_ROWS: 10000000 FILTERED_ROWS: 0 UNSELECTED_ROWS: 0 SINK_ROWS: 10000000 ETL_INFO: TASK_INFO: resource:N/A; timeout(s):300; max_filter_ratio:0.0 CREATE_TIME: 2023-11-09 11:56:01 ETL_START_TIME: 2023-11-09 11:56:01 ETL_FINISH_TIME: 2023-11-09 11:56:01 LOAD_START_TIME: 2023-11-09 11:56:01 LOAD_FINISH_TIME: 2023-11-09 11:56:44 JOB_DETAILS: {All backends:{e3b882f5-7eb3-11ee-ae77-00163e267b60:[10142]},FileNumber:0,FileSize:0,InternalTableLoadBytes:311710786,InternalTableLoadRows:10000000,ScanBytes:581574034,ScanRows:10000000,TaskNumber:1,Unfinished backends:{e3b882f5-7eb3-11ee-ae77-00163e267b60:[]}} ERROR_MSG: NULL TRACKING_URL: NULL TRACKING_SQL: NULL REJECTED_RECORD_PATH: NULL:::tip INSERT 是同步命令。若 INSERT 任务仍在运行需要另开一个会话查看其执行状态。 :::对比两种表在磁盘上的大小下面这条查询用于对比「推断 schema 的表」与「手工声明 schema 的表」的存储差异。由于推断 schema 包含可空列且时间戳字段用了varchar其数据长度更大SELECT TABLE_NAME, TABLE_ROWS, AVG_ROW_LENGTH, DATA_LENGTH FROM information_schema.tables WHERE TABLE_NAME like user_behavior%\G*************************** 1. row *************************** TABLE_NAME: user_behavior_declared TABLE_ROWS: 10000000 AVG_ROW_LENGTH: 10 DATA_LENGTH: 102562516 *************************** 2. row *************************** TABLE_NAME: user_behavior_inferred TABLE_ROWS: 10000000 AVG_ROW_LENGTH: 17 DATA_LENGTH: 176803880 2 rows in set (0.04 sec)两表行数相同各 1000 万行但手工声明 schema 的表DATA_LENGTH约 102 MB而推断 schema 的表约 176 MB——相差约 70 MB这正是varchar时间戳与可空列带来的额外存储开销。为时间字段选择datetime等紧凑类型、按需设置 NOT NULL能显著降低存储成本并提升扫描性能。使用 Broker Load异步加载Broker Load 是异步加载方式由 StarRocks 在后台完成与 MinIO 的连接、数据拉取与落库客户端无需保持连接。它支持 Parquet、ORC、CSV以及自 v3.2.3 起支持的 JSON 格式。Broker Load 的优势后台运行任务提交后客户端无需持续在线适合长任务默认超时时间长达 4 小时可支撑大批量长时间加载格式更全除 Parquet、ORC 外还支持 CSV 与 JSONJSON 自 v3.2.3 起。数据流转过程Broker Load 的整体工作流如下图所示其执行链路分为三步用户创建加载任务LOAD LABEL ...前端 FE 生成查询计划并将计划分发到后端节点 BE 或计算节点 CNBE/CN 节点从数据源拉取数据并将数据写入 StarRocks。典型示例从 MinIO 加载 Parquet创建数据库与表CREATE DATABASE IF NOT EXISTS mydatabase; USE mydatabase;创建目标表建议与 Parquet 文件结构保持一致CREATE TABLE user_behavior ( UserID int(11) NOT NULL, ItemID int(11) NOT NULL, CategoryID int(11) NOT NULL, BehaviorType varchar(65533) NOT NULL, Timestamp datetime NOT NULL ) ENGINE OLAP DUPLICATE KEY(UserID) DISTRIBUTED BY HASH(UserID) PROPERTIES ( replication_num 1 );启动 Broker Load 任务LOAD LABEL UserBehavior ( DATA INFILE(s3://starrocks/user_behavior_ten_million_rows.parquet) INTO TABLE user_behavior ) WITH BROKER ( aws.s3.endpoint http://minio:9000, aws.s3.enable_ssl false, aws.s3.access_key AAAAAAAAAAAAAAAAAAAA, aws.s3.secret_key BBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBBB, aws.s3.use_aws_sdk_default_behavior false, aws.s3.use_instance_profile false, aws.s3.enable_path_style_access true ) PROPERTIES ( timeout 72000 );任务由四个主要部分组成LABEL任务的唯一标识字符串用于后续查询任务状态LOAD声明指定源数据 URI、源数据格式与目标表名BROKER源端连接信息MinIO 的 endpoint、凭据与认证模式PROPERTIES任务的超时时间等属性示例中timeout设为 72000 秒。BROKER段中的aws.s3.*参数与 INSERTFILES()中的同名参数含义完全一致均对应 BE 端 cloud_configuration_factory.h 中定义的配置键并由create_aws()解析为AWSCloudConfiguration后驱动 S3 客户端访问 MinIO。完整的语法与参数说明参见 BROKER LOAD 文档。查询 Broker Load 任务进度从 v3.1 起可通过 loads 视图查询 Broker Load 任务状态SELECT * FROM information_schema.loads;按LABEL过滤定位特定任务注意视图中小写展示 labelSELECT * FROM information_schema.loads WHERE LABEL UserBehavior\G*************************** 1. row *************************** JOB_ID: 10176 LABEL: userbehavior DATABASE_NAME: mydatabase STATE: FINISHED PROGRESS: ETL:100%; LOAD:100% TYPE: BROKER PRIORITY: NORMAL SCAN_ROWS: 10000000 FILTERED_ROWS: 0 UNSELECTED_ROWS: 0 SINK_ROWS: 10000000 ETL_INFO: TASK_INFO: resource:N/A; timeout(s):72000; max_filter_ratio:0.0 CREATE_TIME: 2023-12-19 23:02:41 ETL_START_TIME: 2023-12-19 23:02:44 ETL_FINISH_TIME: 2023-12-19 23:02:44 LOAD_START_TIME: 2023-12-19 23:02:44 LOAD_FINISH_TIME: 2023-12-19 23:02:46 JOB_DETAILS: {All backends:{4aeec563-a91e-4c1e-b169-977b660950d1:[10004]},FileNumber:1,FileSize:132251298,InternalTableLoadBytes:311710786,InternalTableLoadRows:10000000,ScanBytes:132251298,ScanRows:10000000,TaskNumber:1,Unfinished backends:{4aeec563-a91e-4c1e-b169-977b660950d1:[]}} ERROR_MSG: NULL TRACKING_URL: NULL TRACKING_SQL: NULL REJECTED_RECORD_PATH: NULL 1 row in set (0.02 sec)任务状态为FINISHED、SCAN_ROWS与SINK_ROWS均为 1000 万、FILTERED_ROWS为 0说明全部数据无过滤地完成加载。最后抽样验证目标表数据SELECT * from user_behavior LIMIT 3;---------------------------------------------------------------- | UserID | ItemID | CategoryID | BehaviorType | Timestamp | ---------------------------------------------------------------- | 142 | 2869980 | 2939262 | pv | 2017-11-25 03:43:22 | | 142 | 2522236 | 1669167 | pv | 2017-11-25 15:14:12 | | 142 | 3031639 | 3607361 | pv | 2017-11-25 15:19:25 | ----------------------------------------------------------------两种方式的对比与选型建议维度INSERT FILES()Broker Load同步/异步同步命令返回即完成或失败异步后台持续执行默认超时约 300 秒timeout(s):3004 小时可自行调整示例设 72000 秒文件格式Parquet、ORC、CSVv3.3.0 起Parquet、ORC、CSV、JSONv3.2.3 起适用场景中小规模数据、快速预览、CTAS 建表、与 SELECT 联用大批量、长耗时任务或需要 DELETE 等数据变更客户端要求需保持连接等待结果提交后即可断开后台继续实际工程中可按以下思路决策能选 INSERTFILES()就优先选它——语法简单、天然与 SQL 生态融合可先用 SELECT 预览数据、再用 CTAS 或 INSERT 落库需要加载 JSON、或文件格式不在 Parquet/ORC/CSV 之列时切换到 Broker Load任务量大、执行时间长、客户端不便长连接时Broker Load 的异步机制与 4 小时默认超时更稳妥无论哪种方式aws.s3.enable_path_style_access、aws.s3.enable_ssl、aws.s3.use_aws_sdk_default_behavior、aws.s3.use_instance_profile这组认证/访问参数都是连接 MinIO 的关键务必结合你的 MinIO 部署方式HTTP/HTTPS、路径风格正确设置。延伸阅读INSERT 语句参考BROKER LOAD 语句参考FILES()表函数参考loads 信息视图参考StarRocks 表设计指南Primary Key 表加载与数据变更【免费下载链接】starrocksThe worlds fastest open query engine for sub-second analytics both on and off the data lakehouse. With the flexibility to support nearly any scenario, StarRocks provides best-in-class performance for multi-dimensional analytics, real-time analytics, and ad-hoc queries. A Linux Foundation project.项目地址: https://gitcode.com/GitHub_Trending/st/starrocks创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
