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