Updated on 2026-08-29 GMT+08:00

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

  1. 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

  1. Click Save to save the changes and restart expired instances.
  2. Choose Dashboard, click More, and select Download Client to install and use the client.

Configuring Client Parameters

  1. Skip this section if you have already completed the steps in Configuring Server Parameters and installed the new client.
  2. 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

  3. Run the source bigdata_env command in the client directory to update environment variables.
  4. 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:

  1. Enable Spark Native by referring to Configuring Client Parameters.
  2. 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.

  3. 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.

  1. 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.
  2. 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.
  3. The Spark Native engine does not support the SET TIME ZONE INTERVAL syntax.

Common Issues

The Spark Native engine typically consumes more memory than the standard Spark runtime, which can increase the risk of OOM errors. To mitigate this, enable shuffle prefer spill and reduce the batch size.
  1. 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

  2. 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

Memory optimization parameters
  • 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
Decimal overflow check and precision
  • 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