StarRocks 加载时数据转换(ETL in Loading)完整指南:列映射、行过滤与派生列生成 📅 发布时间:2026/9/17 3:07:12 👁 浏览次数: StarRocks 加载时数据转换ETL in Loading完整指南列映射、行过滤与派生列生成【免费下载链接】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/starrocksStarRocks 提供了在**数据加载过程中直接完成 ETL提取与转换**的能力无需在上游预先清洗数据。本文以 CSV 数据为例系统讲解如何通过 Stream Load、Broker Load 与 Routine Load 三种加载方式实现列跳过与列映射、WHERE 行过滤、基于函数/表达式的派生列生成以及从 Hive 文件路径中提取分区字段值并附上源码级的实现原理说明与可复现的完整示例。功能概述与适用范围加载时转换Transform data at loading是 StarRocks 原生支持的一项能力当你把数据文件载入 StarRocks 表时如果数据文件的列与目标表的列无法完全一一对应无需在加载前单独做抽取或转换StarRocks 可以在加载过程中直接完成抽取与转换。该功能支持以下三种加载方式Stream Load从本地文件系统或流式数据源同步加载提交后立即返回作业结果Broker Load从 HDFS 及云存储异步加载作业结果通过SELECT * FROM information_schema.loads查询v3.1 起Routine Load从 Kafka 持续消费并加载数据。不支持Spark Load。权限说明只有对目标表拥有 INSERT 权限的用户才能执行加载。授权语法为GRANT INSERT ON TABLE table_name IN DATABASE database_name TO { ROLE role_name | USER user_identity}参见 GRANT。CSV 分隔符约定对于 CSV 数据可以使用长度不超过 50 字节的 UTF-8 字符串作为文本分隔符例如逗号,、制表符或竖线|。此外各加载方式支持的数据文件格式不同Stream Load 支持 CSV/JSONBroker Load 支持 CSV/Parquet/ORC/JSONRoutine Load 支持 CSV/JSON/Avro请按所选方式确定可用的数据格式。四种典型转换场景当数据文件与 StarRocks 表的列无法完全对应时加载时转换可以覆盖以下四类需求跳过不需要加载的列数据文件中存在无法映射到目标表的列时可直接忽略若数据文件的列顺序与目标表不同可建立列映射。过滤掉不想加载的行通过指定过滤条件让 StarRocks 只加载满足条件的行。从原始列生成新列生成列generated columns是依据数据文件原始列计算出来的特殊列可映射到目标表的列。从文件路径提取分区字段值当数据文件由 Apache Hive 生成时可以从文件路径中提取分区字段值。数据与表结构示例为便于对照后续各场景先在本地创建两个数据文件。创建file1.csv包含四列依次表示用户 ID、用户性别、事件日期、事件类型354,female,2020-05-20,1 465,male,2020-05-21,2 576,female,2020-05-22,1 687,male,2020-05-23,2创建file2.csv仅包含一列表示日期2020-05-20 2020-05-21 2020-05-22 2020-05-23在 StarRocks 数据库test_db中创建两张目标表。创建table1包含三列event_date、event_type、user_idMySQL [test_db] CREATE TABLE table1 ( event_date DATE COMMENT event date, event_type TINYINT COMMENT event type, user_id BIGINT COMMENT user ID ) DISTRIBUTED BY HASH(user_id);创建table2包含四列date、year、month、dayMySQL [test_db] CREATE TABLE table2 ( date DATE COMMENT date, year INT COMMENT year, month TINYINT COMMENT month, day TINYINT COMMENT day ) DISTRIBUTED BY HASH(date);说明自 v2.5.7 起StarRocks 在创建表或添加分区时可自动设置分桶数BUCKETS无需再手动指定详见 数据分布。随后做如下准备将file1.csv与file2.csv上传到 HDFS 集群的/user/starrocks/data/input/路径将file1.csv的数据发布到 Kafka 集群的topic1将file2.csv的数据发布到topic2。跳过不需要加载的列数据文件中可能包含无法映射到 StarRocks 表任何列的列此时只需加载可映射的列即可。该能力支持本地文件系统、HDFS 与云存储下文以 HDFS 为例以及 Kafka 三类数据源。由于多数 CSV 文件的列没有列名即使首行是列名StarRocks 也会将其作为普通数据处理因此必须按顺序在作业创建语句或命令中临时命名 CSV 的列这些临时命名的列再按名称映射到目标表的列。规则如下能映射到目标表列即临时命名与目标表列名一致的数据会被直接加载无法映射到目标表列的列会被忽略其数据不会加载如果某列本可映射到目标表列但未在语句或命令中临时命名加载作业会报错。以file1.csv与table1为例file1.csv的四列依次临时命名为user_id、user_gender、event_date、event_type。其中user_id、event_date、event_type可映射到table1的列而user_gender无法映射因此加载时user_gender被跳过。从本地文件系统加载Stream Loadcurl --location-trusted -u username:password \ -H Expect:100-continue \ -H column_separator:, \ -H columns: user_id, user_gender, event_date, event_type \ -T file1.csv -XPUT \ http://fe_host:fe_http_port/api/test_db/table1/_stream_load注意使用 Stream Load 时必须通过columns参数临时命名数据文件的列以建立数据文件与目标表之间的列映射。完整语法与参数说明见 STREAM LOAD。从 HDFS 集群加载Broker LoadLOAD LABEL test_db.label1 ( DATA INFILE(hdfs://hdfs_host:hdfs_port/user/starrocks/data/input/file1.csv) INTO TABLE table1 FORMAT AS csv COLUMNS TERMINATED BY , (user_id, user_gender, event_date, event_type) ) WITH BROKER;注意使用 Broker Load 时必须通过column_list参数临时命名数据文件的列。column_list中声明的列按名称映射到目标表列如果数据文件列与目标表列按顺序一一对应则无需指定column_list。完整语法见 BROKER LOAD。从 Kafka 集群加载Routine LoadCREATE ROUTINE LOAD test_db.table101 ON table1 COLUMNS TERMINATED BY ,, COLUMNS(user_id, user_gender, event_date, event_type) FROM KAFKA ( kafka_broker_list kafka_broker_host:kafka_broker_port, kafka_topic topic1, property.kafka_default_offsets OFFSET_BEGINNING );注意使用 Routine Load 时必须通过COLUMNS参数临时命名数据文件的列。property.kafka_default_offsets用于指定所有消费分区的默认起始偏移OFFSET_BEGINNING表示从最早偏移开始消费详见 CREATE ROUTINE LOAD。验证加载结果MySQL [test_db] SELECT * FROM table1; --------------------------------- | event_date | event_type | user_id | --------------------------------- | 2020-05-22 | 1 | 576 | | 2020-05-20 | 1 | 354 | | 2020-05-21 | 2 | 465 | | 2020-05-23 | 2 | 687 | --------------------------------- 4 rows in set (0.01 sec)可以看到 4 行数据全部入库user_gender列被成功跳过。过滤掉不想加载的行WHERE 子句当数据文件中的某些行无需加载时可以在作业创建语句或命令中使用 WHERE 子句指定过滤条件StarRocks 会过滤掉不满足条件的行。该能力同样支持本地文件系统、HDFS 与云存储以及 Kafka 三类数据源。以file1.csv与table1为例只加载事件类型event_type为1的行即指定过滤条件event_type 1。Stream Load本地文件系统curl --location-trusted -u username:password \ -H Expect:100-continue \ -H column_separator:, \ -H columns: user_id, user_gender, event_date, event_type \ -H where: event_type1 \ -T file1.csv -XPUT \ http://fe_host:fe_http_port/api/test_db/table1/_stream_load提示即便数据文件与目标表列数相同且按顺序映射只要需要按列做过滤也必须先用columns参数为数据文件列定义临时名称例如 STREAM LOAD 的opt_properties中的where参数STREAM LOAD。where过滤发生在数据预处理之后且被 WHERE 过滤掉的行不计入max_filter_ratio的容错统计。Broker LoadHDFSLOAD LABEL test_db.label2 ( DATA INFILE(hdfs://hdfs_host:hdfs_port/user/starrocks/data/input/file1.csv) INTO TABLE table1 FORMAT AS csv COLUMNS TERMINATED BY , (user_id, user_gender, event_date, event_type) WHERE event_type 1 ) WITH BROKER;WHERE位于data_desc描述符中指定基于源数据包括column_list中定义的列或 SET 生成的列的过滤条件详见 BROKER LOAD。Routine LoadKafkaCREATE ROUTINE LOAD test_db.table102 ON table1 COLUMNS TERMINATED BY ,, COLUMNS (user_id, user_gender, event_date, event_type), WHERE event_type 1 FROM KAFKA ( kafka_broker_list kafka_broker_host:kafka_broker_port, kafka_topic topic1, property.kafka_default_offsets OFFSET_BEGINNING );注意Routine Load 的 WHERE 条件中列可以是源数据列也可以是派生列且被 WHERE 过滤掉的行同样不计入max_error_number与max_filter_ratio的统计详见 CREATE ROUTINE LOAD。验证加载结果MySQL [test_db] SELECT * FROM table1; --------------------------------- | event_date | event_type | user_id | --------------------------------- | 2020-05-20 | 1 | 354 | | 2020-05-22 | 1 | 576 | --------------------------------- 2 rows in set (0.01 sec)仅事件类型为1的两行数据被加载。从原始列生成新列函数与表达式当数据文件中的某些数据在载入目标表前需要转换时可以在作业创建命令或语句中使用函数或表达式实现数据转换即生成列。该能力支持本地文件系统、HDFS 与云存储以及 Kafka 三类数据源。以file2.csv与table2为例file2.csv只有一列日期数据。可使用 year、month、day 函数分别提取年、月、日并加载到table2的year、month、day列。Stream Load本地文件系统curl --location-trusted -u username:password \ -H Expect:100-continue \ -H column_separator:, \ -H columns:date,yearyear(date),monthmonth(date),dayday(date) \ -T file2.csv -XPUT \ http://fe_host:fe_http_port/api/test_db/table2/_stream_load注意在columns参数中必须先临时命名数据文件的全部原始列再临时命名由原始列生成的新列。上例中先将file2.csv的唯一一列临时命名为date再通过yearyear(date)、monthmonth(date)、dayday(date)生成三个新列。当数据文件与目标表列数相同但顺序不同、且无需函数计算时只需在columns中按数据文件列顺序写出目标表列名即可当列数不同或需要函数计算时则需临时命名并引用表达式例如columns: col1, col2, col3, temptemp为被跳过的第 4 列临时名或columns: col, year year(col), monthmonth(col), dayday(col)详见 STREAM LOAD 的列映射。Broker LoadHDFSLOAD LABEL test_db.label3 ( DATA INFILE(hdfs://hdfs_host:hdfs_port/user/starrocks/data/input/file2.csv) INTO TABLE table2 FORMAT AS csv COLUMNS TERMINATED BY , (date) SET(yearyear(date), monthmonth(date), dayday(date)) ) WITH BROKER;注意使用 Broker Load 时必须先用column_list参数临时命名数据文件的全部列再用 SET 子句临时命名由原始列生成的新列。上例中先在column_list中将唯一列临时命名为date再在 SET 子句中调用yearyear(date)、monthmonth(date)、dayday(date)生成三个新列。SET 子句的典型用法还包括对多列求和例如column_list声明(col1,col2,tmp_col3,tmp_col4)SET 子句写(col3tmp_col3tmp_col4)实现数据转换详见 BROKER LOAD。Routine LoadKafkaCREATE ROUTINE LOAD test_db.table201 ON table2 COLUMNS TERMINATED BY ,, COLUMNS(date,yearyear(date),monthmonth(date),dayday(date)) FROM KAFKA ( kafka_broker_list kafka_broker_host:kafka_broker_port, kafka_topic topic2, property.kafka_default_offsets OFFSET_BEGINNING );注意在 Routine Load 的COLUMNS参数中同样必须先临时命名数据文件的全部列再临时命名生成列。Routine Load 将COLUMNS中的列分为两类直接映射列column_name与派生列column_assignment即column_name expr。建议将派生列放在映射列之后因为 StarRocks 先解析映射列详见 CREATE ROUTINE LOAD。验证加载结果MySQL [test_db] SELECT * FROM table2; ------------------------------- | date | year | month | day | ------------------------------- | 2020-05-20 | 2020 | 5 | 20 | | 2020-05-21 | 2020 | 5 | 21 | | 2020-05-22 | 2020 | 5 | 22 | | 2020-05-23 | 2020 | 5 | 23 | ------------------------------- 4 rows in set (0.01 sec)日期列被成功拆分为年、月、日三列。从文件路径提取分区字段值COLUMNS FROM PATH AS当指定的文件路径中包含分区字段时可以使用COLUMNS FROM PATH AS参数从文件路径中提取分区字段。路径中的分区字段等价于数据文件中的列。注意该参数仅在从 HDFS 集群加载数据时受支持。例如需要加载由 Hive 生成的以下四个数据文件它们存放在 HDFS 的/user/starrocks/data/input/路径下按分区字段date分区每个文件包含两列依次表示事件类型和用户 ID/user/starrocks/data/input/date2020-05-20/data 1,354 /user/starrocks/data/input/date2020-05-21/data 2,465 /user/starrocks/data/input/date2020-05-22/data 1,576 /user/starrocks/data/input/date2020-05-23/data 2,687从 HDFS 集群加载Broker Load执行以下语句创建 Broker Load 作业从/user/starrocks/data/input/文件路径中提取date分区字段值并使用通配符*加载该路径下的所有数据文件到table1LOAD LABEL test_db.label4 ( DATA INFILE(hdfs://fe_host:fe_http_port/user/starrocks/data/input/date*/*) INTO TABLE table1 FORMAT AS csv COLUMNS TERMINATED BY , (event_type, user_id) COLUMNS FROM PATH AS (date) SET(event_date date) ) WITH BROKER;注意上例中文件路径中的date分区字段等价于table1的event_date列因此需要使用 SET 子句将date分区字段映射到event_date列如果路径中的分区字段与目标表列同名则无需 SET 子句建立映射。DATA INFILE的文件路径支持通配符?、*、[]、{}或^可用于中间路径例如hdfs://hdfs_host:hdfs_port/user/data/tablename/dt202104*/*可加载指定月份所有分区的文件详见 BROKER LOAD。验证加载结果MySQL [test_db] SELECT * FROM table1; --------------------------------- | event_date | event_type | user_id | --------------------------------- | 2020-05-22 | 1 | 576 | | 2020-05-20 | 1 | 354 | | 2020-05-21 | 2 | 465 | | 2020-05-23 | 2 | 687 | --------------------------------- 4 rows in set (0.01 sec)分区字段date的值被成功提取并加载到event_date列。底层实现原理列映射与路径分区填充上述能力的落地依托于 BEBackend侧的文件扫描器File Scanner实现。在 be/src/connector/file/scanner/file_scanner.cpp 中FileScanner::fill_columns_from_path()负责将文件路径中的分区值填充到扫描得到的 Chunk 中函数遍历columns_from_path列表把每个分区字段对应的值写入对应 slot 位置。也就是说从路径提取分区字段在 BE 端被实现为一种特殊的列填充操作——分区字段在读取时被当作数据列追加到文件字段之后。在 be/src/connector/file/scanner/csv_scanner.cpp 中还可以看到列数的强校验逻辑同一个加载作业内各文件范围FileRange的columns_from_path数量必须一致且必须满足num_of_columns_from_file columns_from_path.size() _src_slot_descriptors.size()即文件本身字段数 路径分区字段数必须等于目标 slot 描述符的数量。这从实现层面印证了本文第四节的规则数据文件列含临时命名与生成列与路径分区列合并后必须能完整对应目标表的列定义若临时命名缺失或数量不匹配加载作业便会报错。CSV 扫描器在处理每个文件时调用fill_columns_from_path()见 csv_scanner.cpp将分区值按序填入数据块随后再经过 WHERE 过滤与 SET 表达式求值完成整条加载时 ETL流水线。此外参数侧同样有据可查Stream Load 的columns、where、max_filter_ratio等参数位于 HTTP 请求头where用于对预处理后的数据过滤max_filter_ratio表示可容忍的数据质量问题行占比0~1默认 0被 WHERE 过滤的行不计入其中详见 STREAM LOADBroker Load 的data_desc中COLUMNS FROM PATH AS明确标注仅当从 HDFS 加载数据时可用SET子句支持任意函数转换WHERE指定源数据过滤条件详见 BROKER LOADRoutine Load 的COLUMNS区分映射列与派生列WHERE条件可引用源列或派生列详见 CREATE ROUTINE LOAD。实战建议与注意事项区分加载方式与数据源本地小文件优先使用 Stream Load单文件建议不超过 10 GBHDFS/云存储大文件或海量文件使用 Broker LoadKafka 持续流式数据使用 Routine Load。各方式的具体格式支持与适用场景可参考 Stream Load 指南、HDFS 加载指南 与 Routine Load 指南。临时命名是列映射的前提CSV 列无名三种加载方式都必须按顺序临时命名能映射的列直接加载、不能映射的列被忽略、应当映射却未命名的列会报错。WHERE 过滤与容错统计分离被 WHERE 过滤掉的行不属于数据质量问题不计入max_filter_ratioStream Load/Broker Load或max_error_number/max_filter_ratioRoutine Load。若作业因脏数据失败可结合log_rejected_record_num与返回的ErrorURLStream Load或information_schema.loadsBroker Load定位问题行。派生列顺序在columns/COLUMNS/SET 中遵循先原始列、后生成列的书写顺序Routine Load 建议将派生列放在映射列之后。分区路径提取仅 HDFS 场景可用COLUMNS FROM PATH AS路径分区字段与目标表列同名时无需 SET 映射不同名时务必用 SET 建立映射否则分区值无法正确入库。综上StarRocks 的加载时数据转换将 ETL 环节内嵌到数据导入流水线中配合列映射、行过滤、函数派生与分区路径提取四项能力可以显著简化数据接入链路减少上游数据预处理的额外开发与运维成本。【免费下载链接】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),仅供参考