Spark向量化
1. 概述
向量化执行引擎是 BMR Spark 内置的查询加速能力,基于 Apache Gluten 与 Velox 构建,把 Spark SQL 的算子计算下沉到 C++ 原生向量化引擎执行。
它完全兼容 Apache Spark API,业务代码与 SQL 无需任何改造;BMR 集群已默认开启,您不需要做额外配置。 对于当前尚未支持的算子、数据类型、文件格式或函数,引擎会自动回退到原生 Spark 执行,结果保持正确。
1.1 支持的版本
| BMR 版本 | Spark 版本 | 说明 |
|---|---|---|
| BMR 3.5.2 | Spark 3.3.2 | 已支持 |
| BMR 3.7.3 及以上 | Spark 3.5.5 | 已支持 |
| 其他 BMR 版本 | — | 暂不提供向量化执行能力 |
1.2 性能表现
在 256 CU 计算规模、5 TB TPC-DS 数据集下,Spark 3.5.5 开启向量化执行后,整体性能相比基线 提升 3.98 倍。
| 测试项 | 说明 |
|---|---|
| 计算规模 | 256 CU |
| 数据集 | TPC-DS 5 TB |
| Spark 版本 | 3.5.5 |
| 对比基线 | 原生 Spark |
| 整体结果 | 性能提升 3.98 倍 |
2. 原理介绍
随着 SSD 与高速网卡的普及,大数据计算的瓶颈已经从 IO 转向 CPU。而 Spark 运行在 JVM 之上, 即使有 Codegen 这类优化,也仍受字节码长度、方法参数个数等限制,难以充分利用现代 CPU 的向量化指令。
向量化执行引擎的做法是只替换计算部分,不改变 Spark 的架构:SQL 解析、Catalyst 优化、 任务调度、失败重试、整个分布式框架都仍由 Spark 负责;Gluten 在物理计划阶段把可下沉的算子改写为原生算子, 序列化成 Substrait 计划后经 JNI 交给 Velox,由 C++ 引擎完成单个 Task 的实际计算。
性能提升主要来自四个方面:
| 加速来源 | 说明 |
|---|---|
| 列式批处理与 SIMD | 数据以列式批(每批默认 4096 行)在内存中连续存放,单条 CPU 指令可处理多个数据, 同时消除了逐行处理带来的虚函数调用与解释执行开销 |
| 原生堆外内存 | 计算全程在堆外进行,避免 JVM 对象开销与 GC 停顿 |
| 原生列式 Shuffle | 直接以列式格式落盘与网络传输,省去行列转换和 Java 序列化 |
| 延迟物化与过滤下推 | 过滤条件下推到文件读取层,未命中的列不做解码,减少无效 IO 与计算 |

3. 使用限制
当查询中出现尚未支持的算子、数据类型、文件格式或函数时,引擎会在生成执行计划的阶段自动把这部分交还给 原生 Spark 执行,这一机制称为回退(Fallback)。回退后结果保持正确,但在向量化与行式执行的交界处需要做行列格式转换, 因此如果一条查询中回退的部分过多,整体耗时可能不优于原生 Spark。建议参考第6章确认作业的实际回退情况。
3.1 功能限制
| 项目 | 支持情况 | 限制说明 |
|---|---|---|
| 文件格式(读) | 支持(受限) | 仅支持 Parquet、ORC、DWRF。CSV、JSON、Hive 文本表的扫描回退到原生 Spark |
| 文件格式(写) | 支持(受限) | 仅支持 Parquet。ORC 写入、CREATE TABLE AS SELECT、分桶表写入回退; 不支持 brotli、lzo、lz4raw 压缩编码; BMR 3.5.3(Spark 3.3.2)上原生写入仅覆盖静态分区写入,动态分区写入建议先小范围验证 |
| ANSI 模式 | 不支持 | spark.sql.ansi.enabled=true 时整个执行计划回退到原生 Spark |
| RDD 作业 | 不适用 | 向量化仅作用于 Spark SQL 与 DataFrame API,纯 RDD 编写的作业不受益 |
| Structured Streaming | 不支持 | 流式作业不下沉到原生引擎 |
| Python UDF | 支持(受限) | ArrowEvalPython(普通 Pandas UDF)已向量化; AggregateInPandas、MapInPandas、FlatMapGroupsInPandas、 WindowInPandas 这几类分组 Pandas UDF 不下沉 |
| Hive UDF / Scala UDF | 支持(受限) | 同一个 Project 中只有 UDF 本身回到 JVM 计算,其余表达式仍在原生引擎执行, 不会因为一个 UDF 导致整段回退 |
| 正则表达式 | 支持(受限) | 模式串必须是常量;含 lookahead / lookbehind 的模式回退到原生 Spark。 s 的匹配范围与原生 Spark 略有差异,见本文3.2节 |
| 大小写敏感 | 不支持 | 开启 spark.sql.caseSensitive=true 时没有回退保护,结果可能与原生 Spark 不一致, 建议保持该参数默认关闭 |
approx_count_distinct |
支持(受限) | 采用与原生 Spark 不同的 HyperLogLog 实现,sketch 二进制格式不互通, 不能跨引擎读写 sketch 中间结果 |
| JSON 字符串格式 | 支持(受限) | JSON 函数只接受双引号字符串,单引号内容会得到非预期结果; get_json_object 对 [*] 形式的路径返回 null |
| Iceberg | 支持(受限) | 读取已支持;原生写入默认关闭,写入走原生 Spark 路径;含 equality-delete 文件的表读取回退 |
3.2 结果差异与精度说明
以下差异不会触发回退、不会报错。表中的样例与两侧输出全部取自引擎自带的回归测试基线 (与 Spark 3.5 官方预期结果逐条比对得出),可直接复现。
精度量级:已确认的数值差异全部出现在 double 类型的浮点聚合上, 差异位于第 16~17 位有效数字,量级为 1 个 ULP(浮点最小精度单位)。 测试基线中未发现 Decimal 类型的结果差异,整型的 sum、count、 avg 也未出现差异。若作业对浮点一致性敏感,可将 spark.gluten.sql.columnar.backend.velox.floatingPointMode 设为 strict, 该模式会关闭 sum(float/double) 与 avg(float/double) 的部分聚合刷写。
| 类别 | 场景与样例 | 原生 Spark | 向量化引擎 | 说明 |
|---|---|---|---|---|
| 精度问题 | corr 完全相关时不返回精确 1.0 |
|||
SELECT corr(DISTINCT x, y)FROM (VALUES (1,1),(2,2),(2,2)) t(x,y) |
1.0 |
0.9999999999999999 |
不要对 corr 结果做等值判断,改用误差范围比较 |
|
| 精度问题 | 统计类聚合末位差异 | |||
SELECT stddev(a) FROM testData |
0.8997354108424372 |
0.8997354108424375 |
variance、skewness、kurtosis 属于同类。同一组数据上 skewness 为 -0.2723801058145729 与 -0.27238010581457284 |
|
| 精度问题 | 线性回归函数末位差异 | |||
SELECT regr_sxx(y, x) FROM testRegression |
288.6666666666667 |
288.66666666666663 |
regr_r2(0.997690531177829 与 0.9976905311778291)、 regr_slope(0.9988445981121533 与 0.9988445981121532)属于同类 |
|
| 算法实现问题 | 相同 seed 的随机序列不同 | |||
SELECT rand(0) |
0.7604953758285915 |
0.5488135024422883 |
两者伪随机数实现不同。依赖固定 seed 复现抽样、随机打散或分流的作业会得到不同的行集合; 需要严格复现时请关闭向量化执行 |
4. 适用范围
以下支持范围基于 Spark 3.3.2 与 Spark 3.5.5。未列出的格式、类型、算子与函数会回退到原生 Spark 执行。
4.1 存储格式
数据格式
| 格式 | 读取 | 说明 |
|---|---|---|
| Parquet | 支持 | 推荐格式。开启 mergeSchema 或读取加密文件时回退 |
| ORC | 支持 | 不支持数组中嵌套结构体或数组、Map 的键为结构体、Map 的值为数组、以及 Timestamp 列; char(n) 类型列默认回退 |
| DWRF | 支持 | — |
| CSV / JSON / 文本 | 不支持 | 扫描阶段回退到原生 Spark |
表格式
| 表格式 | 读取 | 写入 | 说明 |
|---|---|---|---|
| Hive | 支持 | 支持 | 写入要求输出格式为 Parquet |
| Iceberg | 支持 | 不支持 | 原生写入默认关闭,写入走原生 Spark;含 equality-delete 文件的表读取回退 |
| Paimon | 支持 | 不支持 | — |
存储介质
| 介质 | 访问方式 |
|---|---|
| HDFS | 原生支持,无需额外配置 |
| BOS | 通过 S3 兼容协议访问,配置沿用 spark.hadoop.fs.s3a.* 系列参数 |
4.2 数据类型
| 支持情况 | 数据类型 |
|---|---|
| 支持 | Boolean、Byte(TinyInt)、Short(SmallInt)、Int、 Long(BigInt)、Float、Double、Decimal(精度与小数位上限 38)、 String、Binary、Date、Timestamp, 以及 Array、Map、Struct 嵌套类型 |
| 不支持 | TimestampNTZ(不带时区的时间戳)、Interval Day To Second、 CalendarInterval、用户自定义类型(UDT)、Variant |
| 特殊限制 | 聚合算子的分组列与聚合列不支持 Map 类型; 查询的输入或输出中只要包含 TimestampNTZ,相关算子即整体回退; Decimal 与 Timestamp 之间的相互转换不支持 |
4.3 算子
| 类型 | 支持 | 不支持 |
|---|---|---|
| 数据源 | FileSourceScanExec、HiveTableScanExec、BatchScanExec、 InMemoryTableScanExec、RDDScanExec、RangeExec |
— |
| 数据写入 | DataWritingCommandExec、WriteFilesExec3.5、 AppendDataExec、ReplaceDataExec、OverwriteByExpressionExec、 OverwritePartitionsDynamicExec、WriteToDataSourceV2Exec |
— |
| 通用 | ProjectExec、FilterExec、SortExec、UnionExec、 CoalesceExec、ExpandExec、GenerateExec |
SampleExec(默认关闭) |
| 聚合 | HashAggregateExec、ObjectHashAggregateExec、SortAggregateExec |
— |
| 关联 | BroadcastHashJoinExec、ShuffledHashJoinExec、SortMergeJoinExec、 BroadcastNestedLoopJoinExec、CartesianProductExec |
— |
| 窗口 | WindowExec、WindowGroupLimitExec3.5 |
— |
| 数据交换 | ShuffleExchangeExec、BroadcastExchangeExec、SubqueryBroadcastExec |
— |
| 限流 | GlobalLimitExec、LocalLimitExec、TakeOrderedAndProjectExec、 CollectLimitExec、CollectTailExec |
带 OFFSET 的排序取前 N 回退 |
| Python UDF | ArrowEvalPythonExec、BatchEvalPythonExec |
AggregateInPandasExec、MapInPandasExec、 FlatMapGroupsInPandasExec、WindowInPandasExec |
算子使用条件
| 算子 | 条件与限制 |
|---|---|
| 带 3.5 标记的算子 | 由 Spark 3.4 / 3.5 引入,仅在 BMR 3.7.3 及以上(Spark 3.5.5)存在 |
| 数据写入算子 | 原生写入仅支持 Parquet;ORC 写入、CREATE TABLE AS SELECT、分桶表写入回退到原生 Spark, 详见 第3章 |
SortMergeJoinExec |
支持 Inner、LeftOuter、RightOuter、LeftSemi、LeftAnti;不支持 cross join 与 existence join |
BroadcastNestedLoopJoinExec |
支持 Inner、LeftOuter、RightOuter、ExistenceJoin;LeftOuter 不能广播左表,RightOuter 不能广播右表 |
GenerateExec |
支持 explode、posexplode、inline、json_tuple、stack |
WindowGroupLimitExec |
窗口 Top-N 下推,仅支持 row_number、rank、dense_rank |
InMemoryTableScanExec |
缓存表以列式格式存储;表结构含不支持的类型时自动降级为原生 Spark 缓存格式 |
FilterExec、ProjectExec |
单个算子中嵌套表达式数量达到 50 时回退 |
SampleExec |
默认关闭,始终回退到原生 Spark |
4.4 函数
下表以 Spark 3.3.2 的内建函数清单为基准统计,这部分函数在两个 BMR 版本上都可用: 标量函数 323 个中支持 252 个、不支持 71 个;聚合函数 50 个中支持 45 个、不支持 5 个; 窗口函数 9 个与生成器函数 7 个全部支持。支持的函数中有一部分带使用条件,在表中以 * 标出。
带 3.5 标记的函数是 Spark 3.4 / 3.5 新增的,仅在 BMR 3.7.3 及以上(Spark 3.5.5)可用, BMR 3.5.3(Spark 3.3.2)本身不提供这些函数。把这部分计入后,Spark 3.5.5 上标量函数支持 268 个、不支持 88 个, 聚合函数支持 52 个、不支持 10 个。
| 类别 | 支持 | 不支持 |
|---|---|---|
| 数组函数 | array、array_append3.5、array_compact3.5、array_contains、array_distinct、array_except、array_insert3.5、array_intersect、array_join、array_max、array_min、array_position、array_prepend3.5、array_remove、array_repeat、array_union、arrays_overlap、arrays_zip、flatten、get、shuffle、slice、sort_array |
sequence |
| 集合函数 | array_size、cardinality、concat、reverse、size |
— |
| 位运算函数 | &、|、^、~、bit_count、bit_get、getbit、shiftright |
|
| 条件函数 | coalesce、if、ifnull、nanvl、nullif、nvl、nvl2、when |
— |
| 类型转换函数 | bigint、binary、boolean、cast、date、decimal、double、float、int、smallint、string、timestamp、tinyint |
— |
| 谓词函数 | !、<、<>、>、and、between、case、ilike、isnan、isnotnull、isnull、like、not、or、in、regexp、regexp_like、rlike** |
— |
| 哈希函数 | crc32、hash、md5、sha、sha1、sha2、xxhash64 |
— |
| Map 函数 | element_at、map_contains_key、map_entries、map_keys、map_values、map、map_concat、str_to_map* |
map_from_arrays、map_from_entries、try_element_at |
| 结构体函数 | named_struct、struct |
— |
| 高阶函数 | aggregate、array_sort、exists、filter、forall、map_filter、map_zip_with、reduce3.5、transform、transform_keys、transform_values、zip_with |
— |
| 数学函数 | %、*、+、/、abs、acos、acosh、asin、asinh、atan、atan2、atanh、bin、cbrt、conv、cos、cosh、cot、csc、degrees、div、e、exp、expm1、factorial、greatest、hex、hypot、least、log、log10、log1p、log2、mod、negative、pi、pmod、positive、pow、power、rand、random、rint、round、sec、shiftleft、sign、signum、sinh、sqrt、unhex、width_bucket、ceil、ceiling、floor、try_add** |
bround、ln、radians、randn、sin、tan、tanh、try_divide、try_multiply、try_subtract |
| 字符串函数 | ascii、bit_length、btrim、char、char_length、character_length、chr、concat_ws、find_in_set、initcap、instr、lcase、left、len3.5、length、levenshtein、locate、lower、ltrim、luhn_check3.5、mask3.5、overlay、position、repeat、replace、right、rtrim、soundex、split、split_part、substring_index、translate、trim、ucase、upper、substr、substring、contains、endswith、startswith、lpad、rpad、regexp_extract、regexp_extract_all、regexp_replace、base64、unbase64** |
elt、encode、format_number、format_string、octet_length、printf、regexp_count3.5、regexp_instr3.5、regexp_substr3.5、sentences、space、to_binary、to_char3.5、to_number、to_varchar3.5、try_to_binary、try_to_number |
| 日期与时间函数 | add_months、date_add、date_diff3.5、date_format、date_from_unix_date、date_sub、date_trunc、dateadd3.5、datediff、day、dayofmonth、dayofweek、dayofyear、extract、from_unixtime、from_utc_timestamp、hour、last_day、make_date、make_timestamp、make_ym_interval、minute、month、months_between、next_day、quarter、second、timestamp_micros、timestamp_millis、to_utc_timestamp、trunc、unix_date、unix_micros、unix_millis、unix_seconds、unix_timestamp、weekday、weekofyear、year、timestamp_seconds、to_unix_timestamp** |
convert_timezone、date_part、datepart3.5、make_dt_interval、make_interval、make_timestamp_ltz、make_timestamp_ntz、session_window、to_timestamp_ltz、to_timestamp_ntz、try_to_timestamp3.5、window、window_time3.5 |
| JSON 函数 | get_json_object、json_array_length、json_object_keys、json_tuple、from_json、to_json** |
schema_of_json |
| URL 函数 | url_decode3.5、url_encode3.5 |
parse_url |
| CSV 函数 | — | from_csv、schema_of_csv、to_csv |
| XML 函数 | — | xpath、xpath_boolean、xpath_double、xpath_float、xpath_int、xpath_long、xpath_number、xpath_short、xpath_string |
| 其他函数 | assert_true、equal_null3.5、spark_partition_id、uuid、version、| |
|
| 聚合函数 | any、any_value3.5、approx_count_distinct、array_agg、avg、bit_and、bit_or、bit_xor、bool_and、bool_or、collect_list、collect_set、corr、count、count_if、covar_pop、covar_samp、every、first、first_value、grouping、grouping_id、kurtosis、last、last_value、max、max_by、mean、median3.5、min、min_by、regr_avgx、regr_avgy、regr_count、regr_intercept3.5、regr_r2、regr_slope3.5、regr_sxx3.5、regr_sxy3.5、regr_syy3.5、skewness、some、std、stddev、stddev_pop、stddev_samp、sum、try_avg、var_pop、var_samp、variance、try_sum* |
approx_percentile、bitmap_construct_agg3.5、bitmap_or_agg3.5、count_min_sketch、histogram_numeric、hll_sketch_agg3.5、hll_union_agg3.5、mode3.5、percentile、percentile_approx |
| 窗口函数 | cume_dist、dense_rank、lag、lead、nth_value、ntile、percent_rank、rank、row_number |
— |
| 生成器函数 | explode、explode_outer、inline、inline_outer、posexplode、posexplode_outer、stack |
— |
5. 开启与关闭
5.1 默认已开启
BMR 集群已在 spark-defaults.conf 中预置以下配置,向量化执行默认生效,您无需做任何操作。
| 配置项 | 值 |
|---|---|
spark.plugins |
org.apache.gluten.GlutenPlugin |
spark.shuffle.manager |
org.apache.spark.shuffle.sort.ColumnarShuffleManager |
spark.memory.offHeap.enabled |
true |
spark.memory.offHeap.size |
按 5.2 设置 |
前两项是引擎的加载入口,不建议修改。如需关闭向量化执行,请使用5.3节的开关,而不是删除这两项配置。
5.2 内存配置建议
向量化执行的计算过程全部在堆外内存进行,因此需要为 Executor 配置堆外内存。推荐比例为 每 1 个 CPU 核配置 1 GB executor.memory 与 2 GB offHeap.size。
此外,spark.executor.memoryOverhead 若未显式设置,会被自动调整为 max(0.3 × offHeap.size, 384 MiB)。做容量规划时需要把这部分一并算入,否则容易因为 Executor 总内存超出申请值而被 YARN 终止。
| Executor 核数 | spark.executor.memory |
spark.memory.offHeap.size |
自动调整的 memoryOverhead |
单 Executor 实际占用 |
|---|---|---|---|---|
| 2 | 2 GB | 4 GB | 1.2 GB | 约 7.2 GB |
| 4 | 4 GB | 8 GB | 2.4 GB | 约 14.4 GB |
| 8 | 8 GB | 16 GB | 4.8 GB | 约 28.8 GB |
注意:spark.memory.offHeap.enabled 为 false,或 spark.memory.offHeap.size 小于 1 MB 时,作业会在 Driver 启动阶段直接失败。 如果您需要完全不使用堆外内存,请先按5.3节关闭向量化执行。
5.3 关闭方式
关闭后对应部分回落到原生 Spark 的行式执行,查询结果不变。三种生效途径任选其一:
| 途径 | 生效范围 | 操作 |
|---|---|---|
| 修改集群配置 | 集群上所有作业 | 在控制台的配置管理中修改 spark-defaults.conf,添加或调整对应配置项 |
| 提交作业时覆盖 | 仅本次作业 | spark-submit --conf spark.gluten.enabled=false ... |
| 会话中动态设置 | 仅当前会话 | 在 spark-sql 中执行 SET spark.gluten.enabled=false; |
可用开关
| 配置项 | 默认值 | 作用 |
|---|---|---|
spark.gluten.enabled |
true |
总开关。设为 false 后整个作业完全使用原生 Spark 执行 |
spark.plugins |
org.apache.gluten.GlutenPlugin |
向量化引擎加载Plugin。移除配置后可彻底关闭向量化引擎初始化 |
6. 查看回退信息
作业是否真正走了向量化执行、哪些算子发生了回退,都可以在 Spark UI 中直接查看,无需额外配置或改代码。 打开 Spark UI 后切换到 Gluten SQL / DataFrame 页签即可。
| 页面元素 | 含义 |
|---|---|
| Num Gluten Nodes | 该查询中成功下沉到原生引擎执行的算子数量,数值越大说明向量化覆盖越充分 |
| Num Fallback Nodes | 回退到原生 Spark 执行的算子数量,为 0 表示全程向量化 |
+details 展开 |
显示 == Fallback Summary == 块,逐行给出 (算子编号) 算子名: 回退原因; 若该查询全程向量化,则显示 No fallback nodes |
| 页签内的版本信息 | 展示当前使用的 Apache Gluten 与 Velox 版本 |
如果发现回退节点数偏多,可结合第3章的限制清单排查原因,常见情况是表使用了 CSV / JSON / 文本格式、查询开启了 ANSI 模式、或使用了尚未支持的函数。
评价此篇文章
