layout | title | nav_order |
---|---|---|
page |
Configuration |
3 |
The following is the list of options that rapids-plugin-4-spark
supports.
On startup use: --conf [conf key]=[conf value]
. For example:
${SPARK_HOME}/bin/spark --jars 'rapids-4-spark_2.12-0.1-SNAPSHOT.jar,cudf-0.14-SNAPSHOT-cuda10.jar' \
--conf spark.plugins=com.nvidia.spark.SQLPlugin \
--conf spark.rapids.sql.incompatibleOps.enabled=true
At runtime use: spark.conf.set("[conf key]", [conf value])
. For example:
scala> spark.conf.set("spark.rapids.sql.incompatibleOps.enabled", true)
All configs can be set on startup, but some configs, especially for shuffle, will not work if they are set at runtime.
The RAPIDS Accelerator for Apache Spark can be further configured to enable or disable specific expressions and to control what parts of the query execute using the GPU or the CPU.
Please leverage the spark.rapids.sql.explain
setting to get
feedback from the plugin as to why parts of a query may not be executing on the GPU.
NOTE: Setting
spark.rapids.sql.incompatibleOps.enabled=true
will enable all the settings in the table below which are not enabled by default due to
incompatibilities.
Name | Description | Default Value | Incompatibilities |
---|---|---|---|
spark.rapids.sql.expression.Abs | absolute value | true | None |
spark.rapids.sql.expression.Acos | inverse cosine | true | None |
spark.rapids.sql.expression.Acosh | inverse hyperbolic cosine | true | None |
spark.rapids.sql.expression.Add | addition | true | None |
spark.rapids.sql.expression.Alias | gives a column a name | true | None |
spark.rapids.sql.expression.And | logical and | true | None |
spark.rapids.sql.expression.AnsiCast | convert a column of one type of data into another type | true | None |
spark.rapids.sql.expression.Asin | inverse sine | true | None |
spark.rapids.sql.expression.Asinh | inverse hyperbolic sine | true | None |
spark.rapids.sql.expression.AtLeastNNonNulls | checks if number of non null/Nan values is greater than a given value | true | None |
spark.rapids.sql.expression.Atan | inverse tangent | true | None |
spark.rapids.sql.expression.Atanh | inverse hyperbolic tangent | true | None |
spark.rapids.sql.expression.AttributeReference | references an input column | true | None |
spark.rapids.sql.expression.BitwiseAnd | Returns the bitwise AND of the operands | true | None |
spark.rapids.sql.expression.BitwiseNot | Returns the bitwise NOT of the operands | true | None |
spark.rapids.sql.expression.BitwiseOr | Returns the bitwise OR of the operands | true | None |
spark.rapids.sql.expression.BitwiseXor | Returns the bitwise XOR of the operands | true | None |
spark.rapids.sql.expression.CaseWhen | CASE WHEN expression | true | None |
spark.rapids.sql.expression.Cast | convert a column of one type of data into another type | true | None |
spark.rapids.sql.expression.Cbrt | cube root | true | None |
spark.rapids.sql.expression.Ceil | ceiling of a number | true | None |
spark.rapids.sql.expression.Coalesce | Returns the first non-null argument if exists. Otherwise, null. | true | None |
spark.rapids.sql.expression.Concat | String Concatenate NO separator | true | None |
spark.rapids.sql.expression.Contains | Contains | true | None |
spark.rapids.sql.expression.Cos | cosine | true | None |
spark.rapids.sql.expression.Cosh | hyperbolic cosine | true | None |
spark.rapids.sql.expression.Cot | Returns the cotangent | true | None |
spark.rapids.sql.expression.CurrentRow$ | Special boundary for a window frame, indicating stopping at the current row | true | None |
spark.rapids.sql.expression.DateAdd | Returns the date that is num_days after start_date | true | None |
spark.rapids.sql.expression.DateDiff | datediff | true | None |
spark.rapids.sql.expression.DateSub | Returns the date that is num_days before start_date | true | None |
spark.rapids.sql.expression.DayOfMonth | get the day of the month from a date or timestamp | true | None |
spark.rapids.sql.expression.Divide | division | true | None |
spark.rapids.sql.expression.EndsWith | Ends With | true | None |
spark.rapids.sql.expression.EqualNullSafe | check if the values are equal including nulls <=> | true | None |
spark.rapids.sql.expression.EqualTo | check if the values are equal | true | None |
spark.rapids.sql.expression.Exp | Euler's number e raised to a power | true | None |
spark.rapids.sql.expression.Expm1 | Euler's number e raised to a power minus 1 | true | None |
spark.rapids.sql.expression.Floor | floor of a number | true | None |
spark.rapids.sql.expression.FromUnixTime | get the String from a unix timestamp | true | None |
spark.rapids.sql.expression.GreaterThan | > operator | true | None |
spark.rapids.sql.expression.GreaterThanOrEqual | >= operator | true | None |
spark.rapids.sql.expression.Hour | Returns the hour component of the string/timestamp. | true | None |
spark.rapids.sql.expression.If | IF expression | true | None |
spark.rapids.sql.expression.In | IN operator | true | None |
spark.rapids.sql.expression.InSet | INSET operator | true | None |
spark.rapids.sql.expression.InitCap | Returns str with the first letter of each word in uppercase. All other letters are in lowercase | false | This is not 100% compatible with the Spark version because in some cases unicode characters change byte width when changing the case. The GPU string conversion does not support these characters. For a full list of unsupported characters see rapidsai/cudf#3132 Spark also only sees the space character as a word deliminator, but this uses more white space characters. |
spark.rapids.sql.expression.InputFileBlockLength | Returns the length of the block being read, or -1 if not available. | true | None |
spark.rapids.sql.expression.InputFileBlockStart | Returns the start offset of the block being read, or -1 if not available. | true | None |
spark.rapids.sql.expression.InputFileName | Returns the name of the file being read, or empty string if not available. | true | None |
spark.rapids.sql.expression.IntegralDivide | division with a integer result | true | None |
spark.rapids.sql.expression.IsNaN | checks if a value is NaN | true | None |
spark.rapids.sql.expression.IsNotNull | checks if a value is not null | true | None |
spark.rapids.sql.expression.IsNull | checks if a value is null | true | None |
spark.rapids.sql.expression.KnownFloatingPointNormalized | tag to prevent redundant normalization | true | None |
spark.rapids.sql.expression.Length | String Character Length | true | None |
spark.rapids.sql.expression.LessThan | < operator | true | None |
spark.rapids.sql.expression.LessThanOrEqual | <= operator | true | None |
spark.rapids.sql.expression.Like | Like | true | None |
spark.rapids.sql.expression.Literal | holds a static value from the query | true | None |
spark.rapids.sql.expression.Log | natural log | true | None |
spark.rapids.sql.expression.Log10 | log base 10 | true | None |
spark.rapids.sql.expression.Log1p | natural log 1 + expr | true | None |
spark.rapids.sql.expression.Log2 | log base 2 | true | None |
spark.rapids.sql.expression.Logarithm | log variable base | true | None |
spark.rapids.sql.expression.Lower | String lowercase operator | false | This is not 100% compatible with the Spark version because in some cases unicode characters change byte width when changing the case. The GPU string conversion does not support these characters. For a full list of unsupported characters see rapidsai/cudf#3132 |
spark.rapids.sql.expression.Minute | Returns the minute component of the string/timestamp. | true | None |
spark.rapids.sql.expression.MonotonicallyIncreasingID | Returns monotonically increasing 64-bit integers. | true | None |
spark.rapids.sql.expression.Month | get the month from a date or timestamp | true | None |
spark.rapids.sql.expression.Multiply | multiplication | true | None |
spark.rapids.sql.expression.NaNvl | evaluates to left iff left is not NaN, right otherwise. |
true | None |
spark.rapids.sql.expression.Not | boolean not operator | true | None |
spark.rapids.sql.expression.Or | logical or | true | None |
spark.rapids.sql.expression.Pmod | pmod | true | None |
spark.rapids.sql.expression.Pow | lhs ^ rhs | true | None |
spark.rapids.sql.expression.Rand | Generate a random column with i.i.d. uniformly distributed values in [0, 1) | true | None |
spark.rapids.sql.expression.RegExpReplace | RegExpReplace | true | None |
spark.rapids.sql.expression.Remainder | remainder or modulo | true | None |
spark.rapids.sql.expression.Rint | Rounds up a double value to the nearest double equal to an integer | true | None |
spark.rapids.sql.expression.RowNumber | Window function that returns the index for the row within the aggregation window | true | None |
spark.rapids.sql.expression.Second | Returns the second component of the string/timestamp. | true | None |
spark.rapids.sql.expression.ShiftLeft | Bitwise shift left (<<) | true | None |
spark.rapids.sql.expression.ShiftRight | Bitwise shift right (>>) | true | None |
spark.rapids.sql.expression.ShiftRightUnsigned | Bitwise unsigned shift right (>>>) | true | None |
spark.rapids.sql.expression.Signum | Returns -1.0, 0.0 or 1.0 as expr is negative, 0 or positive | true | None |
spark.rapids.sql.expression.Sin | sine | true | None |
spark.rapids.sql.expression.Sinh | hyperbolic sine | true | None |
spark.rapids.sql.expression.SortOrder | sort order | true | None |
spark.rapids.sql.expression.SparkPartitionID | Returns the current partition id. | true | None |
spark.rapids.sql.expression.SpecifiedWindowFrame | specification of the width of the group (or "frame") of input rows around which a window function is evaluated | true | None |
spark.rapids.sql.expression.Sqrt | square root | true | None |
spark.rapids.sql.expression.StartsWith | Starts With | true | None |
spark.rapids.sql.expression.StringLocate | Substring search operator | true | None |
spark.rapids.sql.expression.StringReplace | StringReplace operator | true | None |
spark.rapids.sql.expression.StringTrim | StringTrim operator | true | None |
spark.rapids.sql.expression.StringTrimLeft | StringTrimLeft operator | true | None |
spark.rapids.sql.expression.StringTrimRight | StringTrimRight operator | true | None |
spark.rapids.sql.expression.Substring | Substring operator | true | None |
spark.rapids.sql.expression.SubstringIndex | substring_index operator | true | None |
spark.rapids.sql.expression.Subtract | subtraction | true | None |
spark.rapids.sql.expression.Tan | tangent | true | None |
spark.rapids.sql.expression.Tanh | hyperbolic tangent | true | None |
spark.rapids.sql.expression.TimeSub | Subtracts interval from timestamp | true | None |
spark.rapids.sql.expression.ToDegrees | Converts radians to degrees | true | None |
spark.rapids.sql.expression.ToRadians | Converts degrees to radians | true | None |
spark.rapids.sql.expression.ToUnixTimestamp | Returns the UNIX timestamp of the given time | false | This is not 100% compatible with the Spark version because Incorrectly formatted strings and bogus dates produce garbage data instead of null |
spark.rapids.sql.expression.UnaryMinus | negate a numeric value | true | None |
spark.rapids.sql.expression.UnaryPositive | a numeric value with a + in front of it | true | None |
spark.rapids.sql.expression.UnboundedFollowing$ | Special boundary for a window frame, indicating all rows preceding the current row | true | None |
spark.rapids.sql.expression.UnboundedPreceding$ | Special boundary for a window frame, indicating all rows preceding the current row | true | None |
spark.rapids.sql.expression.UnixTimestamp | Returns the UNIX timestamp of current or specified time | false | This is not 100% compatible with the Spark version because Incorrectly formatted strings and bogus dates produce garbage data instead of null |
spark.rapids.sql.expression.Upper | String uppercase operator | false | This is not 100% compatible with the Spark version because in some cases unicode characters change byte width when changing the case. The GPU string conversion does not support these characters. For a full list of unsupported characters see rapidsai/cudf#3132 |
spark.rapids.sql.expression.WindowExpression | calculates a return value for every input row of a table based on a group (or "window") of rows | true | None |
spark.rapids.sql.expression.WindowSpecDefinition | specification of a window function, indicating the partitioning-expression, the row ordering, and the width of the window | true | None |
spark.rapids.sql.expression.Year | get the year from a date or timestamp | true | None |
spark.rapids.sql.expression.AggregateExpression | aggregate expression | true | None |
spark.rapids.sql.expression.Average | average aggregate operator | true | None |
spark.rapids.sql.expression.Count | count aggregate operator | true | None |
spark.rapids.sql.expression.First | first aggregate operator | true | None |
spark.rapids.sql.expression.Last | last aggregate operator | true | None |
spark.rapids.sql.expression.Max | max aggregate operator | true | None |
spark.rapids.sql.expression.Min | min aggregate operator | true | None |
spark.rapids.sql.expression.Sum | sum aggregate operator | true | None |
spark.rapids.sql.expression.NormalizeNaNAndZero | normalize nan and zero | true | None |
CUDF can compile GPU kernels at runtime using a just-in-time (JIT) compiler. The
resulting kernels are cached on the filesystem. The default location for this cache is
under the .cudf
directory in the user's home directory. When running in an environment
where the user's home directory cannot be written, such as running in a container
environment on a cluster, the JIT cache path will need to be specified explicitly with
the LIBCUDF_KERNEL_CACHE_PATH
environment variable.
The specified kernel cache path should be specific to the user to avoid conflicts with
others running on the same host. For example, the following would specify the path to a
user-specific location under /tmp
:
--conf spark.executorEnv.LIBCUDF_KERNEL_CACHE_PATH="/tmp/cudf-$USER"