baisui1981 opened a new issue #4802:
URL: https://github.com/apache/hudi/issues/4802


   **Describe the problem you faced**
   
   When export MySql rows into Hudi table , If none null value in all columns , 
the inserting process will be success ,that can execute Hive query for imported 
table from MySql.   Another case , having null value in MySql table ,the 
inserting process will be faild.
   
   **To Reproduce**
   For example there a simple MySql table, Table `base`'s format as below:
   
   | Field        | Type         | Null | Key | Default           | Extra |
   |--------------|--------------|------|-----|-------------------|-------|
   | base_id      | int(11)      | NO   | PRI | NULL              |       |
   | start_time   | datetime     | YES  |     | NULL              |       |
   | update_date  | date         | YES  |     | NULL              |       |
   | update_time  | timestamp    | NO   |     | CURRENT_TIMESTAMP |       |
   | price        | decimal(5,2) | YES  |     | NULL              |       |
   | json_content | json         | YES  |     | NULL              |       |
   | col_blob     | blob         | YES  |     | NULL              |       |
   | col_text     | text         | YES  |     | NULL              |       |
   
   The process contains two steps:
   1. Export MySql table to Hdfs with `CSV` format , since just has 4 rows in 
the Table, csv file content show as below:
   ```
   
base_id,start_time,update_date,update_time,price,json_content,col_blob,col_text
   1,"2022-01-14 12:07:28",2022-01-14,"2022-01-14 12:07:28",2.99,{},"wo ai 
beijing","wo ai zuguo"
   2,"2021-12-27 15:51:31",2021-12-27,"2021-12-27 15:51:31",99.2,"{""name"": 
""baisui""}",,col_text
   3,"2022-01-14 12:10:12",2022-01-14,"2022-01-14 12:10:12",2.99,{},"wo ai 
beijing","wo ai zuguo"
   4,"2022-01-14 12:23:03",2022-01-14,"2022-01-14 12:23:03",2.99,{},"wo ai 
beijing","wo ai zuguo"
   ```
    **notice that** ,the row which base_id equals to 2 , col_blob is **null**
   
   2. using 
[hoodie_deltastreamer](https://hudi.apache.org/docs/hoodie_deltastreamer) , 
ingest to base table in HDFS
   
   ``` shell
   spark-submit \
     --class org.apache.hudi.utilities.deltastreamer.HoodieDeltaStreamer 
$HUDI_UTILITIES_BUNDLE \
     --table-type COPY_ON_WRITE \
     --source-class org.apache.hudi.utilities.sources.CsvDFSSource \
     --source-ordering-field update_time  \
     --target-base-path /user/hive/warehouse/base \
     --target-table base --props /var/demo/config/base-source.properties \
     --schemaprovider-class 
org.apache.hudi.utilities.schema.FilebasedSchemaProvider \
     --enable-sync
   ```
   
   content of path: /var/demo/config/base-source.properties
   ```
   hoodie.upsert.shuffle.parallelism=2
   hoodie.insert.shuffle.parallelism=2
   hoodie.delete.shuffle.parallelism=2
   hoodie.bulkinsert.shuffle.parallelism=2
   hoodie.embed.timeline.server=true
   hoodie.filesystem.view.type=EMBEDDED_KV_STORE
   hoodie.compact.inline=false
   
hoodie.deltastreamer.source.dfs.root=hdfs://namenode/user/admin/default/base/data
   hoodie.deltastreamer.csv.header=true
   hoodie.deltastreamer.csv.sep=,
   
hoodie.deltastreamer.schemaprovider.source.schema.file=hdfs://namenode/user/admin/default/base/meta/schema.avsc
   
hoodie.deltastreamer.schemaprovider.target.schema.file=hdfs://namenode/user/admin/default/base/meta/schema.avsc
   hoodie.datasource.hive_sync.database=default
   hoodie.datasource.hive_sync.table=base
   hoodie.datasource.hive_sync.partition_fields=pt
   
hoodie.datasource.hive_sync.partition_extractor_class=org.apache.hudi.hive.MultiPartKeysValueExtractor
   hoodie.datasource.hive_sync.username=root
   hoodie.datasource.hive_sync.password=111111
   hoodie.datasource.hive_sync.jdbcurl=jdbc:hive2://192.168.28.201:10000
   hoodie.datasource.hive_sync.mode=jdbc
   hoodie.datasource.write.recordkey.field=base_id
   hoodie.datasource.write.partitionpath.field=base_id
   ```
   
   content of: `/user/admin/default/base/meta/schema.avsc`
   
   ```json
   {
     "type" : "record",
     "name" : "base",
     "fields" : [ {
       "name" : "base_id",
       "type" : "int"
     }, {
       "name" : "start_time",
       "type" : "string",
       "default" : ""
     }, {
       "name" : "update_date",
       "type" : "string",
       "default" : ""
     }, {
       "name" : "update_time",
       "type" : "string"
     }, {
       "name" : "price",
       "type" : "double"
     }, {
       "name" : "json_content",
       "type" : "string",
       "default" : ""
     }, {
       "name" : "col_blob",
       "type" : "string",
       "default" : ""
     }, {
       "name" : "col_text",
       "type" : "string",
       "default" : ""
     } ]
   }
   ```
   
   after trigger the `hoodie_deltastreamer`,get the error log :
   ``` shell
   Driver stacktrace:
        at 
org.apache.spark.scheduler.DAGScheduler.org$apache$spark$scheduler$DAGScheduler$$failJobAndIndependentStages(DAGScheduler.scala:1889)
        at 
org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1877)
        at 
org.apache.spark.scheduler.DAGScheduler$$anonfun$abortStage$1.apply(DAGScheduler.scala:1876)
        at 
scala.collection.mutable.ResizableArray$class.foreach(ResizableArray.scala:59)
        at scala.collection.mutable.ArrayBuffer.foreach(ArrayBuffer.scala:48)
        at 
org.apache.spark.scheduler.DAGScheduler.abortStage(DAGScheduler.scala:1876)
        at 
org.apache.spark.scheduler.DAGScheduler$$anonfun$handleTaskSetFailed$1.apply(DAGScheduler.scala:926)
        at 
org.apache.spark.scheduler.DAGScheduler$$anonfun$handleTaskSetFailed$1.apply(DAGScheduler.scala:926)
        at scala.Option.foreach(Option.scala:257)
        at 
org.apache.spark.scheduler.DAGScheduler.handleTaskSetFailed(DAGScheduler.scala:926)
        at 
org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.doOnReceive(DAGScheduler.scala:2110)
        at 
org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2059)
        at 
org.apache.spark.scheduler.DAGSchedulerEventProcessLoop.onReceive(DAGScheduler.scala:2048)
        at org.apache.spark.util.EventLoop$$anon$1.run(EventLoop.scala:49)
        at 
org.apache.spark.scheduler.DAGScheduler.runJob(DAGScheduler.scala:737)
        at org.apache.spark.SparkContext.runJob(SparkContext.scala:2061)
        at org.apache.spark.SparkContext.runJob(SparkContext.scala:2082)
        at org.apache.spark.SparkContext.runJob(SparkContext.scala:2101)
        at org.apache.spark.SparkContext.runJob(SparkContext.scala:2126)
        at org.apache.spark.rdd.RDD$$anonfun$collect$1.apply(RDD.scala:945)
        at 
org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
        at 
org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
        at org.apache.spark.rdd.RDD.withScope(RDD.scala:363)
        at org.apache.spark.rdd.RDD.collect(RDD.scala:944)
        at 
org.apache.spark.rdd.PairRDDFunctions$$anonfun$countByKey$1.apply(PairRDDFunctions.scala:370)
        at 
org.apache.spark.rdd.PairRDDFunctions$$anonfun$countByKey$1.apply(PairRDDFunctions.scala:370)
        at 
org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
        at 
org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
        at org.apache.spark.rdd.RDD.withScope(RDD.scala:363)
        at 
org.apache.spark.rdd.PairRDDFunctions.countByKey(PairRDDFunctions.scala:369)
        at 
org.apache.spark.api.java.JavaPairRDD.countByKey(JavaPairRDD.scala:312)
        at 
org.apache.hudi.data.HoodieJavaPairRDD.countByKey(HoodieJavaPairRDD.java:103)
        at 
org.apache.hudi.index.bloom.HoodieBloomIndex.lookupIndex(HoodieBloomIndex.java:115)
        at 
org.apache.hudi.index.bloom.HoodieBloomIndex.tagLocation(HoodieBloomIndex.java:85)
        at 
org.apache.hudi.table.action.commit.SparkWriteHelper.tag(SparkWriteHelper.java:56)
        at 
org.apache.hudi.table.action.commit.SparkWriteHelper.tag(SparkWriteHelper.java:39)
        at 
org.apache.hudi.table.action.commit.AbstractWriteHelper.write(AbstractWriteHelper.java:51)
        ... 22 more
   Caused by: java.io.IOException: Could not create payload for class: 
org.apache.hudi.common.model.OverwriteWithLatestAvroPayload
        at 
org.apache.hudi.DataSourceUtils.createPayload(DataSourceUtils.java:133)
        at 
org.apache.hudi.utilities.deltastreamer.DeltaSync.lambda$readFromSource$d62e16$1(DeltaSync.java:450)
        at 
org.apache.spark.api.java.JavaPairRDD$$anonfun$toScalaFunction$1.apply(JavaPairRDD.scala:1040)
        at scala.collection.Iterator$$anon$11.next(Iterator.scala:410)
        at scala.collection.Iterator$$anon$11.next(Iterator.scala:410)
        at 
org.apache.spark.util.collection.ExternalSorter.insertAll(ExternalSorter.scala:193)
        at 
org.apache.spark.shuffle.sort.SortShuffleWriter.write(SortShuffleWriter.scala:62)
        at 
org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:99)
        at 
org.apache.spark.scheduler.ShuffleMapTask.runTask(ShuffleMapTask.scala:55)
        at org.apache.spark.scheduler.Task.run(Task.scala:123)
        at 
org.apache.spark.executor.Executor$TaskRunner$$anonfun$10.apply(Executor.scala:408)
        at org.apache.spark.util.Utils$.tryWithSafeFinally(Utils.scala:1360)
        at org.apache.spark.executor.Executor$TaskRunner.run(Executor.scala:414)
        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: Unable to instantiate 
class org.apache.hudi.common.model.OverwriteWithLatestAvroPayload
        at 
org.apache.hudi.common.util.ReflectionUtils.loadClass(ReflectionUtils.java:91)
        at 
org.apache.hudi.DataSourceUtils.createPayload(DataSourceUtils.java:130)
        ... 15 more
   Caused by: java.lang.reflect.InvocationTargetException
        at sun.reflect.NativeConstructorAccessorImpl.newInstance0(Native Method)
        at 
sun.reflect.NativeConstructorAccessorImpl.newInstance(NativeConstructorAccessorImpl.java:62)
        at 
sun.reflect.DelegatingConstructorAccessorImpl.newInstance(DelegatingConstructorAccessorImpl.java:45)
        at java.lang.reflect.Constructor.newInstance(Constructor.java:423)
        at 
org.apache.hudi.common.util.ReflectionUtils.loadClass(ReflectionUtils.java:89)
        ... 16 more
   Caused by: java.lang.NullPointerException: null of string in field col_blob 
of hoodie.source.hoodie_source
        at 
org.apache.avro.generic.GenericDatumWriter.npe(GenericDatumWriter.java:145)
        at 
org.apache.avro.generic.GenericDatumWriter.writeWithoutConversion(GenericDatumWriter.java:139)
        at 
org.apache.avro.generic.GenericDatumWriter.write(GenericDatumWriter.java:75)
        at 
org.apache.avro.generic.GenericDatumWriter.write(GenericDatumWriter.java:62)
        at 
org.apache.hudi.avro.HoodieAvroUtils.indexedRecordToBytes(HoodieAvroUtils.java:104)
        at 
org.apache.hudi.avro.HoodieAvroUtils.avroToBytes(HoodieAvroUtils.java:96)
        at 
org.apache.hudi.common.model.BaseAvroPayload.<init>(BaseAvroPayload.java:49)
        at 
org.apache.hudi.common.model.OverwriteWithLatestAvroPayload.<init>(OverwriteWithLatestAvroPayload.java:44)
        ... 21 more
   Caused by: java.lang.NullPointerException
        at org.apache.avro.io.Encoder.writeString(Encoder.java:121)
        at 
org.apache.avro.generic.GenericDatumWriter.writeString(GenericDatumWriter.java:267)
        at 
org.apache.avro.generic.GenericDatumWriter.writeString(GenericDatumWriter.java:262)
        at 
org.apache.avro.generic.GenericDatumWriter.writeWithoutConversion(GenericDatumWriter.java:128)
        at 
org.apache.avro.generic.GenericDatumWriter.write(GenericDatumWriter.java:75)
        at 
org.apache.avro.generic.GenericDatumWriter.writeField(GenericDatumWriter.java:166)
        at 
org.apache.avro.generic.GenericDatumWriter.writeRecord(GenericDatumWriter.java:156)
        at 
org.apache.avro.generic.GenericDatumWriter.writeWithoutConversion(GenericDatumWriter.java:118)
        ... 27 more
   ```
   
   **Expected behavior**
   
   When MySql rows contain null column value, the processing of 
hoodie_deltastreamer can be success
   
   **Environment Description**
   * Running on Docker? (yes) :
   
   
   


-- 
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]


Reply via email to