StarRocks Spark Connector
使用 Spark connector 导入数据(推荐)
StarRocks 提供 Apache Spark™ 连接器 (StarRocks Connector for Apache Spark™),可以通过 Spark 导入数据至 StarRocks(推荐)。 基本原理是对数据攒批后,通过 Stream Load 批量导入StarRocks。Connector 导入数据基于Spark DataSource V2 实现, 可以通过 Spark DataFrame 或 Spark SQL 创建 DataSource,支持 Batch 和 Structured Streaming。
注意
使用 Spark connector 导入数据至 StarRocks 需要目标表的 SELECT 和 INSERT 权限。如果您的用户账号没有这些权限,请参考 GRANT 给用户赋权。
版本要求
| Spark connector | Spark | StarRocks | Java | Scala |
|---|---|---|---|---|
| 1.1.4 | 4.0, 4.1 | 2.5 及以上 | 17 | 2.13 |
| 1.1.4 | 3.3, 3.4, 3.5 | 2.5 及以上 | 8 | 2.12 |
| 1.1.3 | 3.2, 3.3, 3.4, 3.5 | 2.5 及以上 | 8 | 2.12 |
| 1.1.2 | 3.2, 3.3, 3.4, 3.5 | 2.5 及以上 | 8 | 2.12 |
| 1.1.1 | 3.2, 3.3, or 3.4 | 2.5 及以上 | 8 | 2.12 |
| 1.1.0 | 3.2, 3.3, or 3.4 | 2.5 及以上 | 8 | 2.12 |
注意
- 了解不同版本的 Spark connector 之间的行为变化,请查看升级 Spark connector。
- 自 1.1.1 版本起,Spark connector 不再提供 MySQL JDBC 驱动程序,您需要将驱动程序手动放到 Spark 的类路径中。您可以在 MySQL 官网或 Maven 中央仓库上找到该驱动程序。
获取 Connector
您可以通过以下方式获取 connector jar 包
- 直接下载已经编译好的jar
- 通过 Maven 添加 connector 依赖
- 通过源码手动编译
connector jar包的命名格式如下
starrocks-spark-connector-${spark_version}_${scala_version}-${connector_version}.jar
比如,想在 Spark 3.2 和 scala 2.12 上使用 1.1.0 版本的 connector,可以选择 starrocks-spark-connector-3.2_2.12-1.1.0.jar。
注意
一般情况下最新版本的 connector 只维护最近3个版本的 Spark。
直接下载
可以在 Maven Central Repository 获取不同版本的 connector jar。
Maven 依赖
依赖配置的格式如下,需要将 spark_version、scala_version 和 connector_version 替换成对应的版本。
<dependency>
<groupId>com.starrocks</groupId>
<artifactId>starrocks-spark-connector-${spark_version}_${scala_version}</artifactId>
<version>${connector_version}</version>
</dependency>
比如,想在 Spark 3.2 和 scala 2.12 上使用 1.1.0 版本的 connector,可以添加如下依赖
<dependency>
<groupId>com.starrocks</groupId>
<artifactId>starrocks-spark-connector-3.2_2.12</artifactId>
<version>1.1.0</version>
</dependency>
手动编译
-
下载 Spark 连接器代码。
-
通过如下命令进行编译,需要将
spark_version替换成相应的 Spark 版本sh build.sh <spark_version>比如,在 Spark 3.2 上使用,命令如下
sh build.sh 3.2 -
编译完成后,
target/目录下会生成 connector jar 包,比如starrocks-spark-connector-3.2_2.12-1.1-SNAPSHOT.jar。
注意
非正式发布的connector版本会带有
SNAPSHOT后缀。
参数说明
starrocks.fe.http.url
- 是否必填:是
- 默认值:无
- 描述:FE 的 HTTP 地址,支持输入多个FE地址,使用逗号 , 分隔。格式为
<fe_host1>:<fe_http_port1>,<fe_host2>:<fe_http_port2>。自版本 1.1.1 开始,您还可以在 URL 中添加http://前缀,例如http://<fe_host1>:<fe_http_port1>,http://<fe_host2>:<fe_http_port2>。
starrocks.fe.jdbc.url
- 是否必填:是
- 默认值:无
- 描述:FE 的 MySQL Server 连接地址。格式为
jdbc:mysql://<fe_host>:<fe_query_port>。
starrocks.table.identifier
- 是否必填:是
- 默认值:无
- 描述:StarRocks 目标表的名称,格式为
<database_name>.<table_name>。
starrocks.user
- 是否必填:是
- 默认值:无
- 描述:StarRocks 集群账号的用户名。使用 Spark connector 导入数据至 StarRocks 需要目标表的 SELECT 和 INSERT 权限。如果您的用户账号没有这些权限,请参考 GRANT 给用户赋权。
starrocks.password
- 是否必填:是
- 默认值:无
- 描述:StarRocks 集群账号的用户密码。
starrocks.write.label.prefix
- 是否必填:否
- 默认值:
spark- - 描述:指定Stream Load使用的label的前缀。
starrocks.write.enable.transaction-stream-load
- 是否必填:否
- 默认值:
true - 描述:是否使用 Stream Load 事务接口导入数据。要求 StarRocks 版本为 v2.5 或更高。此功能可以在一次导入事务中导入更多数据,同时减少内存使用量,提高性能。
自 1.1.1 版本以来,只有当 starrocks.write.max.retries 的值为非正数时,此参数才会生效,因为 Stream Load 事务接口不支持重试。
starrocks.write.buffer.size
- 是否必填:否
- 默认值:
104857600 - 描述:积攒在内存中的数据量,达到该阈值后数据一次性发送给 StarRocks,支持带单位
k,m,g。增大该值能提高导入性能,但会带来写入延迟。
starrocks.write.buffer.rows
- 是否必填:否
- 默认值:Integer.MAX_VALUE
- 描述:自 1.1.1 版本起支持。积攒在内存中的数据行数,达到该阈值后数据一次性发送给 StarRocks。
starrocks.write.flush.interval.ms
- 是否必填:否
- 默认值:300000
- 描述:数据攒批发送的间隔,用于控制数据写入StarRocks的延迟。
starrocks.write.max.retries
- 是否必填:否
- 默认值:
3 - 描述:自 1.1.1 版本起支持。如果一批数据导入失败,Spark connector 导入该批数据的重试次数上线。
由于 Stream Load 事务接口不支持重试。如果此参数为正数,则 Spark connector 始终使用 Stream Load 接口,并忽略 starrocks.write.enable.transaction-stream-load 的值。
starrocks.write.retry.interval.ms
- 是否必填:否
- 默认值:
10000 - 描述:自 1.1.1 版本起支持。如果一批数据导入失败,Spark connector 尝试再次导入该批数据的时间间隔。
starrocks.write.use_bitmap_hash64
- 是否必填:否
- 默认值:
false - 描述:自 1.1.3 版本起支持。是否使用 64 位哈希函数生成位图。默认使用 32 位哈希函数。
starrocks.columns
- 是否必填:否
- 默认值:无
- 描述:支持向 StarRocks 表中写入部分列,通过该参数指定列名,多个列名之间使用逗号 (,) 分隔,例如
"c0,c1,c2"。
starrocks.write.properties.*
- 是否必填:否
- 默认值:无
- 描述:指定 Stream Load 的参数,用于控制导入行为,例如使用
starrocks.write.properties.format指定导入数据的格式为 CSV 或者 JSON。更多参数和说明,请参见 Stream Load。
starrocks.write.properties.format
- 是否必填:否
- 默认值:
CSV - 描述:指定导入数据的格式,取值为 CSV 和 JSON。connector 会将每批数据转换成相应的格式发送给 StarRocks。
starrocks.write.properties.row_delimiter
- 是否必填:否
- 默认值:
\n - 描述:使用CSV格式导入时,用于指定行分隔符。
starrocks.write.properties.column_separator
- 是否必填:否
- 默认值:
\t - 描述:使用CSV格式导入时,用于指定列分隔符。
starrocks.write.properties.partial_update
- 是否必填:否
- 默认值:
FALSE - 描述:是否使用部分列更新。取值包括
TRUE和FALSE。
starrocks.write.properties.partial_update_mode
- 是否必填:否
- 默认值:
row - 描述:指定部分更新的模式,取值包括
row和column。row(默认值),指定使用行模式执行部分更新,比较适用于较多列且小批量的实时更新场景。column,指定使用列模式执行部分更新,比较适用于少数列并且大量行的批处理更新场景。在该场景,开启列模式,更新速度更快。例如,在一个包含 100 列的表中,每次更新 10 列(占比 10%)并更新所有行,则开启列模式,更新性能将提高 10 倍。
starrocks.write.num.partitions
- 是否必填:否
- 默认值:无
- 描述:Spark用于并行写入的分区数,数据量小时可以通过减少 分区数降低导入并发和频率,默认分区数由Spark决定。使用该功能可能会引入 Spark Shuffle cost。
starrocks.write.partition.columns
- 是否必填:否
- 默认值:无
- 描述:用于Spark分区的列,只有指定 starrocks.write.num.partitions 后才有效,如果不指定则使用所有写入的列进行分区。
starrocks.timezone
- 是否必填:否
- 默认值:JVM 默认时区
- 描述:自 1.1.1 版本起支持。StarRocks 的时区。用于将 Spark 的
TimestampType类型的值转换为 StarRocks 的DATETIME类型的值。默认为ZoneId#systemDefault()返回的 JVM 时区。格式可以是时区名称,例如 Asia/Shanghai,或时区偏移,例如 +08:00。
数据类型映射
-
数据类型映射默认如下:
Spark 数据类型 StarRocks 数据类型 BooleanType BOOLEAN ByteType TINYINT ShortType SMALLINT IntegerType INT LongType BIGINT StringType LARGEINT FloatType FLOAT DoubleType DOUBLE DecimalType DECIMAL StringType CHAR StringType VARCHAR StringType STRING StringType JSON DateType DATE TimestampType DATETIME ArrayType ARRAY
说明:
自版本 1.1.1 开始支持。详细步骤, 请参见 导入至 ARRAY 类型的列。MapType MAP
说明:
自版本 1.1.3 开始支持。详细步骤, 请参见 导入嵌套列。StructType STRUCT
说明:
自版本 1.1.3 开始支持。详细步骤, 请参见 导入嵌套列。 -
您还可以自定义数据类型映射。
例如,一个 StarRocks 表包含了 BITMAP 和 HLL 类型的列,但 Spark 不支持这两种数据类型。则您需要在 Spark 中设置其支持的数据类型,并且自定义数据类型映射关系。详细步骤,参见导入至 BITMAP 和 HLL 类型的列。自版本 1.1.1 起支持导入至 BITMAP 和 HLL 类型的列。