[ 
https://issues.apache.org/jira/browse/HUDI-7028?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=17782369#comment-17782369
 ] 

Lin Liu commented on HUDI-7028:
-------------------------------

To reproduce the second error:
{code:java}
import org.apache.hudi.QuickstartUtils._
import scala.collection.JavaConversions._
import org.apache.spark.sql.SaveMode._
import org.apache.hudi.DataSourceReadOptions._
import org.apache.hudi.DataSourceWriteOptions._
import org.apache.hudi.config.HoodieWriteConfig._
import org.apache.hudi.common.model.HoodieRecordval tableName = "hudi_trips_cow"
val basePath = "file:///tmp/hudi_trips_cow"
val dataGen = new DataGenerator
val inserts = convertToStringList(dataGen.generateInserts(10))
val df = spark.read.json(spark.sparkContext.parallelize(inserts, 2))
df.write.format("hudi").
  options(getQuickstartWriteConfigs).
  option(PRECOMBINE_FIELD_OPT_KEY, "ts").
  option(RECORDKEY_FIELD_OPT_KEY, "uuid").
  option(PARTITIONPATH_FIELD_OPT_KEY, "partitionpath").
  option(TABLE_NAME, tableName).
  option("hoodie.datasource.write.table.type", "MERGE_ON_READ").
  option("hoodie.logfile.data.block.format", "parquet").
  option("hoodie.datasource.write.record.merger.impls", 
"org.apache.hudi.HoodieSparkRecordMerger").
  option("hoodie.datasource.read.use.new.parquet.file.format", "true").
  option("hoodie.file.group.reader.enabled", "true").
  option("hoodie.write.record.positions", "true").
  mode(Overwrite).
  save(basePath)val tripsSnapshotDF = spark.
  read.
  option("hoodie.datasource.write.table.type", "MERGE_ON_READ").
  option("hoodie.logfile.data.block.format", "parquet").
  option("hoodie.datasource.write.record.merger.impls", 
"org.apache.hudi.HoodieSparkRecordMerger").
  option("hoodie.datasource.read.use.new.parquet.file.format", "true").
  option("hoodie.file.group.reader.enabled", "true").
  option("hoodie.write.record.positions", "true").
  format("hudi").
  load(basePath)
tripsSnapshotDF.createOrReplaceTempView("hudi_trips_snapshot")spark.sql("select 
fare, begin_lon, begin_lat, ts from  hudi_trips_snapshot where fare > 
20.0").show()
spark.sql("select _hoodie_commit_time, _hoodie_record_key, 
_hoodie_partition_path, rider, driver, fare from  
hudi_trips_snapshot").show()val updates = 
convertToStringList(dataGen.generateUpdates(10))
val df = spark.read.json(spark.sparkContext.parallelize(updates, 2))
df.write.format("hudi").
  options(getQuickstartWriteConfigs).
  option(PRECOMBINE_FIELD_OPT_KEY, "ts").
  option(RECORDKEY_FIELD_OPT_KEY, "uuid").
  option(PARTITIONPATH_FIELD_OPT_KEY, "partitionpath").
  option(TABLE_NAME, tableName).
  option("hoodie.datasource.write.table.type", "MERGE_ON_READ").
  option("hoodie.logfile.data.block.format", "parquet").
  option("hoodie.datasource.write.record.merger.impls", 
"org.apache.hudi.HoodieSparkRecordMerger").
  option("hoodie.datasource.read.use.new.parquet.file.format", "true").
  option("hoodie.file.group.reader.enabled", "true").
  option("hoodie.write.record.positions", "true").
  mode(Append).
  save(basePath)
spark.
  read.
  option("hoodie.datasource.write.table.type", "MERGE_ON_READ").
  option("hoodie.logfile.data.block.format", "parquet").
  option("hoodie.datasource.write.record.merger.impls", 
"org.apache.hudi.HoodieSparkRecordMerger").
  option("hoodie.datasource.read.use.new.parquet.file.format", "true").
  option("hoodie.file.group.reader.enabled", "true").
  option("hoodie.write.record.positions", "true").
  format("hudi").
  load(basePath).
  createOrReplaceTempView("hudi_trips_snapshot")val commits = spark.sql("select 
distinct(_hoodie_commit_time) as commitTime from  hudi_trips_snapshot order by 
commitTime").map(k => k.getString(0)).take(50) {code}

> Fix Spark Quick Start
> ---------------------
>
>                 Key: HUDI-7028
>                 URL: https://issues.apache.org/jira/browse/HUDI-7028
>             Project: Apache Hudi
>          Issue Type: Bug
>            Reporter: Lin Liu
>            Priority: Major
>             Fix For: 1.0.0
>
>
> Fix the bugs for Spark quick start when turning on file group reader and 
> positional merging flag.
>  
> List some issues found so far:
>  # [compatibility]When no positions are stored in the header, the read query 
> failed. Idea behavior: use key based merging instead of failing.
>  # [compatibility]When a parquet file contains Avro records, the file group 
> reader of spark job will check if the payload is the expected type; 
> otherwise, it will throw.
>  #  



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

Reply via email to