ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

使用 StarRocks Spark Connector 将数据批量加载至 StarRocks(Stream Load 指南)

2026/9/17 1:26:58 拓冰建站 浏览量
使用 StarRocks Spark Connector 将数据批量加载至 StarRocks(Stream Load 指南) 使用 StarRocks Spark Connector 将数据批量加载至 StarRocksStream 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/starrocksStarRocks 为 Apache Spark™ 提供自研的官方连接器 StarRocks Connector for Apache Spark下文简称 Spark connector其核心原理是先在内存量累积数据再通过 STREAM LOAD 一次性将整批数据写入 StarRocks。该连接器基于 Spark DataSource V2 实现可通过 Spark DataFrame 或 Spark SQL 两种方式创建数据源同时支持 Batch批式与 Structured Streaming结构化流式两种写入模式。读完本文你将掌握 Spark connector 的版本选型、获取方式、全部写入参数、类型映射规则以及主键表部分更新/条件更新、BITMAP/HLL/ARRAY/STRUCT/MAP 等复杂类型的实战加载方案。版本要求Spark connector 与 Spark、StarRocks、Java、Scala 的版本兼容关系如下Spark connectorSparkStarRocksJavaScala1.1.44.0, 4.12.5 及以上172.131.1.43.3, 3.4, 3.52.5 及以上82.121.1.33.2, 3.3, 3.4, 3.52.5 及以上82.121.1.23.2, 3.3, 3.4, 3.52.5 及以上82.121.1.13.2, 3.3 或 3.42.5 及以上82.121.1.03.2, 3.3 或 3.42.5 及以上82.12需要注意的兼容性事项不同版本 Spark connector 之间的行为差异参见下文 升级 Spark connector 一节。自 1.1.1 起Spark connector 不再内置 MySQL JDBC 驱动mysql-connector-java你需要手动将驱动加入 Spark classpath可从 MySQL 官网或 Maven Central 获取。通常 Spark connector 的最新版本仅维护对最近三个 Spark 大版本的兼容性。获取 Spark connectorSpark connector 的 JAR 文件命名格式为starrocks-spark-connector-${spark_version}_${scala_version}-${connector_version}.jar。例如环境安装 Spark 3.2 与 Scala 2.12使用 connector 1.1.0对应 JAR 为starrocks-spark-connector-3.2_2.12-1.1.0.jar。可通过以下三种方式获取直接下载已编译的 JAR从 Maven Central Repository 的com.starrocks组下下载对应版本。Maven 依赖引入在pom.xml中添加依赖将spark_version、scala_version、connector_version替换为实际版本dependency groupIdcom.starrocks/groupId artifactIdstarrocks-spark-connector-${spark_version}_${scala_version}/artifactId version${connector_version}/version /dependency例如 Spark 3.2、Scala 2.12、connector 1.1.0dependency groupIdcom.starrocks/groupId artifactIdstarrocks-spark-connector-3.2_2.12/artifactId version1.1.0/version /dependency自行编译源码下载 Spark connector 源码包StarRocks 官方 Apache Spark 连接器仓库执行sh build.sh spark_version例如 Spark 版本为 3.2sh build.sh 3.2编译完成后在target/目录下找到生成的 JAR如starrocks-spark-connector-3.2_2.12-1.1.0-SNAPSHOT.jar。注意未正式发布的 connector 名称带有SNAPSHOT后缀。参数详解Spark connector 的全部配置均通过df.write.option(...)DataFrame 场景或OPTIONS(...)Spark SQL 场景传入下表为各参数的用途、必填性与默认值。必填参数参数必填默认值说明starrocks.fe.http.url是无StarRocks 集群 FE 的 HTTP 地址多个地址用逗号,分隔格式为fe_host1:fe_http_port1,fe_host2:fe_http_port2。自 1.1.1 起支持在地址前加http://前缀starrocks.fe.jdbc.url是无连接 FE MySQL 服务的地址格式为jdbc:mysql://fe_host:fe_query_portstarrocks.table.identifier是无StarRocks 表名格式为database_name.table_namestarrocks.user是无StarRocks 集群账号用户名该用户需要具备目标表的 SELECT 与 INSERT 权限starrocks.password是无集群账号密码写入行为参数参数必填默认值说明starrocks.write.label.prefix否spark-Stream Load 使用的 label 前缀starrocks.write.enable.transaction-stream-load否TRUE是否使用 Stream Load 事务接口 加载数据要求 StarRocks v2.5 及以上。该特性可在单个事务中加载更多数据、占用更少内存并提升性能starrocks.write.buffer.size否104857600内存中累积的数据达到该大小字节后一次性发送至 StarRocks。值越大加载性能越高但会增大加载延迟starrocks.write.buffer.rows否Integer.MAX_VALUE自 1.1.1 起内存中累积的行数达到该值后一次性发送至 StarRocksstarrocks.write.flush.interval.ms否300000数据发送至 StarRocks 的间隔时间毫秒用于控制加载延迟starrocks.write.max.retries否3自 1.1.1 起同一批数据 Stream Load 失败后的重试次数starrocks.write.retry.interval.ms否10000自 1.1.1 起同一批数据 Stream Load 失败后的重试间隔毫秒starrocks.write.use_bitmap_hash64否false自 1.1.3 起是否使用 64 位哈希函数生成 bitmap默认使用 32 位哈希函数starrocks.write.num.partitions否无Spark 并行写入的分区数。数据量小时可调小分区数以降低加载并发与频率默认值由 Spark 决定但此方式可能带来 Spark Shuffle 开销starrocks.write.partition.columns否无Spark 中的分区列仅在指定了starrocks.write.num.partitions时生效未指定时使用所有待写入列参与分区starrocks.timezone否JVM 默认时区自 1.1.1 起用于将 SparkTimestampType转换为 StarRocksDATETIME的时区默认取ZoneId#systemDefault()。格式可为时区名如Asia/Shanghai或偏移量如08:00starrocks.write.socket.timeout.ms否-1自 1.1.3 起HTTP 客户端等待数据的时长毫秒-1表示不超时关于两个事务接口参数的联动自 1.1.1 起starrocks.write.enable.transaction-stream-load仅在starrocks.write.max.retries为非正值时生效因为 Stream Load 事务接口不支持重试。若starrocks.write.max.retries为正数连接器始终使用普通 Stream Load 接口并忽略前者取值。列与格式参数参数必填默认值说明starrocks.columns否无要加载数据的 StarRocks 表列多列用逗号分隔如col0,col1,col2starrocks.column.types否无自 1.1.1 起自定义 Spark 侧列数据类型替代从 StarRocks 表与默认映射推断的类型。取值为 SparkStructType#toDDL输出同款 DDL 格式 schema如col0 INT, col1 STRING, col2 BIGINT只需指定需要自定义的列。典型场景是向 BITMAP 或 HLL 列加载数据starrocks.write.properties.*否无控制 Stream Load 行为的参数例如starrocks.write.properties.format指定加载数据格式CSV 或 JSON完整参数清单见 STREAM LOADstarrocks.write.properties.format否CSVconnector 在发送数据前将每批数据转换成的文件格式合法值CSV、JSONstarrocks.write.properties.row_delimiter否\nCSV 格式数据的行分隔符starrocks.write.properties.column_separator否\tCSV 格式数据的列分隔符starrocks.write.properties.partial_update否FALSE是否启用部分更新合法值TRUE/FALSEstarrocks.write.properties.partial_update_mode否row部分更新模式合法值row与columnrow默认适用于列多、批量小的实时更新场景column适用于列少、行多的批量更新场景例如 100 列的表仅更新其中 10% 的列时column 模式更新速度约为 row 模式的 10 倍starrocks.write.properties.compression否无自 1.1.3 起Stream Load 使用的压缩算法合法值lz4_frame。JSON 格式压缩要求数据库 v3.2.7 及以上CSV 格式压缩无数据库版本要求Spark 与 StarRocks 的数据类型映射默认类型映射如下Spark 数据类型StarRocks 数据类型BooleanTypeBOOLEANByteTypeTINYINTShortTypeSMALLINTIntegerTypeINTLongTypeBIGINTStringTypeLARGEINTFloatTypeFLOATDoubleTypeDOUBLEDecimalTypeDECIMALStringTypeCHARStringTypeVARCHARStringTypeSTRINGStringTypeJSONDateTypeDATETimestampTypeDATETIMEArrayTypeARRAY自 1.1.1 起支持详见加载 ARRAY 类型列MapTypeMAP自 1.1.3 起支持详见加载嵌套列StructTypeSTRUCT自 1.1.3 起支持详见加载嵌套列同时支持自定义映射。例如 StarRocks 表包含 BITMAP、HLL 列而 Spark 不支持这两种类型时需要在 Spark 侧自定义对应列类型参见 BITMAP 与 HLL两者均自 1.1.1 起支持。升级 Spark connector从 1.1.0 升级到 1.1.1由于mysql-connector-java采用 GPL 许可证的限制自 1.1.1 起 Spark connector 不再内置该 MySQL JDBC 驱动。但 connector 仍需要通过 MySQL JDBC 驱动连接 StarRocks 获取表元数据因此你需要手动将驱动加入 Spark classpath。自 1.1.1 起connector 默认使用普通 Stream Load 接口而 1.1.0 默认使用 Stream Load 事务接口。若仍想使用事务接口可将starrocks.write.max.retries设为0详见starrocks.write.enable.transaction-stream-load与starrocks.write.max.retries的说明。基础示例以 DataFrame 与 Spark SQL 加载数据准备工作1. 创建 StarRocks 表。创建数据库test与主键表score_boardCREATE DATABASE test; CREATE TABLE test.score_board ( id int(11) NOT NULL COMMENT , name varchar(65533) NULL DEFAULT COMMENT , score int(11) NOT NULL DEFAULT 0 COMMENT ) ENGINEOLAP PRIMARY KEY(id) COMMENT OLAP DISTRIBUTED BY HASH(id);2. 网络配置。确保 Spark 所在机器能通过 FE 的http_port默认8030与query_port默认9030访问 FE 节点并能通过 BE 的be_http_port默认8040访问 BE 节点。这些默认端口均可在仓库配置文件中确认例如 conf/fe.conf 中定义了http_port 8030、query_port 9030conf/be.conf 中定义了be_http_port 8040。3. 搭建 Spark 环境。以下示例基于 Spark 3.2.4使用spark-shell、pyspark与spark-sql。运行前请将 Spark connector JAR 放入$SPARK_HOME/jars目录。使用 Spark DataFrame 加载Batch 模式在内存中构造数据并加载进 StarRocks 表。Scalaspark-shell// 1. 从序列创建 DataFrame val data Seq((1, starrocks, 100), (2, spark, 100)) val df data.toDF(id, name, score) // 2. 配置 format 为 starrocks 并设置如下选项写入 StarRocks // 请根据实际环境修改选项值 df.write.format(starrocks) .option(starrocks.fe.http.url, 127.0.0.1:8030) .option(starrocks.fe.jdbc.url, jdbc:mysql://127.0.0.1:9030) .option(starrocks.table.identifier, test.score_board) .option(starrocks.user, root) .option(starrocks.password, ) .mode(append) .save()Pythonpysparkfrom pyspark.sql import SparkSession spark SparkSession \ .builder \ .appName(StarRocks Example) \ .getOrCreate() # 1. 从序列创建 DataFrame data [(1, starrocks, 100), (2, spark, 100)] df spark.sparkContext.parallelize(data) \ .toDF([id, name, score]) # 2. 配置 format 为 starrocks 并设置如下选项写入 StarRocks # 请根据实际环境修改选项值 df.write.format(starrocks) \ .option(starrocks.fe.http.url, 127.0.0.1:8030) \ .option(starrocks.fe.jdbc.url, jdbc:mysql://127.0.0.1:9030) \ .option(starrocks.table.identifier, test.score_board) \ .option(starrocks.user, root) \ .option(starrocks.password, ) \ .mode(append) \ .save()查询验证MySQL [test] SELECT * FROM score_board; ------------------------ | id | name | score | ------------------------ | 1 | starrocks | 100 | | 2 | spark | 100 | ------------------------ 2 rows in set (0.00 sec)使用 Spark DataFrame 加载Structured Streaming 模式构造从 CSV 文件流式读取数据并加载进 StarRocks 表的任务。1. 准备数据。在目录csv-data下创建test.csv3,starrocks,100 4,spark,1002. 编写应用。Scalaspark-shellimport org.apache.spark.sql.types.StructType // 1. 从 CSV 创建 DataFrame val schema (new StructType() .add(id, integer) .add(name, string) .add(score, integer) ) val df (spark.readStream .option(sep, ,) .schema(schema) .format(csv) // 替换为你的 csv-data 目录路径 .load(/path/to/csv-data) ) // 2. 配置 format 为 starrocks 并设置如下选项写入 StarRocks // 请根据实际环境修改选项值 val query (df.writeStream.format(starrocks) .option(starrocks.fe.http.url, 127.0.0.1:8030) .option(starrocks.fe.jdbc.url, jdbc:mysql://127.0.0.1:9030) .option(starrocks.table.identifier, test.score_board) .option(starrocks.user, root) .option(starrocks.password, ) // 替换为你的 checkpoint 目录 .option(checkpointLocation, /path/to/checkpoint) .outputMode(append) .start() )Pythonpysparkfrom pyspark.sql import SparkSession from pyspark.sql.types import IntegerType, StringType, StructType, StructField spark SparkSession \ .builder \ .appName(StarRocks SS Example) \ .getOrCreate() # 1. 从 CSV 创建 DataFrame schema StructType([ StructField(id, IntegerType()), StructField(name, StringType()), StructField(score, IntegerType()) ]) df ( spark.readStream .option(sep, ,) .schema(schema) .format(csv) # 替换为你的 csv-data 目录路径 .load(/path/to/csv-data) ) # 2. 配置 format 为 starrocks 并设置如下选项写入 StarRocks # 请根据实际环境修改选项值 query ( df.writeStream.format(starrocks) .option(starrocks.fe.http.url, 127.0.0.1:8030) .option(starrocks.fe.jdbc.url, jdbc:mysql://127.0.0.1:9030) .option(starrocks.table.identifier, test.score_board) .option(starrocks.user, root) .option(starrocks.password, ) # 替换为你的 checkpoint 目录 .option(checkpointLocation, /path/to/checkpoint) .outputMode(append) .start() )3. 查询验证MySQL [test] select * from score_board; ------------------------ | id | name | score | ------------------------ | 4 | spark | 100 | | 3 | starrocks | 100 | ------------------------ 2 rows in set (0.67 sec)使用 Spark SQL 加载在spark-sql中通过INSERT INTO语句写入-- 1. 配置数据源为 starrocks 并设置如下选项建表 -- 请根据实际环境修改选项值 CREATE TABLE score_board USING starrocks OPTIONS( starrocks.fe.http.url127.0.0.1:8030, starrocks.fe.jdbc.urljdbc:mysql://127.0.0.1:9030, starrocks.table.identifiertest.score_board, starrocks.userroot, starrocks.password ); -- 2. 向表中插入两行数据 INSERT INTO score_board VALUES (5, starrocks, 100), (6, spark, 100);查询验证MySQL [test] select * from score_board; ------------------------ | id | name | score | ------------------------ | 6 | spark | 100 | | 5 | starrocks | 100 | ------------------------ 2 rows in set (0.00 sec)最佳实践向主键表加载数据部分更新与条件更新本部分展示如何向 StarRocks 主键表加载数据以实现部分更新与条件更新功能详细介绍见 通过加载变更数据示例使用 Spark SQL。准备创建数据库test与主键表score_boardDDL 与上文一致。部分更新——仅更新name列在 MySQL 客户端插入初始数据mysql INSERT INTO score_board VALUES (1, starrocks, 100), (2, spark, 100);在 Spark SQL 客户端创建 Spark 表将starrocks.write.properties.partial_update设为true开启部分更新并将starrocks.columns设为id,name指定写入列CREATE TABLE score_board USING starrocks OPTIONS( starrocks.fe.http.url127.0.0.1:8030, starrocks.fe.jdbc.urljdbc:mysql://127.0.0.1:9030, starrocks.table.identifiertest.score_board, starrocks.userroot, starrocks.password, starrocks.write.properties.partial_updatetrue, starrocks.columnsid,name );在 Spark SQL 客户端插入数据仅更新name列INSERT INTO score_board VALUES (1, starrocks-update), (2, spark-update);查询验证——只有name变化score保持不变mysql select * from score_board; ------------------------------- | id | name | score | ------------------------------- | 1 | starrocks-update | 100 | | 2 | spark-update | 100 | ------------------------------- 2 rows in set (0.02 sec)条件更新——以score列为条件仅当新值大于等于旧值时更新生效在 MySQL 客户端插入初始数据同上。在 Spark SQL 客户端建表将starrocks.write.properties.merge_condition设为score指定条件列并确保 connector 使用普通 Stream Load 接口而非事务接口后者不支持该特性CREATE TABLE score_board USING starrocks OPTIONS( starrocks.fe.http.url127.0.0.1:8030, starrocks.fe.jdbc.urljdbc:mysql://127.0.0.1:9030, starrocks.table.identifiertest.score_board, starrocks.userroot, starrocks.password, starrocks.write.properties.merge_conditionscore );插入数据id1使用更小的 scoreid2使用更大的 scoreINSERT INTO score_board VALUES (1, starrocks-update, 99), (2, spark-update, 101);查询验证——只有id2的行变化id1的行保持不变mysql select * from score_board; --------------------------- | id | name | score | --------------------------- | 1 | starrocks | 100 | | 2 | spark-update | 101 | --------------------------- 2 rows in set (0.03 sec)加载数据到 BITMAP 类型列BITMAP常用于加速精确去重计数如 UV 统计详见 使用 Bitmap 实现精确去重。以下以 UV 统计为例BITMAP 自 1.1.1 起支持在数据库test中创建聚合表page_uvvisit_users列定义为BITMAP类型并配置聚合函数BITMAP_UNIONCREATE TABLE test.page_uv ( page_id INT NOT NULL COMMENT page ID, visit_date datetime NOT NULL COMMENT access time, visit_users BITMAP BITMAP_UNION NOT NULL COMMENT user ID ) ENGINEOLAP AGGREGATE KEY(page_id, visit_date) DISTRIBUTED BY HASH(page_id);创建 Spark 表。Spark 表的 schema 从 StarRocks 表推断但 Spark 不支持BITMAP类型因此需通过starrocks.column.typesvisit_users BIGINT将对应列自定义为BIGINT。Stream Load 写入时 connector 会使用to_bitmap函数把BIGINT数据转换为BITMAP类型CREATE TABLE page_uv USING starrocks OPTIONS( starrocks.fe.http.url127.0.0.1:8030, starrocks.fe.jdbc.urljdbc:mysql://127.0.0.1:9030, starrocks.table.identifiertest.page_uv, starrocks.userroot, starrocks.password, starrocks.column.typesvisit_users BIGINT );在spark-sql中加载数据INSERT INTO page_uv VALUES (1, CAST(2020-06-23 01:30:30 AS TIMESTAMP), 13), (1, CAST(2020-06-23 01:30:30 AS TIMESTAMP), 23), (1, CAST(2020-06-23 01:30:30 AS TIMESTAMP), 33), (1, CAST(2020-06-23 02:30:30 AS TIMESTAMP), 13), (2, CAST(2020-06-23 01:30:30 AS TIMESTAMP), 23);计算各页面的 UVMySQL [test] SELECT page_id, COUNT(DISTINCT visit_users) FROM page_uv GROUP BY page_id; -------------------------------------- | page_id | count(DISTINCT visit_users) | -------------------------------------- | 2 | 1 | | 1 | 3 | -------------------------------------- 2 rows in set (0.01 sec)类型转换细节connector 对 Spark 的TINYINT、SMALLINT、INTEGER、BIGINT类型数据使用to_bitmap转换为 StarRocks 的BITMAP对其他 Spark 数据类型使用bitmap_hash或bitmap_hash64函数。加载数据到 HLL 类型列HLL用于近似去重计数详见 使用 HLL 实现近似去重。以下同样以 UV 统计为例HLL 自 1.1.1 起支持在数据库test中创建聚合表hll_uvvisit_users列定义为HLL类型并配置聚合函数HLL_UNIONCREATE TABLE hll_uv ( page_id INT NOT NULL COMMENT page ID, visit_date datetime NOT NULL COMMENT access time, visit_users HLL HLL_UNION NOT NULL COMMENT user ID ) ENGINEOLAP AGGREGATE KEY(page_id, visit_date) DISTRIBUTED BY HASH(page_id);创建 Spark 表通过starrocks.column.typesvisit_users BIGINT自定义列类型。Stream Load 写入时 connector 使用hll_hash函数将BIGINT数据转换为HLL类型CREATE TABLE hll_uv USING starrocks OPTIONS( starrocks.fe.http.url127.0.0.1:8030, starrocks.fe.jdbc.urljdbc:mysql://127.0.0.1:9030, starrocks.table.identifiertest.hll_uv, starrocks.userroot, starrocks.password, starrocks.column.typesvisit_users BIGINT );在spark-sql中加载数据INSERT INTO hll_uv VALUES (3, CAST(2023-07-24 12:00:00 AS TIMESTAMP), 78), (4, CAST(2023-07-24 13:20:10 AS TIMESTAMP), 2), (3, CAST(2023-07-24 12:30:00 AS TIMESTAMP), 674);计算各页面的 UVMySQL [test] SELECT page_id, COUNT(DISTINCT visit_users) FROM hll_uv GROUP BY page_id; -------------------------------------- | page_id | count(DISTINCT visit_users) | -------------------------------------- | 4 | 1 | | 3 | 2 | -------------------------------------- 2 rows in set (0.01 sec)加载数据到 ARRAY 类型列以下示例展示如何加载ARRAY类型列在数据库test中创建主键表array_tbl包含一个INT列与两个ARRAY列CREATE TABLE array_tbl ( id INT NOT NULL, a0 ARRAYSTRING, a1 ARRAYARRAYINT ) ENGINEOLAP PRIMARY KEY(id) DISTRIBUTED BY HASH(id) ;写入数据。由于部分 StarRocks 版本不提供ARRAY列的元数据connector 无法推断对应 Spark 数据类型需在starrocks.column.types中显式指定本例配置为a0 ARRAYSTRING,a1 ARRAYARRAYINT。在spark-shell中执行val data Seq( | (1, Seq(hello, starrocks), Seq(Seq(1, 2), Seq(3, 4))), | (2, Seq(hello, spark), Seq(Seq(5, 6, 7), Seq(8, 9, 10))) | ) val df data.toDF(id, a0, a1) df.write .format(starrocks) .option(starrocks.fe.http.url, 127.0.0.1:8030) .option(starrocks.fe.jdbc.url, jdbc:mysql://127.0.0.1:9030) .option(starrocks.table.identifier, test.array_tbl) .option(starrocks.user, root) .option(starrocks.password, ) .option(starrocks.column.types, a0 ARRAYSTRING,a1 ARRAYARRAYINT) .mode(append) .save()查询验证MySQL [test] SELECT * FROM array_tbl; ------------------------------------------------- | id | a0 | a1 | ------------------------------------------------- | 1 | [hello,starrocks] | [[1,2],[3,4]] | | 2 | [hello,spark] | [[5,6,7],[8,9,10]] | ------------------------------------------------- 2 rows in set (0.01 sec)加载嵌套列STRUCT、ARRAY 和 MAP嵌套列加载自 1.1.3 起支持。Spark connector 支持写入 StarRocks 的STRUCT、ARRAY、MAP类型列但必须通过starrocks.column.types显式声明 StarRocks 列类型connector 才能正确序列化数据。嵌套数据类型映射嵌套值在通过 Stream Load 发送至 StarRocks 前会被序列化为兼容 JSON 的字符串STRUCT列序列化为 JSON 对象{field1: value1, field2: value2}。ARRAY列序列化为 JSON 数组[value1, value2, ...]。MAP列序列化为 JSON 对象{key1: value1, key2: value2}。示例给定如下 StarRocks 表CREATE TABLE nested_tbl ( id INT, info STRUCTtype STRING, phone BIGINT, tags ARRAYSTRING, attributes MAPSTRING, STRING ) ENGINEOLAP DUPLICATE KEY(id) DISTRIBUTED BY HASH(id) BUCKETS 4;从 Spark 写入数据import org.apache.spark.sql.types._ import org.apache.spark.sql.Row val schema StructType(Seq( StructField(id, IntegerType), StructField(info, StructType(Seq( StructField(type, StringType), StructField(phone, LongType) ))), StructField(tags, ArrayType(StringType)), StructField(attributes, MapType(StringType, StringType)) )) val data Seq( Row(1, Row(admin, 123456789L), Seq(spark, starrocks), Map(env - prod)) ) val df spark.createDataFrame(spark.sparkContext.parallelize(data), schema) val columnTypes info STRUCTtype STRING, phone BIGINT, tags ARRAYSTRING, attributes MAPSTRING, STRING df.write.format(starrocks) .option(starrocks.fenodes, 127.0.0.1:8030) .option(starrocks.table.identifier, test.nested_tbl) .option(starrocks.user, root) .option(starrocks.password, ) .option(starrocks.column.types, columnTypes) .mode(append) .save()底层原理从连接器到 Stream Load 的调用链Spark connector 的所有写入最终都落在 StarRocks 的 Stream Load 服务上这一点可以在仓库源码中得到印证BE 在启动时通过 HTTP 服务注册了 Stream Load 相关处理器。在 be/src/service/service_be/http_service.cpp 中可以看到PUT /api/{db}/{table}/_stream_load路由被注册到StreamLoadAction同时注册了事务接口的POST/PUT /api/transaction/{txn_op}begin/prepare/commit/rollback 等操作与PUT /api/transaction/load数据写入路由后者对应TransactionStreamLoadAction。be/src/http/action/stream_load.h 定义了StreamLoadAction其注释明确说明鉴权身份校验 表级 INSERT 权限由 FE 在loadTxnBegin/streamLoadPutRPC 流程中完成。这正是文档要求账号必须拥有目标表 SELECT 与 INSERT 权限见 GRANT的底层原因。理解这条调用链有助于排查问题Spark 端出现连接失败时应首先确认三个端口链路FEhttp_port、FEquery_port、BEbe_http_port是否畅通出现权限报错时则应回到账号的 INSERT/SELECT 权限上。小结本文完整覆盖了 StarRocks Spark connector 的版本兼容矩阵、JAR 获取与自编译方式、全部写入参数与类型映射规则并给出了 DataFrameBatch/Structured Streaming、Spark SQL 以及主键表部分更新/条件更新、BITMAP/HLL/ARRAY/嵌套列STRUCT/ARRAY/MAP等典型场景的可运行示例。无论你是搭建批式离线入仓管道还是构建流式实时写入链路都可以直接以本文的示例为模板结合自身环境修改连接参数后落地使用。更多示例可参考 Spark connector 官方源码仓库的测试目录src/test/java/com/starrocks/connector/spark/examples。【免费下载链接】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),仅供参考