[ 
https://issues.apache.org/jira/browse/HUDI-9629?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Davis Zhang resolved HUDI-9629.
-------------------------------

> NPE after DS write with schema evolution
> ----------------------------------------
>
>                 Key: HUDI-9629
>                 URL: https://issues.apache.org/jira/browse/HUDI-9629
>             Project: Apache Hudi
>          Issue Type: Bug
>            Reporter: Davis Zhang
>            Priority: Blocker
>             Fix For: 1.1.0
>
>
> {code:java}
>   test("Test schema evolution basic case") {
>     withTempDir { tmp =>
>       val tableName = generateTableName
>       val basePath = s"${tmp.getCanonicalPath}/$tableName"
>       // Create table with initial schema where quantity is int
>       spark.sql(
>         s"""
>             |create table $tableName (
>             |  id int,
>             |  name string,
>             |  quantity int,
>             |  ts long
>             |) using hudi
>             | options (
>             |  primaryKey ='id',
>             |  type = 'cow',
>             |  preCombineField = 'ts',
>             |  hoodie.metadata.enable = 'true',
>             |  hoodie.metadata.record.index.enable = 'true'
>             | )
>             | location '$basePath'
>     """.stripMargin)
>       // Insert initial data with integer quantities
>       spark.sql(s"insert into $tableName values(1, 'a1', 10, 1000)")
>       spark.sql(s"insert into $tableName values(2, 'a2', 20, 1001)")
>       // Now the schema evolution should work using data source write
>       import spark.implicits._
>       val evolvedData = Seq(
>         (3, "a3", 30.5, 1002L) // quantity is double
>       ).toDF("id", "name", "quantity", "ts")
>       evolvedData.write
>         .format("hudi")
>         .option("hoodie.table.name", tableName)
>         .option("hoodie.datasource.write.table.type", "COPY_ON_WRITE")
>         .option("hoodie.datasource.write.recordkey.field", "id")
>         .option("hoodie.datasource.write.precombine.field", "ts")
>         .option("hoodie.datasource.write.operation", "upsert")
>         .option("hoodie.schema.on.read.enable", "true")
>         .mode("append")
>         .save(basePath)
>      // this hits NPE on executor
>       spark.sql(s"SELECT id, name, quantity, ts FROM $tableName ORDER BY 
> id").collect()
>     }
>   }got NPE on spark executor code path which is very surprising
> Caused by: java.lang.NullPointerException
>       at 
> org.apache.spark.sql.execution.vectorized.OnHeapColumnVector.getInt(OnHeapColumnVector.java:328)
>       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 scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
>       at scala.collection.Iterator$$anon$10.hasNext(Iterator.scala:460)
>       at 
> org.apache.spark.util.random.SamplingUtils$.reservoirSampleAndCount(SamplingUtils.scala:41)
>       at 
> org.apache.spark.RangePartitioner$.$anonfun$sketch$1(Partitioner.scala:322)
>       at 
> org.apache.spark.RangePartitioner$.$anonfun$sketch$1$adapted(Partitioner.scala:320)
>       at 
> org.apache.spark.rdd.RDD.$anonfun$mapPartitionsWithIndex$2(RDD.scala:910)
>       at 
> org.apache.spark.rdd.RDD.$anonfun$mapPartitionsWithIndex$2$adapted(RDD.scala:910)
>       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.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
>       at 
> java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
>       at java.lang.Thread.run(Thread.java:750) {code}



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to