[
https://issues.apache.org/jira/browse/HUDI-7028?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=17782367#comment-17782367
]
Lin Liu commented on HUDI-7028:
-------------------------------
To reproduce the first error:
{code:java}
import scala.collection.JavaConversions._
import org.apache.spark.sql.SaveMode._
import org.apache.hudi.DataSourceReadOptions._
import org.apache.hudi.DataSourceWriteOptions._
import org.apache.hudi.common.table.HoodieTableConfig._
import org.apache.hudi.config.HoodieWriteConfig._
import org.apache.hudi.keygen.constant.KeyGeneratorOptions._
import org.apache.hudi.common.model.HoodieRecord
import spark.implicits._val tableName = "trips_table"
val basePath = "file:///tmp/trips_table_1"val columns =
Seq("ts","uuid","rider","driver","fare","city")
val data =
Seq((1695159649087L,"334e26e9-8355-45cc-97c6-c31daf0df330","rider-A","driver-K",19.10,"san_francisco"),
(1695091554788L,"e96c4396-3fad-413a-a942-4cb36106d721","rider-C","driver-M",27.70
,"san_francisco"),
(1695046462179L,"9909a8b1-2d15-4d3d-8ec9-efc48c536a00","rider-D","driver-L",33.90
,"san_francisco"),
(1695516137016L,"e3cf430c-889d-4015-bc98-59bdce1e530c","rider-F","driver-P",34.15,"sao_paulo"
),
(1695115999911L,"c8abbe79-8d89-47ea-b4ce-4d224bae5bfa","rider-J","driver-T",17.85,"chennai"));var
inserts = spark.createDataFrame(data).toDF(columns:_*)
inserts.write.format("hudi").
option(PARTITIONPATH_FIELD_NAME.key(), "city").
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 tripsDF = spark.read.
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)
tripsDF.createOrReplaceTempView("trips_table")spark.sql("SELECT uuid, fare, ts,
rider, driver, city FROM trips_table WHERE fare > 20.0").show()
spark.sql("SELECT _hoodie_commit_time, _hoodie_record_key,
_hoodie_partition_path, rider, driver, fare FROM trips_table").show(1000,
false)
val updatesDf = spark.read.
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).filter($"rider" ===
"rider-D").withColumn("fare", col("fare") * 10)updatesDf.write.format("hudi").
option(OPERATION_OPT_KEY, "upsert").
option(PARTITIONPATH_FIELD_NAME.key(), "city").
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-shell
val adjustedFareDF = spark.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).limit(2).
withColumn("fare", col("fare") * 10)adjustedFareDF.write.format("hudi").
option("hoodie.datasource.write.payload.class","com.payloads.CustomMergeIntoConnector").
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)
{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)