apache / apache/hudi

[SUPPORT] org.apache.avro.SchemaParseException: Can't redefine: array When there are Top level variables , Struct and Array[struct] (no complex datatype within array[struct])

Open
#7,717 11 comments 0 reactions 0 assignees View on GitHub
area:schema area:sql priority:medium
Dominant language
Java
Stars
6.2k
Forks
2.5k
Avg merge
2d 8h
Merged PRs (30d)
111

Description

### Describe the problem you faced

When storing a data structure with the following layout into a copy-on-write table:
```
root
|-- personDetails: struct (nullable = true)
| |-- id: integer (nullable = false)
|-- idInfo: struct (nullable = true)
| |-- adhaarId: integer (nullable = false)
|-- addressInfo: array (nullable = true)
| |-- element: struct (containsNull = true)
| | |-- addressId: integer (nullable = false)
|-- employmentInfo: array (nullable = true)
| |-- element: struct (containsNull = true)
| | |-- employmenCd: integer (nullable = false)
|-- src_load_ts: timestamp (nullable = false)
|-- load_ts: timestamp (nullable = false)
|-- load_dt: date (nullable = false)
```
the first write will succeed, but then subsequent writes will fail with the error included in the stacktrace.

### To Reproduce

Steps to reproduce the behavior:

```
case class Person(personDetails: PersonDetails,
idInfo: IdInfo,
addressInfo: Array[AddressInfo] = Array.empty[AddressInfo],
employmentInfo: Array[EmploymentInfo] = Array.empty[EmploymentInfo])

case class PersonDetails(id: Int)

case class IdInfo(adhaarId: Int)

case class AddressInfo(addressId: Int)

case class EmploymentInfo(employmenCd: Int)

def maskedParquetBugTest(spark: SparkSession): Unit = {

import spark.implicits._

val personDetails1 = PersonDetails(1)
val idInfo1 = IdInfo(1)
val addressInfo1 = AddressInfo(1)
val employmentInfo1 = EmploymentInfo(1)

val item1 = Person(personDetails1, idInfo1, Array(addressInfo1), Array(employmentInfo1))
val parquetBugDs = Seq(item1).toDF()
.withColumn("src_load_ts", current_timestamp())
.withColumn("load_ts", timestampInCst).withColumn("load_dt", to_date(col("load_ts")))

parquetBugDs.printSchema()

writeHudi(parquetBugDs, "parquet_bug_ds",
"load_dt",
"personDetails.id",
"src_load_ts")
}

def writeHudi(ds: DataFrame, tableName: String, partitionPath: String, recordKey: String, precombineKey: String): Unit = {

val hoodieConfigs: util.Map[String, String] = new java.util.HashMap[String, String]
hoodieConfigs.put("hoodie.table.name", tableName)
hoodieConfigs.put("hoodie.datasource.write.keygenerator.class", classOf[SimpleKeyGenerator].getName)
hoodieConfigs.put("hoodie.datasource.write.partitionpath.field", partitionPath)
hoodieConfigs.put("hoodie.datasource.write.recordkey.field", recordKey)
hoodieConfigs.put("hoodie.datasource.write.precombine.field", precombineKey)
hoodieConfigs.put("hoodie.payload.ordering.field", precombineKey)
hoodieConfigs.put("hoodie.index.type", "GLOBAL_SIMPLE")
hoodieConfigs.put("hoodie.insert.shuffle.parallelism", "1")
hoodieConfigs.put("hoodie.upsert.shuffle.parallelism", "1")
hoodieConfigs.put("hoodie.bulkinsert.shuffle.parallelism", "1")
hoodieConfigs.put("hoodie.delete.shuffle.parallelism", "1")
hoodieConfigs.put("hoodie.simple.index.update.partition.path", "false")
hoodieConfigs.put("hoodie.datasource.write.payload.class", classOf[DefaultHoodieRecordPayload].getName)
hoodieConfigs.put("hoodie.datasource.write.hive_style_partitioning", "false")
hoodieConfigs.put("hoodie.datasource.write.table.type", COW_TABLE_TYPE_OPT_VAL)
hoodieConfigs.put("hoodie.datasource.write.row.writer.enable", "true")
hoodieConfigs.put("hoodie.combine.before.upsert", "true")
hoodieConfigs.put("hoodie.datasource.write.keygenerator.consistent.logical.timestamp.enabled", "true")
hoodieConfigs.put("hoodie.schema.on.read.enable", "true")
hoodieConfigs.put("hoodie.datasource.write.reconcile.schema", "true")
hoodieConfigs.put("hoodie.datasource.write.operation", "upsert")

ds.toDF().write.format("hudi").
options(hoodieConfigs).
mode("append").
save(s"/tmp/data/hudi/$tableName")
}

maskedParquetBugTest(spark)

maskedParquetBugTest(spark)
```

### Expected behavior

The second write succeeds.

### Environment Description

Hudi version (hudi-spark3.1-bundle_2.12) : 0.12.2 , 0.12.1, 0.12.0

Spark version : 3.1.3

Hive version : -

Hadoop version : -

Storage (HDFS/S3/GCS..) : Local storage

Running on Docker? (yes/no) : No

### Stack Trace
```
Driver stacktrace:
at org.apache.spark.scheduler.DAGScheduler.failJobAndIndependentStages(DAGScheduler.scala:2303)
at org.apache.spark.scheduler.DAGScheduler.$anonfun$abortStage$2(DAGScheduler.scala:2252)
at org.apache.spark.scheduler.DAGScheduler.$anonfun$abortStage$2$adapted(DAGScheduler.scala:2251)
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:2251)
at org.apache.spark.scheduler.DAGScheduler.$anonfun$handleTaskSetFailed$1(DAGScheduler.scala:1124)
at org.apache.spark.scheduler.DAGScheduler.$anonfun$handleTaskSetFailed$1$adapted(DAGScheduler.scala:1124)
at scala.Option.foreach(Option.scala:407)
at org.apache.spark.scheduler.DAGScheduler.handleTaskSetFailed(DAGScheduler.scala:1124)
at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:2490)
at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2432)
at org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2421)
at org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:49)
at org.apache.spark.scheduler.DAGScheduler.runJob(DAGScheduler.scala:902)
at org.apache.spark.SparkContext.runJob(SparkContext.scala:2196)
at org.apache.spark.SparkContext.runJob(SparkContext.scala:2217)
at org.apache.spark.SparkContext.runJob(SparkContext.scala:2236)
at org.apache.spark.SparkContext.runJob(SparkContext.scala:2261)
at org.apache.spark.rdd.RDD.count(RDD.scala:1253)
at org.apache.hudi.HoodieSparkSqlWriter$.commitAndPerformPostOperations(HoodieSparkSqlWriter.scala:693)
at org.apache.hudi.HoodieSparkSqlWriter$.write(HoodieSparkSqlWriter.scala:345)
at org.apache.hudi.DefaultSource.createRelation(DefaultSource.scala:145)
at org.apache.spark.sql.execution.datasources.SaveIntoDataSourceCommand.run(SaveIntoDataSourceCommand.scala:46)
at org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult$lzycompute(commands.scala:70)
at org.apache.spark.sql.execution.command.ExecutedCommandExec.sideEffectResult(commands.scala:68)
at org.apache.spark.sql.execution.command.ExecutedCommandExec.doExecute(commands.scala:90)
at org.apache.spark.sql.execution.SparkPlan.$anonfun$execute$1(SparkPlan.scala:180)
at org.apache.spark.sql.execution.SparkPlan.$anonfun$executeQuery$1(SparkPlan.scala:218)
at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
at org.apache.spark.sql.execution.SparkPlan.executeQuery(SparkPlan.scala:215)
at org.apache.spark.sql.execution.SparkPlan.execute(SparkPlan.scala:176)
at org.apache.spark.sql.execution.QueryExecution.toRdd$lzycompute(QueryExecution.scala:132)
at org.apache.spark.sql.execution.QueryExecution.toRdd(QueryExecution.scala:131)
at org.apache.spark.sql.DataFrameWriter.$anonfun$runCommand$1(DataFrameWriter.scala:989)
at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$5(SQLExecution.scala:103)
at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:163)
at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$1(SQLExecution.scala:90)
at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:775)
at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:64)
at org.apache.spark.sql.DataFrameWriter.runCommand(DataFrameWriter.scala:989)
at org.apache.spark.sql.DataFrameWriter.saveToV1Source(DataFrameWriter.scala:438)
at org.apache.spark.sql.DataFrameWriter.saveInternal(DataFrameWriter.scala:415)
at org.apache.spark.sql.DataFrameWriter.save(DataFrameWriter.scala:293)
at com.test.run.hudi.upsert.complex.utils.HudiParquetBugTest$.writeHudi(HudiParquetBugTest.scala:99)
at com.test.run.hudi.upsert.complex.utils.HudiParquetBugTest$.maskedParquetBugTest(HudiParquetBugTest.scala:68)
at com.test.run.hudi.upsert.complex.utils.HudiParquetBugTest$.delayedEndpoint$com$walmart$hnw$datafoundations$techmod$ingestion$utils$HudiParquetBugTest$1(HudiParquetBugTest.scala:32)
at com.test.run.hudi.upsert.complex.utils.HudiParquetBugTest$delayedInit$body.apply(HudiParquetBugTest.scala:14)
at scala.Function0.apply$mcV$sp(Function0.scala:39)
at scala.Function0.apply$mcV$sp$(Function0.scala:39)
at scala.runtime.AbstractFunction0.apply$mcV$sp(AbstractFunction0.scala:17)
at scala.App.$anonfun$main$1$adapted(App.scala:80)
at scala.collection.immutable.List.foreach(List.scala:431)
at scala.App.main(App.scala:80)
at scala.App.main$(App.scala:78)
at com.test.run.hudi.upsert.complex.utils.HudiParquetBugTest$.main(HudiParquetBugTest.scala:14)
at com.test.run.hudi.upsert.complex.utils.HudiParquetBugTest.main(HudiParquetBugTest.scala)
Caused by: org.apache.hudi.exception.HoodieUpsertException: Error upserting bucketType UPDATE for partition :0
at org.apache.hudi.table.action.commit.BaseSparkCommitActionExecutor.handleUpsertPartition(BaseSparkCommitActionExecutor.java:329)
at org.apache.hudi.table.action.commit.BaseSparkCommitActionExecutor.lambda$mapPartitionsAsRDD$a3ab3c4$1(BaseSparkCommitActionExecutor.java:244)
at org.apache.spark.api.java.JavaRDDLike.$anonfun$mapPartitionsWithIndex$1(JavaRDDLike.scala:102)
at org.apache.spark.api.java.JavaRDDLike.$anonfun$mapPartitionsWithIndex$1$adapted(JavaRDDLike.scala:102)
at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsWithIndex$2(RDD.scala:915)
at org.apache.spark.rdd.RDD.$anonfun$mapPartitionsWithIndex$2$adapted(RDD.scala:915)
at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:373)
at org.apache.spark.rdd.RDD.iterator(RDD.scala:337)
at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:373)
at org.apache.spark.rdd.RDD.$anonfun$getOrCompute$1(RDD.scala:386)
at org.apache.spark.storage.BlockManager.$anonfun$doPutIterator$1(BlockManager.scala:1440)
at org.apache.spark.storage.BlockManager.org$apache$spark$storage$BlockManager$$doPut(BlockManager.scala:1350)
at org.apache.spark.storage.BlockManager.doPutIterator(BlockManager.scala:1414)
at org.apache.spark.storage.BlockManager.getOrElseUpdate(BlockManager.scala:1237)
at org.apache.spark.rdd.RDD.getOrCompute(RDD.scala:384)
at org.apache.spark.rdd.RDD.iterator(RDD.scala:335)
at org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:52)
at org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:373)
at org.apache.spark.rdd.RDD.iterator(RDD.scala:337)
at org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:90)
at org.apache.spark.scheduler.Task.run(Task.scala:131)
at org.apache.spark.executor.Executor$TaskRunner.$anonfun$run$3(Executor.scala:498)
at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1439)
at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:501)
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:748)
Caused by: org.apache.hudi.exception.HoodieException: org.apache.avro.SchemaParseException: Can't redefine: array
at org.apache.hudi.table.action.commit.HoodieMergeHelper.runMerge(HoodieMergeHelper.java:166)
at org.apache.hudi.table.action.commit.BaseSparkCommitActionExecutor.handleUpdateInternal(BaseSparkCommitActionExecutor.java:358)
at org.apache.hudi.table.action.commit.BaseSparkCommitActionExecutor.handleUpdate(BaseSparkCommitActionExecutor.java:349)
at org.apache.hudi.table.action.commit.BaseSparkCommitActionExecutor.handleUpsertPartition(BaseSparkCommitActionExecutor.java:322)
... 28 more
Caused by: org.apache.avro.SchemaParseException: Can't redefine: array
at org.apache.avro.Schema$Names.put(Schema.java:1550)
at org.apache.avro.Schema$NamedSchema.writeNameRef(Schema.java:813)
at org.apache.avro.Schema$RecordSchema.toJson(Schema.java:975)
at org.apache.avro.Schema$ArraySchema.toJson(Schema.java:1137)
at org.apache.avro.Schema$UnionSchema.toJson(Schema.java:1242)
at org.apache.avro.Schema$RecordSchema.fieldsToJson(Schema.java:1003)
at org.apache.avro.Schema$RecordSchema.toJson(Schema.java:987)
at org.apache.avro.Schema.toString(Schema.java:426)
at org.apache.avro.Schema.toString(Schema.java:398)
at org.apache.avro.Schema.toString(Schema.java:389)
at org.apache.parquet.avro.AvroReadSupport.setAvroReadSchema(AvroReadSupport.java:69)
at org.apache.hudi.io.storage.HoodieParquetReader.getRecordIterator(HoodieParquetReader.java:69)
at org.apache.hudi.io.storage.HoodieFileReader.getRecordIterator(HoodieFileReader.java:43)
at org.apache.hudi.table.action.commit.HoodieMergeHelper.runMerge(HoodieMergeHelper.java:149)
```

Contributor guide

No contributing guide indexed for this repository

Research direction

Reproduce the failure with the maskedParquetBugTest and writeHudi example from the issue, then trace the stack from HoodieSparkSqlWriter and BaseSparkCommitActionExecutor into HoodieMergeHelper. Compare the first and second writes involving the top-level structs and arrays. Done means the second append/upsert succeeds without SchemaParseException: Can't redefine: array.

Written by the indexing model from the issue text.

Assessment

Tech stack
java, scala
Domain
data-engineering, databases
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
32/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.