Configuring the Spark Native Engine
Scenarios
The Spark Native engine uses the vectorized C++ acceleration library to accelerate Spark operators. Traditional Spark SQL is based on row data and uses JVM codegen to accelerate query. The JVM has a range of restrictions on the generated Java code, such as the method length and number of parameters, and the memory bandwidth utilization of row data is low. The performance needs to be improved. When the mature vectorized C++ acceleration library is used, data is stored in the memory in vectorized format, which improves bandwidth utilization and speeds up queries by processing data in batches.
You can enable the Spark Native engine to accelerate Spark SQL queries. After the Spark Native engine is configured, large Spark SQL statements cannot be queried or deleted.
Constraints
- This section applies only to MRS 3.6.0.1 and later.
- The system supports tasks executed via Spark SQL and Spark Beeline.
- Both MRS and open-source readers are supported. This is controlled by the spark.gluten.sql.columnar.backend.ch.mrs.reader parameter (default value: true).
- The Spark large SQL inspection feature is not supported.
- Data types
Data Type
Supported by MRS Reader
Supported by Open-Source Reader
BooleanType, ByteType, CharType, DateType, DoubleType, DecimalType, FloatType, IntegerType, LongType, StringType, ShortType, TimestampType, and VarcharType
√
√
BinaryType, DayTimeIntervalType, TimestampNTZType, and YearMonthIntervalType
x
√
ArrayType, MapType, and StructType
x
√
- Data formats
Format
Supported by MRS Reader
Supported by Open-Source Reader
Parquet and ORC
√
√
CSV, JSON, and TXT
x
√
- Operators
Operator
Supported or Not
AdaptiveSparkPlanExec, AQEShuffleReadExec, BroadcastHashJoinExec, CoalesceExec, ExpandExec, FileSourceScan, FilterExec, GenerateExec, GlobalLimitExec, HashAggregateExec, HiveTableScanExec, InMemoryTableScanExec, LimitExec, LocalTableScanExec, ObjectHashAggregateExec, ProjectExec, RangeExec, ReusedSubqueryExec, ReusedExchangeExec, RowToColumnarExec, SortExec, SubqueryExec, SubqueryBroadcastExec, SortMergeJoinExec, ShuffleExchangeExec, ShuffledHashJoinExec, SubqueryAdaptiveBroadcastExec, TakeOrderedAndProjectExec, UnionExec, WholeStageCodeGenExec, WindowExec, BroadcastExchangeExec, and BroadcastNestedLoopJoinExec
√
AppendColumnsExec, CollectLimitExec, CollectMetricsExec, CoGroupExec, DataWritingCommandExec, EvalPythonExec, ExternalRDDScanExec, MapElementsExec, ParallelInsertUnionExec, SampleExec, and V2CommandExec
x
- Functions Only functions contained in the table are supported.
Type
Function Name
Logical operation functions
is_not_null, is_null, gte, gt, lte, lt, equal, and, or, not, xor, extract, cast, alias, and nullif
Time functions
get_timestamp, quarter, to_unix_timestamp, unix_timestamp, date_format, from_unixtime, date_add, date_sub, datediff, second, add_months, trunc, date_trunc, and floor_datetime
Mathematical operation functions
subtract, multiply, add, divide, positive, negative, modulus, pmod, abs, ceil, floor, round, bround, exp, power, cos, cosh, sin, sinh, tan, tanh, acos, asin, atan, atan2, asinh, acosh, atanh, bitwise_not, bitwise_and, bitwise_or, bitwise_xor, sqrt, cbrt, degrees, e, pi, hex, unhex, hypot, sign, log10, log1p, log2, log, radians, greatest, least, shiftleft, shiftright, check_overflow, factorial, rand, and isnan
String functions
like, not_like, starts_with, ends_with, contains, substring, lower, upper, trim, ltrim, rtrim, concat, strpos, char_length, replace, regexp_replace, regexp_extract, regexp_extract_all, chr, rlike, ascii, split, concat_ws, base64, unbase64, lpad, rpad, reverse, md5, translate, repeat, position, locate, and space
Hash functions
sha1, sha2, crc32, murmur3hash, and xxhash64
Table data generation functions
explode and posexplode
Function using the IN operator
In
Null value handling function
coalesce
Array operation functions
array, size, get_array_item, element_at, and range
Map operation functions
map, get_map_value, map_keys, map_values, and map_from_arrays
Tuple-related functions
get_struct_field, and named_struct
JSON-related functions
get_json_object, to_json, from_json, json_tuple, json_array_length, make_decimal, and unscaled_value
- File systems
File System
Supported or Not
HDFS, OBS, HDFS, local file system, or hybrid deployment combining OBS
√
- Platforms
Platform
Supported or Not
x86, Arm, or hybrid deployment of x86 and Arm
√
- Usage methods
Method
Supported or Not
Spark SQL in YARN mode
√
Spark SQL in local mode
√
- Time zone format
Format
Supported or Not
Region-based ID, for example, Asia/Beijing.
√
Format based on the second-level interval, for example, SET TIME ZONE INTERVAL interval_literal
x
- Native Shuffle(spark.shuffle.manager=org.apache.spark.shuffle.sort.ColumnarShuffleManager): Supported compression algorithm (spark.io.compression.codec)
Name
Supported or Not
lz4
√
lzf
x
snappy
x
zstd
√
gzip
x
bzip2
x
- MemArtsCC:
Scenario
Supported by Open-Source Reader
Supported by MRS Reader
MemArtsCC in clusters in security mode
√
x
MemArtsCC in clusters in normal mode
√
√
Configuring Server Parameters
- On MRS Manager, choose Cluster > Services > Spark, click Configurations and then All Configurations, choose Spark(Service) > Native, and configure the following parameters.
Parameter
Description
Default Value
spark.plugins
Plug-in used by Spark. Set this parameter to org.apache.gluten.GlutenPlugin.
NOTE:If spark.plugins has been configured, add org.apache.gluten.GlutenPlugin to the file and separate them with commas (,).
N/A
spark.gluten.memory.unify.enabled
Whether to enable unified memory management for the Spark Native engine. If this parameter is set to true, unified memory management is required for Native acceleration.
true
spark.executorEnv.LD_PRELOAD
LD_PRELOAD environment variable for the executor. This parameter is required when the Spark Native engine is enabled.
Select ${PWD}/native/${CPU_TYPE}/libch.so ${PWD}/native/${CPU_TYPE}/lib_gsasl.so ${PWD}/native/${CPU_TYPE}/lib_lemmagen.so for x86_64.
Select ${PWD}/native/${CPU_TYPE}/libch.so ${PWD}/native/${CPU_TYPE}/lib_gsasl.so ${PWD}/native/${CPU_TYPE}/lib_lemmagen.so ${JAVA_HOME_21}/lib/libjsig.so for aarch64 J.
None
spark.yarn.appMasterEnv.LD_PRELOAD
LD_PRELOAD environment variable for the YARN AM. This parameter is required when the Spark Native engine is enabled.
Select ${PWD}/native/${CPU_TYPE}/libch.so ${PWD}/native/${CPU_TYPE}/lib_gsasl.so ${PWD}/native/${CPU_TYPE}/lib_lemmagen.so for x86_64.
Select ${PWD}/native/${CPU_TYPE}/libch.so ${PWD}/native/${CPU_TYPE}/lib_gsasl.so ${PWD}/native/${CPU_TYPE}/lib_lemmagen.so ${JAVA_HOME_21}/lib/libjsig.so for aarch64 J.
None
spark.shuffle.manager
Shuffle manager. Set this parameter to org.apache.spark.shuffle.sort.ColumnarShuffleManager.
sort
spark.sql.orc.impl
To enable the Native engine to read and write ORC tables, set this parameter to native.
hive
- Click Save to save the changes and restart expired instances.
- Choose Dashboard, click More, and select Download Client to install and use the client.
Configuring Client Parameters
- Skip this section if you have already completed the steps in Configuring Server Parameters and installed the new client.
- Modify the following parameters in the {Client installation directory}/Spark/spark/conf/spark-defaults.conf file on the Spark client.
Parameter
Description
Default Value
spark.plugins
Plug-in used by Spark. Set this parameter to org.apache.gluten.GlutenPlugin.
NOTE:If spark.plugins has been configured, you can add org.apache.gluten.GlutenPlugin to it and separate them with a comma (,).
Left blank
spark.gluten.memory.unify.enabled
Whether to enable unified memory management for the Spark Native engine. If this parameter is set to true, unified memory management is required for Native acceleration.
true
spark.executorEnv.LD_PRELOAD
LD_PRELOAD environment variable for the executor. This parameter is required when the Spark Native engine is enabled.
Set x86_64 to ${PWD}/native/${CPU_TYPE}/libch.so ${PWD}/native/${CPU_TYPE}/lib_gsasl.so ${PWD}/native/${CPU_TYPE}/lib_lemmagen.so.
Set aarch64 to ${PWD}/native/${CPU_TYPE}/libch.so ${PWD}/native/${CPU_TYPE}/lib_gsasl.so ${PWD}/native/${CPU_TYPE}/lib_lemmagen.so ${JAVA_HOME_21}/lib/libjsig.so.
None
spark.yarn.appMasterEnv.LD_PRELOAD
LD_PRELOAD environment variable for the YARN AM. This parameter is required when the Spark Native engine is enabled.
This parameter is mandatory when --deploy-mode cluster is configured using the spark-submit command. Keep the value of this parameter to be the same as that of spark.executorEnv.LD_PRELOAD.
Set x86_64 to ${PWD}/native/${CPU_TYPE}/libch.so ${PWD}/native/${CPU_TYPE}/lib_gsasl.so ${PWD}/native/${CPU_TYPE}/lib_lemmagen.so.
Set aarch64 to ${PWD}/native/${CPU_TYPE}/libch.so ${PWD}/native/${CPU_TYPE}/lib_gsasl.so ${PWD}/native/${CPU_TYPE}/lib_lemmagen.so ${JAVA_HOME_21}/lib/libjsig.so.
None
spark.shuffle.manager
Shuffle manager. Set this parameter to org.apache.spark.shuffle.sort.ColumnarShuffleManager.
sort
spark.sql.orc.impl
To enable the Native engine to read and write ORC tables, set this parameter to native.
hive
- Run the source bigdata_env command in the client directory to update environment variables.
- To connect the MRS reader to MemArtsCC, add the following parameters to the {Client installation directory}/Spark/spark/conf/spark-defaults.conf file on the Spark client:
You must install the MemArtsCC service in the cluster, configure Spark to connect to MemArtsCC, and enable the Spark Native engine. The interworking between Spark Native and MemArtsCC takes effect only when spark.gluten.sql.columnar.backend.ch.mrs.reader is set to true.
Table 1 Parameter description Parameter
Description
Default Value
spark.gluten.sql.columnar.backend.ch.memartscc.enable
Whether to enable the MemArtsCC client for the Spark Native engine.
false
fs.obs.memartscc.config.init_ip
IP address of the MemArtsCC worker server. (Only one IP address can be configured.)
127.0.0.1
fs.obs.memartscc.config.init_port
Port of the MemArtsCC worker server.
21840
fs.obs.memartscc.config.zk_endpoint
IP address and port list for the ZooKeeper server.
127.0.0.1:2181
fs.obs.memartscc.config.zk_root_node
Root node path for MemArtsCC in ZooKeeper.
/memartccnew
fs.obs.memartscc.config.zk_cluster_name
Cluster name registered for MemArtsCC in ZooKeeper.
cc01
fs.obs.memartscc.config.log_level
Log level for the MemArtsCC client.
INFO
fs.obs.memartscc.config.client_type
Type of the MemArtsCC client.
obsa
Usage Instructions
You can run the following statement to check whether the Native Engine is enabled by checking the execution plan:
spark-sql -e "explain select * from database.data_source"
If "CHNativeColumnarToRow" is displayed in the execution plan, the Native Engine is enabled.

Data Consistency Verification Tool
The data consistency verification tool uses shell scripts to start the Scala program by invoking the Spark Shell. This Scala program implements the core comparison logic. It compares table data generated by standard Spark queries against table data generated by Spark Native queries. You can perform the following steps to check the query results from Spark and Spark Native:
- Enable Spark Native by referring to Configuring Client Parameters.
- Switch to the {Client installation directory}/Spark/spark/tool directory and run the sh md5compare.sh --help command to view the usage details:
md5compare.sh --dbname <name> --query-path <path> --spark_tbl_location <path> --native_tbl_location <path> --spark_tbl_name <name> --native_tbl_name <name> --clear_tbl_data <boolean> --driver_mem <MEM> --executor_mem <MEM> --executor_num <NUM> --executor_core <NUM>
Table 2 Parameter description Parameter
Mandatory
Description
--dbname <Database name>
Yes
Target database for SQL queries.
--query-path <File path>
Yes
Path to the SQL query script. This file can contain only one SQL statement.
--spark_tbl_location <Directory path>
No
Directory of the external tables for storing Spark query results. The default path is /tmp/spark.
--native_tbl_location <Directory path>
No
Directory of the external tables for storing native query results. The default value is /tmp/native.
--spark_tbl_name <Table name>
No
Name of the temporary table for storing Spark query results. The default name is spark_data.
--native_tbl_name <Table name>
No
Name of the temporary table for storing native query results. The default name is native_data.
--clear_tbl_data <Boolean value>
No
Whether to delete the temporary table and external table path after the comparison. The default value is true.
--driver_mem <Memory value>
No
Memory allocated to the driver. The default value is 4 GB.
--executor_mem <Memory value>
No
Memory allocated to each executor. The default value is 2 GB.
--executor_num <Number>
No
Total number of executors to start. The default value is 4.
--executor_core <Number>
No
Number of CPU cores allocated to each executor. The default value is 1.
- Run the following example command:
sh md5compare.sh --dbname tpcds_hive_spark2x_2 --query-path /opt/query/q1.sql --spark_tbl_location /tmp/spark --native_tbl_location /tmp/native --spark_tbl_name spark_data --native_tbl_name native_data --clear_tbl_data true --driver_mem 8G --executor_mem 4G --executor_num 2 --executor_core 2
Because Spark and the Native engine may produce slightly different decimal values for decimal, float, and double types, the consistency check tool compares values rounded to one decimal place. For example, the consistency check tool rounds a value like 3.14159 to 3.1 for comparison.
Differences
After the Spark Native engine is enabled, execution results may differ slightly from the standard Spark runtime.
- Rounding rules for decimal values differ between the standard Spark runtime and the Spark Native engine. Consequently, the final digit of a decimal result may vary.
- Spark performs Base64 encoding when writing binary data to a JSON file (for example, '0101' becomes 'MDEwMQ=='). While the standard Spark runtime automatically decodes this data during a read operation, the Spark Native engine does not. To read this data correctly in the Spark Native engine, you must explicitly use the unbase64() function when reading binary data from a JSON file.
- The Spark Native engine does not support the SET TIME ZONE INTERVAL syntax.
Common Issues
- Server configuration: Log in to MRS Manager, choose Cluster > Services > Spark, click Configurations and then All Configurations, choose Spark(Service) > Native, and configure the following parameters.
Parameter
Description
Default Value
spark.gluten.sql.columnar.backend.ch.shuffle.preferSpill
Whether to spill data to the disk when the shuffle memory reaches the threshold.
Change the parameter value to true.
false
spark.gluten.sql.columnar.backend.ch.spillThreshold
Shuffle memory threshold. When the shuffle memory exceeds this value, data is spilled to the disk.
You are advised to set this parameter to 128 MB.
0
spark.gluten.sql.columnar.maxBatchSize
Number of rows processed by the shuffle operator per batch. This setting directly impacts the memory consumption of the shuffle operator in the Spark Native engine.
Set this parameter to 1024.
4096
- Client configuration:
Add the following parameters in the {Client installation directory}/Spark/spark/conf/spark-defaults.conf file on the Spark client node:
spark.gluten.sql.columnar.backend.ch.shuffle.preferSpill=true spark.gluten.sql.columnar.backend.ch.spillThreshold=128M spark.gluten.sql.columnar.maxBatchSize=1024
If an SQL task triggers too many fallback nodes, performance may drop below that of the standard Spark runtime. In such cases, use the dynamic switch to disable the Spark Native engine.
set spark.gluten.enabled=false;
Optimization Guidelines
- Specifies the maximum size of a file segment after file splitting. The default value is 64 MB.
spark.gluten.sql.columnar.backend.ch.file.split.size=268435456
- Specifies whether to throw an exception when memory usage exceeds the executor limit. If this parameter is set to false, the Spark Native engine may exceed the limit to optimize performance, but this increases the risk of the executor process being terminated by the OOM killer.
spark.gluten.sql.columnar.backend.ch.throwIfMemoryExceed = false
- Specifies whether to allow precision loss for the var_sample and covar_sample functions.
spark.gluten.sql.columnar.precision.loss.allowed = true
- Specifies whether to enable the decimal overflow check. Disabling this function can improve decimal calculation performance.
spark.gluten.sql.columnar.backend.ch.runtime_settings.decimal_check_overflow=false
Just-in-time compilation
spark.gluten.sql.columnar.backend.ch.runtime_config.compile_expressions=1
What is your overall rating for this page?
Thank you very much for your feedback. We will continue working to improve the documentation.See the reply and handling status in My Cloud VOC.
For any further questions, feel free to contact us through the chatbot.
Chatbot