Zand100 opened a new issue, #8724:
URL: https://github.com/apache/incubator-gluten/issues/8724

   ### Backend
   
   VL (Velox)
   
   ### Bug description
   
   [Expected behavior] and [actual behavior].
   When we took Gluten 1.2.1, Gluten could call a Java UDF that (correctly) 
threw an exception, and handle it gracefully. We tried to take 1.3.0, and it 
gives an error `org.apache.spark.SparkException: [FAILED_EXECUTE_UDF] Failed to 
execute user defined function`
   
   I ran into this while running integration tests in a more complex setup and 
that stacktrace included the error below. But I believe it wasn't actually a 
memory leak, because it would hang for a few seconds before that error. The 
expected behavior is that after the exception is thrown, the code to catch the 
exception executes successfully (and the test passes). In 1.3.0, it stops after 
the exception is thrown. 
   
   `project/ep/build-velox/build/velox_ep/velox/exec/Driver.cpp:601, 
Function:operator(), Expression:  Operator::getOutput failed for [operator: 
ValueStream, plan node ID: 0]: Error during calling Java code from native code: 
org.apache.spark.SparkException: [FAILED_EXECUTE_UDF] Failed to execute user 
defined function (`functions$$$Lambda/0x00007f0049900ba0`: (string, string) => 
boolean).`
   ...
   
   `, Source: RUNTIME, ErrorCode: INVALID_STATE
   E20250211 22:55:13.386138 19912 VeloxMemoryManager.cc:401] Failed to release 
Velox memory manager after 43350ms as there are still outstanding memory 
resources. 
   E20250211 22:55:13.386179 19912 MemoryPool.cpp:442] [MEM] Memory leak (Used 
memory): Memory Pool[default_leaf LEAF root[root] parent[root] MALLOC 
track-usage thread-safe]<unlimited max capacity capacity 8.00MB used 512B 
available 1023.50KB reservation [used 512B, reserved 1.00MB, min 0B] counters 
[allocs 4, frees 0, reserves 0, releases 0, collisions 0])>
   E20250211 22:55:13.386235 19912 Exceptions.h:66] Line: 
/home/circleci/project/cpp/velox/memory/VeloxMemoryManager.cc:102, 
Function:removePool, Expression: pool->reservedBytes() == 0 (1048576 vs. 0), 
Source: RUNTIME, ErrorCode: INVALID_STATE
   terminate called after throwing an instance of 
'facebook::velox::VeloxRuntimeError'`
   
   I think this is a minimal repro. It's TPCDS q9.
   
   ```
   export 
gluten_jar=~/gluten-velox-bundle-spark3.5_2.12-centos_7_x86_64-1.3.0.jar
   export SPARK_HOME=~/spark-3.5.4-bin-hadoop3
   export PATH=$PATH:/$SPARK_HOME/bin
   
   spark-shell   --conf spark.plugins=org.apache.gluten.GlutenPlugin   --conf 
spark.memory.offHeap.enabled=true   --conf spark.memory.offHeap.size=20g   
--conf spark.driver.extraClassPath=${gluten_jar}   --conf 
spark.executor.extraClassPath=${gluten_jar}   --conf 
spark.shuffle.manager=org.apache.spark.shuffle.sort.ColumnarShuffleManager   
--jars ${gluten_jar} --conf spark.gluten.enabled=true --conf 
spark.driver.extraJavaOptions="-Dio.netty.tryReflectionSetAccessible=true"
   
   scala> val storeSalesDF = spark.read.format("csv").option("header", 
"true").load("/home/coder/store_sales.csv")
   scala> val reasonsDF = spark.read.format("csv").option("header", 
"true").load("/home/coder/reason.csv")
   scala> storeSalesDF.createOrReplaceTempView("store_sales")
   scala>  reasonsDF.createOrReplaceTempView("reason")
   
   scala> val resultDF = spark.sql("""
   
   
   
        |      |   SELECT CASE 
   
        |      |            WHEN (SELECT Count(*) 
   
   
        |      |                  FROM   store_sales 
   
   
        |      |                  WHERE  ss_quantity BETWEEN 1 AND 20) > 3672 
THEN 
   
   
        |      |            (SELECT Avg(ss_ext_list_price) 
   
   
        |      |             FROM   store_sales 
   
   
        |      |             WHERE  ss_quantity BETWEEN 1 AND 20) 
   
   
        |      |            ELSE (SELECT Avg(ss_net_profit) 
   
   
        |      |                  FROM   store_sales 
   
   
        |      |                  WHERE  ss_quantity BETWEEN 1 AND 20) 
   
   
        |      |          END bucket1, 
   
   
        |      |          CASE 
   
   
        |      |            WHEN (SELECT Count(*) 
   
   
        |      |                  FROM   store_sales 
   
   
        |      |                  WHERE  ss_quantity BETWEEN 21 AND 40) > 3392 
THEN 
   
   
        |      |            (SELECT Avg(ss_ext_list_price) 
   
   
        |      |             FROM   store_sales 
   
   
        |      |             WHERE  ss_quantity BETWEEN 21 AND 40) 
   
   
        |      |            ELSE (SELECT Avg(ss_net_profit) 
   
   
        |      |                  FROM   store_sales 
   
   
        |      |                  WHERE  ss_quantity BETWEEN 21 AND 40) 
   
   
        |      |          END bucket2, 
   
   
        |      |          CASE 
   
   
        |      |            WHEN (SELECT Count(*) 
   
   
        |      |                  FROM   store_sales 
   
   
        |      |                  WHERE  ss_quantity BETWEEN 41 AND 60) > 32784 
THEN 
   
   
        |      |            (SELECT Avg(ss_ext_list_price) 
   
   
        |      |             FROM   store_sales 
   
   
        |      |             WHERE  ss_quantity BETWEEN 41 AND 60) 
   
   
        |      |            ELSE (SELECT Avg(ss_net_profit) 
   
   
        |      |                  FROM   store_sales 
   
   
        |      |                  WHERE  ss_quantity BETWEEN 41 AND 60) 
   
   
        |      |          END bucket3, 
   
   
        |      |          CASE 
   
   
        |      |            WHEN (SELECT Count(*) 
   
   
        |      |                  FROM   store_sales 
   
   
        |      |                  WHERE  ss_quantity BETWEEN 61 AND 80) > 26032 
THEN 
   
   
        |      |            (SELECT Avg(ss_ext_list_price) 
   
   
        |      |             FROM   store_sales 
   
   
        |      |             WHERE  ss_quantity BETWEEN 61 AND 80) 
   
   
        |      |            ELSE (SELECT Avg(ss_net_profit) 
   
   
        |      |                  FROM   store_sales 
   
   
        |      |                  WHERE  ss_quantity BETWEEN 61 AND 80) 
   
   
        |      |          END bucket4, 
   
   
        |      |          CASE 
   
   
        |      |            WHEN (SELECT Count(*) 
   
   
        |      |                  FROM   store_sales 
   
   
        |      |                  WHERE  ss_quantity BETWEEN 81 AND 100) > 
23982 THEN 
   
   
        |      |            (SELECT Avg(ss_ext_list_price) 
   
   
        |      |             FROM   store_sales 
   
   
        |      |             WHERE  ss_quantity BETWEEN 81 AND 100) 
   
   
        |      |            ELSE (SELECT Avg(ss_net_profit) 
   
   
        |      |                  FROM   store_sales 
   
   
        |      |                  WHERE  ss_quantity BETWEEN 81 AND 100) 
   
   
        |      |          END bucket5 
   
   
        |      |   FROM   reason 
   
   
        |      |   WHERE  r_reason_sk = 1
   
   
        |      | """)
   
   scala> def throwException() = {throw new IllegalStateException("test error")}
   throwException: ()Nothing
   
   scala> val throwExceptionUdf = udf(throwException _)
   throwExceptionUdf: org.apache.spark.sql.expressions.UserDefinedFunction = 
SparkUserDefinedFunction($Lambda$6066/0x00007fa6ed56a3f8@1c650d81,NullType,List(),Some(class[value[0]:
 void]),None,true,true)
   
   scala>  val resultDF2 = resultDF.withColumn("exception", throwExceptionUdf())
   resultDF2: org.apache.spark.sql.DataFrame = [bucket1: double, bucket2: 
double ... 4 more fields]
   
   scala> resultDF2.show()
   
   ```
   
   
   
   
   
   
   
   
   
   ### Spark version
   
   None
   
   ### Spark configurations
   
   Spark 3.5.4 
   
   I downloaded Spark and Gluten from 
https://gluten.apache.org/#31-use-a-pre-built-jar  
   
   `spark-shell   --conf spark.plugins=org.apache.gluten.GlutenPlugin   --conf 
spark.memory.offHeap.enabled=true   --conf spark.memory.offHeap.size=20g   
--conf spark.driver.extraClassPath=${gluten_jar}   --conf 
spark.executor.extraClassPath=${gluten_jar}   --conf 
spark.shuffle.manager=org.apache.spark.shuffle.sort.ColumnarShuffleManager   
--jars ${gluten_jar} --conf spark.gluten.enabled=true --conf 
spark.driver.extraJavaOptions="-Dio.netty.tryReflectionSetAccessible=true"`
   
   ### System information
   
   Ubuntu 20.04
   
   ### Relevant logs
   
   ```bash
   25/02/13 02:46:51 WARN ProjectExecTransformer: Validation failed with 
exception for plan: ProjectExecTransformer, due to: Not supported to map spark 
function name to substrait function name: toprettystring(CASE WHEN (Subquery 
subquery#86, [id=#1775].count(1) > 3672) THEN ReusedSubquery Subquery 
subquery#86, [id=#1775].avg(ss_ext_list_price) ELSE ReusedSubquery Subquery 
subquery#86, [id=#1775].avg(ss_net_profit) END, Some(Etc/UTC)), class name: 
ToPrettyString.
   
   
   25/02/13 02:46:51 WARN ProjectExecTransformer: Validation failed with 
exception for plan: ProjectExecTransformer, due to: Not supported to map spark 
function name to substrait function name: toprettystring(CASE WHEN (Subquery 
subquery#86, [id=#1775].count(1) > 3672) THEN ReusedSubquery Subquery 
subquery#86, [id=#1775].avg(ss_ext_list_price) ELSE ReusedSubquery Subquery 
subquery#86, [id=#1775].avg(ss_net_profit) END, Some(Etc/UTC)), class name: 
ToPrettyString.
   
   
   25/02/13 02:46:51 WARN GlutenFallbackReporter: Validation failed for plan: 
Scan csv [QueryId=5], due to: Unsupported file format TextReadFormat..
   
   
   25/02/13 02:46:51 WARN GlutenFallbackReporter: Validation failed for plan: 
Project[QueryId=5], due to: Not supported to map spark function name to 
substrait function name: toprettystring(CASE WHEN (Subquery subquery#86, 
[id=#1775].count(1) > 3672) THEN ReusedSubquery Subquery subquery#86, 
[id=#1775].avg(ss_ext_list_price) ELSE ReusedSubquery Subquery subquery#86, 
[id=#1775].avg(ss_net_profit) END, Some(Etc/UTC)), class name: ToPrettyString..
   
   
   25/02/13 02:46:51 WARN GlutenFallbackReporter: Validation failed for plan: 
Scan csv [QueryId=5], due to: Unsupported file format TextReadFormat..
   
   
   25/02/13 02:46:51 WARN GlutenFallbackReporter: Validation failed for plan: 
Scan csv [QueryId=5], due to: Unsupported file format TextReadFormat..
   
   
   25/02/13 02:46:51 WARN GlutenFallbackReporter: Validation failed for plan: 
Scan csv [QueryId=5], due to: Unsupported file format TextReadFormat..
   
   
   25/02/13 02:46:51 WARN GlutenFallbackReporter: Validation failed for plan: 
Scan csv [QueryId=5], due to: Unsupported file format TextReadFormat..
   
   
   25/02/13 02:46:51 WARN GlutenFallbackReporter: Validation failed for plan: 
Scan csv [QueryId=5], due to: Unsupported file format TextReadFormat..
   
   
   25/02/13 02:46:52 ERROR TaskResources: Task 23 failed by error: 
   
   
   org.apache.spark.SparkException: [FAILED_EXECUTE_UDF] Failed to execute user 
defined function 
(`$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$Lambda$6066/0x00007fa6ed56a3f8`: () => 
void).
   
   
       at 
org.apache.spark.sql.errors.QueryExecutionErrors$.failedExecuteUserDefinedFunctionError(QueryExecutionErrors.scala:198)
   
   
       at 
org.apache.spark.sql.errors.QueryExecutionErrors.failedExecuteUserDefinedFunctionError(QueryExecutionErrors.scala)
   
   
       at 
org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage1.processNext(Unknown
 Source)
   
   
       at 
org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
   
   
       at 
org.apache.spark.sql.execution.WholeStageCodegenEvaluatorFactory$WholeStageCodegenPartitionEvaluator$$anon$1.hasNext(WholeStageCodegenEvaluatorFactory.scala:43)
   
   
       at 
org.apache.spark.sql.execution.SparkPlan.$anonfun$getByteArrayRdd$1(SparkPlan.scala:388)
   
   
       at 
org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2(RDD.scala:893)
   
   
       at 
org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2$adapted(RDD.scala:893)
   
   
       at 
org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
   
   
       at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:367)
   
   
       at org.apache.spark.rdd.RDD.iterator(RDD.scala:331)
   
   
       at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:93)
   
   
       at 
org.apache.spark.TaskContext.runTaskWithListeners(TaskContext.scala:166)
   
   
       at org.apache.spark.scheduler.Task.run(Task.scala:141)
   
   
       at 
org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$4(Executor.scala:620)
   
   
       at 
org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally(SparkErrorUtils.scala:64)
   
   
       at 
org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally$(SparkErrorUtils.scala:61)
   
   
       at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:94)
   
   
       at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:623)
   
   
       at 
java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)
   
   
       at 
java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635)
   
   
       at java.base/java.lang.Thread.run(Thread.java:840)
   
   
   Caused by: java.lang.IllegalStateException: test error
   
   
       at 
$line24.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw.throwException(<console>:22)
   
   
       at 
$line25.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw.$anonfun$throwExceptionUdf$1(<console>:23)
   
   
       ... 20 more
   
   
   E20250213 02:46:52.092774 97905 Exceptions.h:66] Line: 
/root/src/weiting/gluten/ep/build-velox/build/velox_ep/velox/exec/Task.cpp:2039,
 Function:terminate, Expression:  Cancelled, Source: RUNTIME, ErrorCode: 
INVALID_STATE
   
   
   25/02/13 02:46:52 ERROR Executor: Exception in task 0.0 in stage 33.0 (TID 
23)
   
   
   org.apache.spark.SparkException: [FAILED_EXECUTE_UDF] Failed to execute user 
defined function 
(`$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$Lambda$6066/0x00007fa6ed56a3f8`: () => 
void).
   
   
       at 
org.apache.spark.sql.errors.QueryExecutionErrors$.failedExecuteUserDefinedFunctionError(QueryExecutionErrors.scala:198)
   
   
       at 
org.apache.spark.sql.errors.QueryExecutionErrors.failedExecuteUserDefinedFunctionError(QueryExecutionErrors.scala)
   
   
       at 
org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage1.processNext(Unknown
 Source)
   
   
       at 
org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
   
   
       at 
org.apache.spark.sql.execution.WholeStageCodegenEvaluatorFactory$WholeStageCodegenPartitionEvaluator$$anon$1.hasNext(WholeStageCodegenEvaluatorFactory.scala:43)
   
   
       at 
org.apache.spark.sql.execution.SparkPlan.$anonfun$getByteArrayRdd$1(SparkPlan.scala:388)
   
   
       at 
org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2(RDD.scala:893)
   
   
       at 
org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2$adapted(RDD.scala:893)
   
   
       at 
org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
   
   
       at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:367)
   
   
       at org.apache.spark.rdd.RDD.iterator(RDD.scala:331)
   
   
       at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:93)
   
   
       at 
org.apache.spark.TaskContext.runTaskWithListeners(TaskContext.scala:166)
   
   
       at org.apache.spark.scheduler.Task.run(Task.scala:141)
   
   
       at 
org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$4(Executor.scala:620)
   
   
       at 
org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally(SparkErrorUtils.scala:64)
   
   
       at 
org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally$(SparkErrorUtils.scala:61)
   
   
       at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:94)
   
   
       at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:623)
   
   
       at 
java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)
   
   
       at 
java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635)
   
   
       at java.base/java.lang.Thread.run(Thread.java:840)
   
   
   Caused by: java.lang.IllegalStateException: test error
   
   
       at 
$line24.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw.throwException(<console>:22)
   
   
       at 
$line25.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw.$anonfun$throwExceptionUdf$1(<console>:23)
   
   
       ... 20 more
   
   
   25/02/13 02:46:52 WARN TaskSetManager: Lost task 0.0 in stage 33.0 (TID 23) 
(za-gluten-repro-udf2 executor driver): org.apache.spark.SparkException: 
[FAILED_EXECUTE_UDF] Failed to execute user defined function 
(`$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$Lambda$6066/0x00007fa6ed56a3f8`: () => 
void).
   
   
       at 
org.apache.spark.sql.errors.QueryExecutionErrors$.failedExecuteUserDefinedFunctionError(QueryExecutionErrors.scala:198)
   
   
       at 
org.apache.spark.sql.errors.QueryExecutionErrors.failedExecuteUserDefinedFunctionError(QueryExecutionErrors.scala)
   
   
       at 
org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage1.processNext(Unknown
 Source)
   
   
       at 
org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
   
   
       at 
org.apache.spark.sql.execution.WholeStageCodegenEvaluatorFactory$WholeStageCodegenPartitionEvaluator$$anon$1.hasNext(WholeStageCodegenEvaluatorFactory.scala:43)
   
   
       at 
org.apache.spark.sql.execution.SparkPlan.$anonfun$getByteArrayRdd$1(SparkPlan.scala:388)
   
   
       at 
org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2(RDD.scala:893)
   
   
       at 
org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2$adapted(RDD.scala:893)
   
   
       at 
org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
   
   
       at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:367)
   
   
       at org.apache.spark.rdd.RDD.iterator(RDD.scala:331)
   
   
       at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:93)
   
   
       at 
org.apache.spark.TaskContext.runTaskWithListeners(TaskContext.scala:166)
   
   
       at org.apache.spark.scheduler.Task.run(Task.scala:141)
   
   
       at 
org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$4(Executor.scala:620)
   
   
       at 
org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally(SparkErrorUtils.scala:64)
   
   
       at 
org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally$(SparkErrorUtils.scala:61)
   
   
       at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:94)
   
   
       at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:623)
   
   
       at 
java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)
   
   
       at 
java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635)
   
   
       at java.base/java.lang.Thread.run(Thread.java:840)
   
   
   Caused by: java.lang.IllegalStateException: test error
   
   
       at 
$line24.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw.throwException(<console>:22)
   
   
       at 
$line25.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw.$anonfun$throwExceptionUdf$1(<console>:23)
   
   
       ... 20 more
   
   
   
   
   
   
   25/02/13 02:46:52 ERROR TaskSetManager: Task 0 in stage 33.0 failed 1 times; 
aborting job
   
   
   org.apache.spark.SparkException: Job aborted due to stage failure: Task 0 in 
stage 33.0 failed 1 times, most recent failure: Lost task 0.0 in stage 33.0 
(TID 23) (za-gluten-repro-udf2 executor driver): 
org.apache.spark.SparkException: [FAILED_EXECUTE_UDF] Failed to execute user 
defined function (`$Lambda$6066/0x00007fa6ed56a3f8`: () => void).
   
   
       at 
org.apache.spark.sql.errors.QueryExecutionErrors$.failedExecuteUserDefinedFunctionError(QueryExecutionErrors.scala:198)
   
   
       at 
org.apache.spark.sql.errors.QueryExecutionErrors.failedExecuteUserDefinedFunctionError(QueryExecutionErrors.scala)
   
   
       at 
org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage1.processNext(Unknown
 Source)
   
   
       at 
org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
   
   
       at 
org.apache.spark.sql.execution.WholeStageCodegenEvaluatorFactory$WholeStageCodegenPartitionEvaluator$$anon$1.hasNext(WholeStageCodegenEvaluatorFactory.scala:43)
   
   
       at 
org.apache.spark.sql.execution.SparkPlan.$anonfun$getByteArrayRdd$1(SparkPlan.scala:388)
   
   
       at 
org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2(RDD.scala:893)
   
   
       at 
org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2$adapted(RDD.scala:893)
   
   
       at 
org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
   
   
       at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:367)
   
   
       at org.apache.spark.rdd.RDD.iterator(RDD.scala:331)
   
   
       at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:93)
   
   
       at 
org.apache.spark.TaskContext.runTaskWithListeners(TaskContext.scala:166)
   
   
       at org.apache.spark.scheduler.Task.run(Task.scala:141)
   
   
       at 
org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$4(Executor.scala:620)
   
   
       at 
org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally(SparkErrorUtils.scala:64)
   
   
       at 
org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally$(SparkErrorUtils.scala:61)
   
   
       at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:94)
   
   
       at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:623)
   
   
       at 
java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)
   
   
       at 
java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635)
   
   
       at java.base/java.lang.Thread.run(Thread.java:840)
   
   
   Caused by: java.lang.IllegalStateException: test error
   
   
       at throwException(<console>:22)
   
   
       at $anonfun$throwExceptionUdf$1(<console>:23)
   
   
       ... 20 more
   
   
   
   
   
   
   Driver stacktrace:
   
   
     at 
org.apache.spark.scheduler.DAGScheduler.failJobAndIndependentStages(DAGScheduler.scala:2856)
   
   
     at 
org.apache.spark.scheduler.DAGScheduler.$anonfun$abortStage$2(DAGScheduler.scala:2792)
   
   
     at 
org.apache.spark.scheduler.DAGScheduler.$anonfun$abortStage$2$adapted(DAGScheduler.scala:2791)
   
   
     at scala.collection.mutable.ResizableArray.foreach(ResizableArray.scala:62)
   
   
     at 
scala.collection.mutable.ResizableArray.foreach$(ResizableArray.scala:55)
   
   
     at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:49)
   
   
     at 
org.apache.spark.scheduler.DAGScheduler.abortStage(DAGScheduler.scala:2791)
   
   
     at 
org.apache.spark.scheduler.DAGScheduler.$anonfun$handleTaskSetFailed$1(DAGScheduler.scala:1247)
   
   
     at 
org.apache.spark.scheduler.DAGScheduler.$anonfun$handleTaskSetFailed$1$adapted(DAGScheduler.scala:1247)
   
   
     at scala.Option.foreach(Option.scala:407)
   
   
     at 
org.apache.spark.scheduler.DAGScheduler.handleTaskSetFailed(DAGScheduler.scala:1247)
   
   
     at 
org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:3060)
   
   
     at 
org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2994)
   
   
     at 
org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2983)
   
   
     at org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:49)
   
   
     at org.apache.spark.scheduler.DAGScheduler.runJob(DAGScheduler.scala:989)
   
   
     at org.apache.spark.SparkContext.runJob(SparkContext.scala:2393)
   
   
     at org.apache.spark.SparkContext.runJob(SparkContext.scala:2414)
   
   
     at org.apache.spark.SparkContext.runJob(SparkContext.scala:2433)
   
   
     at 
org.apache.spark.sql.execution.SparkPlan.executeTake(SparkPlan.scala:530)
   
   
     at 
org.apache.spark.sql.execution.SparkPlan.executeTake(SparkPlan.scala:483)
   
   
     at 
org.apache.spark.sql.execution.CollectLimitExec.executeCollect(limit.scala:61)
   
   
     at 
org.apache.spark.sql.execution.adaptive.AdaptiveSparkPlanExec.$anonfun$executeCollect$1(AdaptiveSparkPlanExec.scala:392)
   
   
     at 
org.apache.spark.sql.execution.adaptive.AdaptiveSparkPlanExec.withFinalPlanUpdate(AdaptiveSparkPlanExec.scala:420)
   
   
     at 
org.apache.spark.sql.execution.adaptive.AdaptiveSparkPlanExec.executeCollect(AdaptiveSparkPlanExec.scala:392)
   
   
     at org.apache.spark.sql.Dataset.collectFromPlan(Dataset.scala:4333)
   
   
     at org.apache.spark.sql.Dataset.$anonfun$head$1(Dataset.scala:3316)
   
   
     at org.apache.spark.sql.Dataset.$anonfun$withAction$2(Dataset.scala:4323)
   
   
     at 
org.apache.spark.sql.execution.QueryExecution$.withInternalError(QueryExecution.scala:546)
   
   
     at org.apache.spark.sql.Dataset.$anonfun$withAction$1(Dataset.scala:4321)
   
   
     at 
org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$6(SQLExecution.scala:125)
   
   
     at 
org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:201)
   
   
     at 
org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$1(SQLExecution.scala:108)
   
   
     at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:900)
   
   
     at 
org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:66)
   
   
     at org.apache.spark.sql.Dataset.withAction(Dataset.scala:4321)
   
   
     at org.apache.spark.sql.Dataset.head(Dataset.scala:3316)
   
   
     at org.apache.spark.sql.Dataset.take(Dataset.scala:3539)
   
   
     at org.apache.spark.sql.Dataset.getRows(Dataset.scala:280)
   
   
     at org.apache.spark.sql.Dataset.showString(Dataset.scala:315)
   
   
     at org.apache.spark.sql.Dataset.show(Dataset.scala:838)
   
   
     at org.apache.spark.sql.Dataset.show(Dataset.scala:797)
   
   
     at org.apache.spark.sql.Dataset.show(Dataset.scala:806)
   
   
     ... 47 elided
   
   
   Caused by: org.apache.spark.SparkException: [FAILED_EXECUTE_UDF] Failed to 
execute user defined function (`$Lambda$6066/0x00007fa6ed56a3f8`: () => void).
   
   
     at 
org.apache.spark.sql.errors.QueryExecutionErrors$.failedExecuteUserDefinedFunctionError(QueryExecutionErrors.scala:198)
   
   
     at 
org.apache.spark.sql.errors.QueryExecutionErrors.failedExecuteUserDefinedFunctionError(QueryExecutionErrors.scala)
   
   
     at 
org.apache.spark.sql.catalyst.expressions.GeneratedClass$GeneratedIteratorForCodegenStage1.processNext(Unknown
 Source)
   
   
     at 
org.apache.spark.sql.execution.BufferedRowIterator.hasNext(BufferedRowIterator.java:43)
   
   
     at 
org.apache.spark.sql.execution.WholeStageCodegenEvaluatorFactory$WholeStageCodegenPartitionEvaluator$$anon$1.hasNext(WholeStageCodegenEvaluatorFactory.scala:43)
   
   
     at 
org.apache.spark.sql.execution.SparkPlan.$anonfun$getByteArrayRdd$1(SparkPlan.scala:388)
   
   
     at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2(RDD.scala:893)
   
   
     at 
org.apache.spark.rdd.RDD.$anonfun$mapPartitionsInternal$2$adapted(RDD.scala:893)
   
   
     at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
   
   
     at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:367)
   
   
     at org.apache.spark.rdd.RDD.iterator(RDD.scala:331)
   
   
     at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:93)
   
   
     at org.apache.spark.TaskContext.runTaskWithListeners(TaskContext.scala:166)
   
   
     at org.apache.spark.scheduler.Task.run(Task.scala:141)
   
   
     at 
org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$4(Executor.scala:620)
   
   
     at 
org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally(SparkErrorUtils.scala:64)
   
   
     at 
org.apache.spark.util.SparkErrorUtils.tryWithSafeFinally$(SparkErrorUtils.scala:61)
   
   
     at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:94)
   
   
     at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:623)
   
   
     at 
java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136)
   
   
     at 
java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635)
   
   
     at java.base/java.lang.Thread.run(Thread.java:840)
   
   
   Caused by: java.lang.IllegalStateException: test error
   
   
     at throwException(<console>:22)
   
   
     at $anonfun$throwExceptionUdf$1(<console>:23)
   
   
     ... 20 more
   
   
   
   
   
   
   scala> resultDF2.explain()
   
   
   == Physical Plan ==
   
   
   AdaptiveSparkPlan isFinalPlan=false
   
   
   +- Project [CASE WHEN (Subquery subquery#86, [id=#2896].count(1) > 3672) 
THEN Subquery subquery#87, [id=#2897].avg(ss_ext_list_price) ELSE Subquery 
subquery#88, [id=#2898].avg(ss_net_profit) END AS bucket1#89, CASE WHEN 
(Subquery subquery#90, [id=#2899].count(1) > 3392) THEN Subquery subquery#91, 
[id=#2900].avg(ss_ext_list_price) ELSE Subquery subquery#92, 
[id=#2901].avg(ss_net_profit) END AS bucket2#93, CASE WHEN (Subquery 
subquery#94, [id=#2902].count(1) > 32784) THEN Subquery subquery#95, 
[id=#2903].avg(ss_ext_list_price) ELSE Subquery subquery#96, 
[id=#2904].avg(ss_net_profit) END AS bucket3#97, CASE WHEN (Subquery 
subquery#98, [id=#2905].count(1) > 26032) THEN Subquery subquery#99, 
[id=#2906].avg(ss_ext_list_price) ELSE Subquery subquery#100, 
[id=#2907].avg(ss_net_profit) END AS bucket4#101, CASE WHEN (Subquery 
subquery#102, [id=#2908].count(1) > 23982) THEN Subquery subquery#103, 
[id=#2909].avg(ss_ext_list_price) ELSE Subquery subquery#104, 
[id=#2910].avg(ss_net_profit) END 
 AS bucket5#105, UDF() AS exception#826]
   
   
      :  :- Subquery subquery#86, [id=#2896]
   
   
      :  :  +- AdaptiveSparkPlan isFinalPlan=false
   
   
      :  :     +- Project [named_struct(count(1), count(1)#107L, 
avg(ss_ext_list_price), avg(ss_ext_list_price)#109, avg(ss_net_profit), 
avg(ss_net_profit)#111) AS mergedValue#1188]
   
   
      :  :        +- HashAggregate(keys=[], functions=[count(1), 
avg(cast(ss_ext_list_price#33 as double)), avg(cast(ss_net_profit#38 as 
double))])
   
   
      :  :           +- Exchange SinglePartition, ENSURE_REQUIREMENTS, 
[plan_id=2669]
   
   
      :  :              +- HashAggregate(keys=[], functions=[partial_count(1), 
partial_avg(cast(ss_ext_list_price#33 as double)), 
partial_avg(cast(ss_net_profit#38 as double))])
   
   
      :  :                 +- Project [ss_ext_list_price#33, ss_net_profit#38]
   
   
      :  :                    +- Filter ((isnotnull(ss_quantity#26) AND 
(cast(ss_quantity#26 as int) >= 1)) AND (cast(ss_quantity#26 as int) <= 20))
   
   
      :  :                       +- FileScan csv 
[ss_quantity#26,ss_ext_list_price#33,ss_net_profit#38] Batched: false, 
DataFilters: [isnotnull(ss_quantity#26), (cast(ss_quantity#26 as int) >= 1), 
(cast(ss_quantity#26 as int) <= 20)], Format: CSV, Location: 
InMemoryFileIndex(1 paths)[file:/home/coder/store_sales.csv], PartitionFilters: 
[], PushedFilters: [IsNotNull(ss_quantity)], ReadSchema: 
struct<ss_quantity:string,ss_ext_list_price:string,ss_net_profit:string>
   
   
      :  :- Subquery subquery#87, [id=#2897]
   
   
      :  :  +- AdaptiveSparkPlan isFinalPlan=false
   
   
      :  :     +- Project [named_struct(count(1), count(1)#107L, 
avg(ss_ext_list_price), avg(ss_ext_list_price)#109, avg(ss_net_profit), 
avg(ss_net_profit)#111) AS mergedValue#1188]
   
   
      :  :        +- HashAggregate(keys=[], functions=[count(1), 
avg(cast(ss_ext_list_price#33 as double)), avg(cast(ss_net_profit#38 as 
double))])
   
   
      :  :           +- Exchange SinglePartition, ENSURE_REQUIREMENTS, 
[plan_id=2685]
   
   
      :  :              +- HashAggregate(keys=[], functions=[partial_count(1), 
partial_avg(cast(ss_ext_list_price#33 as double)), 
partial_avg(cast(ss_net_profit#38 as double))])
   
   
      :  :                 +- Project [ss_ext_list_price#33, ss_net_profit#38]
   
   
      :  :                    +- Filter ((isnotnull(ss_quantity#26) AND 
(cast(ss_quantity#26 as int) >= 1)) AND (cast(ss_quantity#26 as int) <= 20))
   
   
      :  :                       +- FileScan csv 
[ss_quantity#26,ss_ext_list_price#33,ss_net_profit#38] Batched: false, 
DataFilters: [isnotnull(ss_quantity#26), (cast(ss_quantity#26 as int) >= 1), 
(cast(ss_quantity#26 as int) <= 20)], Format: CSV, Location: 
InMemoryFileIndex(1 paths)[file:/home/coder/store_sales.csv], PartitionFilters: 
[], PushedFilters: [IsNotNull(ss_quantity)], ReadSchema: 
struct<ss_quantity:string,ss_ext_list_price:string,ss_net_profit:string>
   
   
      :  :- Subquery subquery#88, [id=#2898]
   
   
      :  :  +- AdaptiveSparkPlan isFinalPlan=false
   
   
      :  :     +- Project [named_struct(count(1), count(1)#107L, 
avg(ss_ext_list_price), avg(ss_ext_list_price)#109, avg(ss_net_profit), 
avg(ss_net_profit)#111) AS mergedValue#1188]
   
   
      :  :        +- HashAggregate(keys=[], functions=[count(1), 
avg(cast(ss_ext_list_price#33 as double)), avg(cast(ss_net_profit#38 as 
double))])
   
   
      :  :           +- Exchange SinglePartition, ENSURE_REQUIREMENTS, 
[plan_id=2701]
   
   
      :  :              +- HashAggregate(keys=[], functions=[partial_count(1), 
partial_avg(cast(ss_ext_list_price#33 as double)), 
partial_avg(cast(ss_net_profit#38 as double))])
   
   
      :  :                 +- Project [ss_ext_list_price#33, ss_net_profit#38]
   
   
      :  :                    +- Filter ((isnotnull(ss_quantity#26) AND 
(cast(ss_quantity#26 as int) >= 1)) AND (cast(ss_quantity#26 as int) <= 20))
   
   
      :  :                       +- FileScan csv 
[ss_quantity#26,ss_ext_list_price#33,ss_net_profit#38] Batched: false, 
DataFilters: [isnotnull(ss_quantity#26), (cast(ss_quantity#26 as int) >= 1), 
(cast(ss_quantity#26 as int) <= 20)], Format: CSV, Location: 
InMemoryFileIndex(1 paths)[file:/home/coder/store_sales.csv], PartitionFilters: 
[], PushedFilters: [IsNotNull(ss_quantity)], ReadSchema: 
struct<ss_quantity:string,ss_ext_list_price:string,ss_net_profit:string>
   
   
      :  :- Subquery subquery#90, [id=#2899]
   
   
      :  :  +- AdaptiveSparkPlan isFinalPlan=false
   
   
      :  :     +- Project [named_struct(count(1), count(1)#113L, 
avg(ss_ext_list_price), avg(ss_ext_list_price)#115, avg(ss_net_profit), 
avg(ss_net_profit)#117) AS mergedValue#1189]
   
   
      :  :        +- HashAggregate(keys=[], functions=[count(1), 
avg(cast(ss_ext_list_price#198 as double)), avg(cast(ss_net_profit#203 as 
double))])
   
   
      :  :           +- Exchange SinglePartition, ENSURE_REQUIREMENTS, 
[plan_id=2717]
   
   
      :  :              +- HashAggregate(keys=[], functions=[partial_count(1), 
partial_avg(cast(ss_ext_list_price#198 as double)), 
partial_avg(cast(ss_net_profit#203 as double))])
   
   
      :  :                 +- Project [ss_ext_list_price#198, ss_net_profit#203]
   
   
      :  :                    +- Filter ((isnotnull(ss_quantity#191) AND 
(cast(ss_quantity#191 as int) >= 21)) AND (cast(ss_quantity#191 as int) <= 40))
   
   
      :  :                       +- FileScan csv 
[ss_quantity#191,ss_ext_list_price#198,ss_net_profit#203] Batched: false, 
DataFilters: [isnotnull(ss_quantity#191), (cast(ss_quantity#191 as int) >= 21), 
(cast(ss_quantity#191 as int) ..., Format: CSV, Location: InMemoryFileIndex(1 
paths)[file:/home/coder/store_sales.csv], PartitionFilters: [], PushedFilters: 
[IsNotNull(ss_quantity)], ReadSchema: 
struct<ss_quantity:string,ss_ext_list_price:string,ss_net_profit:string>
   
   
      :  :- Subquery subquery#91, [id=#2900]
   
   
      :  :  +- AdaptiveSparkPlan isFinalPlan=false
   
   
      :  :     +- Project [named_struct(count(1), count(1)#113L, 
avg(ss_ext_list_price), avg(ss_ext_list_price)#115, avg(ss_net_profit), 
avg(ss_net_profit)#117) AS mergedValue#1189]
   
   
      :  :        +- HashAggregate(keys=[], functions=[count(1), 
avg(cast(ss_ext_list_price#198 as double)), avg(cast(ss_net_profit#203 as 
double))])
   
   
      :  :           +- Exchange SinglePartition, ENSURE_REQUIREMENTS, 
[plan_id=2733]
   
   
      :  :              +- HashAggregate(keys=[], functions=[partial_count(1), 
partial_avg(cast(ss_ext_list_price#198 as double)), 
partial_avg(cast(ss_net_profit#203 as double))])
   
   
      :  :                 +- Project [ss_ext_list_price#198, ss_net_profit#203]
   
   
      :  :                    +- Filter ((isnotnull(ss_quantity#191) AND 
(cast(ss_quantity#191 as int) >= 21)) AND (cast(ss_quantity#191 as int) <= 40))
   
   
      :  :                       +- FileScan csv 
[ss_quantity#191,ss_ext_list_price#198,ss_net_profit#203] Batched: false, 
DataFilters: [isnotnull(ss_quantity#191), (cast(ss_quantity#191 as int) >= 21), 
(cast(ss_quantity#191 as int) ..., Format: CSV, Location: InMemoryFileIndex(1 
paths)[file:/home/coder/store_sales.csv], PartitionFilters: [], PushedFilters: 
[IsNotNull(ss_quantity)], ReadSchema: 
struct<ss_quantity:string,ss_ext_list_price:string,ss_net_profit:string>
   
   
      :  :- Subquery subquery#92, [id=#2901]
   
   
      :  :  +- AdaptiveSparkPlan isFinalPlan=false
   
   
      :  :     +- Project [named_struct(count(1), count(1)#113L, 
avg(ss_ext_list_price), avg(ss_ext_list_price)#115, avg(ss_net_profit), 
avg(ss_net_profit)#117) AS mergedValue#1189]
   
   
      :  :        +- HashAggregate(keys=[], functions=[count(1), 
avg(cast(ss_ext_list_price#198 as double)), avg(cast(ss_net_profit#203 as 
double))])
   
   
      :  :           +- Exchange SinglePartition, ENSURE_REQUIREMENTS, 
[plan_id=2749]
   
   
      :  :              +- HashAggregate(keys=[], functions=[partial_count(1), 
partial_avg(cast(ss_ext_list_price#198 as double)), 
partial_avg(cast(ss_net_profit#203 as double))])
   
   
      :  :                 +- Project [ss_ext_list_price#198, ss_net_profit#203]
   
   
      :  :                    +- Filter ((isnotnull(ss_quantity#191) AND 
(cast(ss_quantity#191 as int) >= 21)) AND (cast(ss_quantity#191 as int) <= 40))
   
   
      :  :                       +- FileScan csv 
[ss_quantity#191,ss_ext_list_price#198,ss_net_profit#203] Batched: false, 
DataFilters: [isnotnull(ss_quantity#191), (cast(ss_quantity#191 as int) >= 21), 
(cast(ss_quantity#191 as int) ..., Format: CSV, Location: InMemoryFileIndex(1 
paths)[file:/home/coder/store_sales.csv], PartitionFilters: [], PushedFilters: 
[IsNotNull(ss_quantity)], ReadSchema: 
struct<ss_quantity:string,ss_ext_list_price:string,ss_net_profit:string>
   
   
      :  :- Subquery subquery#94, [id=#2902]
   
   
      :  :  +- AdaptiveSparkPlan isFinalPlan=false
   
   
      :  :     +- Project [named_struct(count(1), count(1)#119L, 
avg(ss_ext_list_price), avg(ss_ext_list_price)#121, avg(ss_net_profit), 
avg(ss_net_profit)#123) AS mergedValue#1190]
   
   
      :  :        +- HashAggregate(keys=[], functions=[count(1), 
avg(cast(ss_ext_list_price#267 as double)), avg(cast(ss_net_profit#272 as 
double))])
   
   
      :  :           +- Exchange SinglePartition, ENSURE_REQUIREMENTS, 
[plan_id=2765]
   
   
      :  :              +- HashAggregate(keys=[], functions=[partial_count(1), 
partial_avg(cast(ss_ext_list_price#267 as double)), 
partial_avg(cast(ss_net_profit#272 as double))])
   
   
      :  :                 +- Project [ss_ext_list_price#267, ss_net_profit#272]
   
   
      :  :                    +- Filter ((isnotnull(ss_quantity#260) AND 
(cast(ss_quantity#260 as int) >= 41)) AND (cast(ss_quantity#260 as int) <= 60))
   
   
      :  :                       +- FileScan csv 
[ss_quantity#260,ss_ext_list_price#267,ss_net_profit#272] Batched: false, 
DataFilters: [isnotnull(ss_quantity#260), (cast(ss_quantity#260 as int) >= 41), 
(cast(ss_quantity#260 as int) ..., Format: CSV, Location: InMemoryFileIndex(1 
paths)[file:/home/coder/store_sales.csv], PartitionFilters: [], PushedFilters: 
[IsNotNull(ss_quantity)], ReadSchema: 
struct<ss_quantity:string,ss_ext_list_price:string,ss_net_profit:string>
   
   
      :  :- Subquery subquery#95, [id=#2903]
   
   
      :  :  +- AdaptiveSparkPlan isFinalPlan=false
   
   
      :  :     +- Project [named_struct(count(1), count(1)#119L, 
avg(ss_ext_list_price), avg(ss_ext_list_price)#121, avg(ss_net_profit), 
avg(ss_net_profit)#123) AS mergedValue#1190]
   
   
      :  :        +- HashAggregate(keys=[], functions=[count(1), 
avg(cast(ss_ext_list_price#267 as double)), avg(cast(ss_net_profit#272 as 
double))])
   
   
      :  :           +- Exchange SinglePartition, ENSURE_REQUIREMENTS, 
[plan_id=2781]
   
   
      :  :              +- HashAggregate(keys=[], functions=[partial_count(1), 
partial_avg(cast(ss_ext_list_price#267 as double)), 
partial_avg(cast(ss_net_profit#272 as double))])
   
   
      :  :                 +- Project [ss_ext_list_price#267, ss_net_profit#272]
   
   
      :  :                    +- Filter ((isnotnull(ss_quantity#260) AND 
(cast(ss_quantity#260 as int) >= 41)) AND (cast(ss_quantity#260 as int) <= 60))
   
   
      :  :                       +- FileScan csv 
[ss_quantity#260,ss_ext_list_price#267,ss_net_profit#272] Batched: false, 
DataFilters: [isnotnull(ss_quantity#260), (cast(ss_quantity#260 as int) >= 41), 
(cast(ss_quantity#260 as int) ..., Format: CSV, Location: InMemoryFileIndex(1 
paths)[file:/home/coder/store_sales.csv], PartitionFilters: [], PushedFilters: 
[IsNotNull(ss_quantity)], ReadSchema: 
struct<ss_quantity:string,ss_ext_list_price:string,ss_net_profit:string>
   
   
      :  :- Subquery subquery#96, [id=#2904]
   
   
      :  :  +- AdaptiveSparkPlan isFinalPlan=false
   
   
      :  :     +- Project [named_struct(count(1), count(1)#119L, 
avg(ss_ext_list_price), avg(ss_ext_list_price)#121, avg(ss_net_profit), 
avg(ss_net_profit)#123) AS mergedValue#1190]
   
   
      :  :        +- HashAggregate(keys=[], functions=[count(1), 
avg(cast(ss_ext_list_price#267 as double)), avg(cast(ss_net_profit#272 as 
double))])
   
   
      :  :           +- Exchange SinglePartition, ENSURE_REQUIREMENTS, 
[plan_id=2797]
   
   
      :  :              +- HashAggregate(keys=[], functions=[partial_count(1), 
partial_avg(cast(ss_ext_list_price#267 as double)), 
partial_avg(cast(ss_net_profit#272 as double))])
   
   
      :  :                 +- Project [ss_ext_list_price#267, ss_net_profit#272]
   
   
      :  :                    +- Filter ((isnotnull(ss_quantity#260) AND 
(cast(ss_quantity#260 as int) >= 41)) AND (cast(ss_quantity#260 as int) <= 60))
   
   
      :  :                       +- FileScan csv 
[ss_quantity#260,ss_ext_list_price#267,ss_net_profit#272] Batched: false, 
DataFilters: [isnotnull(ss_quantity#260), (cast(ss_quantity#260 as int) >= 41), 
(cast(ss_quantity#260 as int) ..., Format: CSV, Location: InMemoryFileIndex(1 
paths)[file:/home/coder/store_sales.csv], PartitionFilters: [], PushedFilters: 
[IsNotNull(ss_quantity)], ReadSchema: 
struct<ss_quantity:string,ss_ext_list_price:string,ss_net_profit:string>
   
   
      :  :- Subquery subquery#98, [id=#2905]
   
   
      :  :  +- AdaptiveSparkPlan isFinalPlan=false
   
   
      :  :     +- Project [named_struct(count(1), count(1)#125L, 
avg(ss_ext_list_price), avg(ss_ext_list_price)#127, avg(ss_net_profit), 
avg(ss_net_profit)#129) AS mergedValue#1191]
   
   
      :  :        +- HashAggregate(keys=[], functions=[count(1), 
avg(cast(ss_ext_list_price#336 as double)), avg(cast(ss_net_profit#341 as 
double))])
   
   
      :  :           +- Exchange SinglePartition, ENSURE_REQUIREMENTS, 
[plan_id=2813]
   
   
      :  :              +- HashAggregate(keys=[], functions=[partial_count(1), 
partial_avg(cast(ss_ext_list_price#336 as double)), 
partial_avg(cast(ss_net_profit#341 as double))])
   
   
      :  :                 +- Project [ss_ext_list_price#336, ss_net_profit#341]
   
   
      :  :                    +- Filter ((isnotnull(ss_quantity#329) AND 
(cast(ss_quantity#329 as int) >= 61)) AND (cast(ss_quantity#329 as int) <= 80))
   
   
      :  :                       +- FileScan csv 
[ss_quantity#329,ss_ext_list_price#336,ss_net_profit#341] Batched: false, 
DataFilters: [isnotnull(ss_quantity#329), (cast(ss_quantity#329 as int) >= 61), 
(cast(ss_quantity#329 as int) ..., Format: CSV, Location: InMemoryFileIndex(1 
paths)[file:/home/coder/store_sales.csv], PartitionFilters: [], PushedFilters: 
[IsNotNull(ss_quantity)], ReadSchema: 
struct<ss_quantity:string,ss_ext_list_price:string,ss_net_profit:string>
   
   
      :  :- Subquery subquery#99, [id=#2906]
   
   
      :  :  +- AdaptiveSparkPlan isFinalPlan=false
   
   
      :  :     +- Project [named_struct(count(1), count(1)#125L, 
avg(ss_ext_list_price), avg(ss_ext_list_price)#127, avg(ss_net_profit), 
avg(ss_net_profit)#129) AS mergedValue#1191]
   
   
      :  :        +- HashAggregate(keys=[], functions=[count(1), 
avg(cast(ss_ext_list_price#336 as double)), avg(cast(ss_net_profit#341 as 
double))])
   
   
      :  :           +- Exchange SinglePartition, ENSURE_REQUIREMENTS, 
[plan_id=2829]
   
   
      :  :              +- HashAggregate(keys=[], functions=[partial_count(1), 
partial_avg(cast(ss_ext_list_price#336 as double)), 
partial_avg(cast(ss_net_profit#341 as double))])
   
   
      :  :                 +- Project [ss_ext_list_price#336, ss_net_profit#341]
   
   
      :  :                    +- Filter ((isnotnull(ss_quantity#329) AND 
(cast(ss_quantity#329 as int) >= 61)) AND (cast(ss_quantity#329 as int) <= 80))
   
   
      :  :                       +- FileScan csv 
[ss_quantity#329,ss_ext_list_price#336,ss_net_profit#341] Batched: false, 
DataFilters: [isnotnull(ss_quantity#329), (cast(ss_quantity#329 as int) >= 61), 
(cast(ss_quantity#329 as int) ..., Format: CSV, Location: InMemoryFileIndex(1 
paths)[file:/home/coder/store_sales.csv], PartitionFilters: [], PushedFilters: 
[IsNotNull(ss_quantity)], ReadSchema: 
struct<ss_quantity:string,ss_ext_list_price:string,ss_net_profit:string>
   
   
      :  :- Subquery subquery#100, [id=#2907]
   
   
      :  :  +- AdaptiveSparkPlan isFinalPlan=false
   
   
      :  :     +- Project [named_struct(count(1), count(1)#125L, 
avg(ss_ext_list_price), avg(ss_ext_list_price)#127, avg(ss_net_profit), 
avg(ss_net_profit)#129) AS mergedValue#1191]
   
   
      :  :        +- HashAggregate(keys=[], functions=[count(1), 
avg(cast(ss_ext_list_price#336 as double)), avg(cast(ss_net_profit#341 as 
double))])
   
   
      :  :           +- Exchange SinglePartition, ENSURE_REQUIREMENTS, 
[plan_id=2845]
   
   
      :  :              +- HashAggregate(keys=[], functions=[partial_count(1), 
partial_avg(cast(ss_ext_list_price#336 as double)), 
partial_avg(cast(ss_net_profit#341 as double))])
   
   
      :  :                 +- Project [ss_ext_list_price#336, ss_net_profit#341]
   
   
      :  :                    +- Filter ((isnotnull(ss_quantity#329) AND 
(cast(ss_quantity#329 as int) >= 61)) AND (cast(ss_quantity#329 as int) <= 80))
   
   
      :  :                       +- FileScan csv 
[ss_quantity#329,ss_ext_list_price#336,ss_net_profit#341] Batched: false, 
DataFilters: [isnotnull(ss_quantity#329), (cast(ss_quantity#329 as int) >= 61), 
(cast(ss_quantity#329 as int) ..., Format: CSV, Location: InMemoryFileIndex(1 
paths)[file:/home/coder/store_sales.csv], PartitionFilters: [], PushedFilters: 
[IsNotNull(ss_quantity)], ReadSchema: 
struct<ss_quantity:string,ss_ext_list_price:string,ss_net_profit:string>
   
   
      :  :- Subquery subquery#102, [id=#2908]
   
   
      :  :  +- AdaptiveSparkPlan isFinalPlan=false
   
   
      :  :     +- Project [named_struct(count(1), count(1)#131L, 
avg(ss_ext_list_price), avg(ss_ext_list_price)#133, avg(ss_net_profit), 
avg(ss_net_profit)#135) AS mergedValue#1192]
   
   
      :  :        +- HashAggregate(keys=[], functions=[count(1), 
avg(cast(ss_ext_list_price#405 as double)), avg(cast(ss_net_profit#410 as 
double))])
   
   
      :  :           +- Exchange SinglePartition, ENSURE_REQUIREMENTS, 
[plan_id=2861]
   
   
      :  :              +- HashAggregate(keys=[], functions=[partial_count(1), 
partial_avg(cast(ss_ext_list_price#405 as double)), 
partial_avg(cast(ss_net_profit#410 as double))])
   
   
      :  :                 +- Project [ss_ext_list_price#405, ss_net_profit#410]
   
   
      :  :                    +- Filter ((isnotnull(ss_quantity#398) AND 
(cast(ss_quantity#398 as int) >= 81)) AND (cast(ss_quantity#398 as int) <= 100))
   
   
      :  :                       +- FileScan csv 
[ss_quantity#398,ss_ext_list_price#405,ss_net_profit#410] Batched: false, 
DataFilters: [isnotnull(ss_quantity#398), (cast(ss_quantity#398 as int) >= 81), 
(cast(ss_quantity#398 as int) ..., Format: CSV, Location: InMemoryFileIndex(1 
paths)[file:/home/coder/store_sales.csv], PartitionFilters: [], PushedFilters: 
[IsNotNull(ss_quantity)], ReadSchema: 
struct<ss_quantity:string,ss_ext_list_price:string,ss_net_profit:string>
   
   
      :  :- Subquery subquery#103, [id=#2909]
   
   
      :  :  +- AdaptiveSparkPlan isFinalPlan=false
   
   
      :  :     +- Project [named_struct(count(1), count(1)#131L, 
avg(ss_ext_list_price), avg(ss_ext_list_price)#133, avg(ss_net_profit), 
avg(ss_net_profit)#135) AS mergedValue#1192]
   
   
      :  :        +- HashAggregate(keys=[], functions=[count(1), 
avg(cast(ss_ext_list_price#405 as double)), avg(cast(ss_net_profit#410 as 
double))])
   
   
      :  :           +- Exchange SinglePartition, ENSURE_REQUIREMENTS, 
[plan_id=2877]
   
   
      :  :              +- HashAggregate(keys=[], functions=[partial_count(1), 
partial_avg(cast(ss_ext_list_price#405 as double)), 
partial_avg(cast(ss_net_profit#410 as double))])
   
   
      :  :                 +- Project [ss_ext_list_price#405, ss_net_profit#410]
   
   
      :  :                    +- Filter ((isnotnull(ss_quantity#398) AND 
(cast(ss_quantity#398 as int) >= 81)) AND (cast(ss_quantity#398 as int) <= 100))
   
   
      :  :                       +- FileScan csv 
[ss_quantity#398,ss_ext_list_price#405,ss_net_profit#410] Batched: false, 
DataFilters: [isnotnull(ss_quantity#398), (cast(ss_quantity#398 as int) >= 81), 
(cast(ss_quantity#398 as int) ..., Format: CSV, Location: InMemoryFileIndex(1 
paths)[file:/home/coder/store_sales.csv], PartitionFilters: [], PushedFilters: 
[IsNotNull(ss_quantity)], ReadSchema: 
struct<ss_quantity:string,ss_ext_list_price:string,ss_net_profit:string>
   
   
      :  +- Subquery subquery#104, [id=#2910]
   
   
      :     +- AdaptiveSparkPlan isFinalPlan=false
   
   
      :        +- Project [named_struct(count(1), count(1)#131L, 
avg(ss_ext_list_price), avg(ss_ext_list_price)#133, avg(ss_net_profit), 
avg(ss_net_profit)#135) AS mergedValue#1192]
   
   
      :           +- HashAggregate(keys=[], functions=[count(1), 
avg(cast(ss_ext_list_price#405 as double)), avg(cast(ss_net_profit#410 as 
double))])
   
   
      :              +- Exchange SinglePartition, ENSURE_REQUIREMENTS, 
[plan_id=2893]
   
   
      :                 +- HashAggregate(keys=[], functions=[partial_count(1), 
partial_avg(cast(ss_ext_list_price#405 as double)), 
partial_avg(cast(ss_net_profit#410 as double))])
   
   
      :                    +- Project [ss_ext_list_price#405, ss_net_profit#410]
   
   
      :                       +- Filter ((isnotnull(ss_quantity#398) AND 
(cast(ss_quantity#398 as int) >= 81)) AND (cast(ss_quantity#398 as int) <= 100))
   
   
      :                          +- FileScan csv 
[ss_quantity#398,ss_ext_list_price#405,ss_net_profit#410] Batched: false, 
DataFilters: [isnotnull(ss_quantity#398), (cast(ss_quantity#398 as int) >= 81), 
(cast(ss_quantity#398 as int) ..., Format: CSV, Location: InMemoryFileIndex(1 
paths)[file:/home/coder/store_sales.csv], PartitionFilters: [], PushedFilters: 
[IsNotNull(ss_quantity)], ReadSchema: 
struct<ss_quantity:string,ss_ext_list_price:string,ss_net_profit:string>
   
   
      +- Filter (isnotnull(r_reason_sk#80) AND (cast(r_reason_sk#80 as int) = 
1))
   
   
         +- FileScan csv [r_reason_sk#80] Batched: false, DataFilters: 
[isnotnull(r_reason_sk#80), (cast(r_reason_sk#80 as int) = 1)], Format: CSV, 
Location: InMemoryFileIndex(1 paths)[file:/home/coder/reason.csv], 
PartitionFilters: [], PushedFilters: [IsNotNull(r_reason_sk)], ReadSchema: 
struct<r_reason_sk:string>
   ```


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to