[
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)