nsivabalan commented on code in PR #20077:
URL: https://github.com/apache/hudi/pull/20077#discussion_r4171022177
##########
hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/execution/datasources/parquet/HoodieFileGroupReaderBasedFileFormat.scala:
##########
@@ -356,97 +339,35 @@ class HoodieFileGroupReaderBasedFileFormat(tablePath:
String,
}
val broadcastedStorageConf = spark.sparkContext.broadcast(new
SerializableConfiguration(augmentedStorageConf.unwrap()))
- val fileIndexProps: TypedProperties =
HoodieFileIndex.getConfigProperties(spark, options, null)
+ val cdcProps: TypedProperties = HoodieFileIndex.getConfigProperties(spark,
options, null)
+ cdcProps.setProperty(HoodieTableConfig.HOODIE_TABLE_NAME_KEY, tableName)
val engineContext = new HoodieSparkEngineContext(new
JavaSparkContext(spark.sparkContext))
val maxMemoryPerCompaction =
MergeUtils.getMaxMemoryPerCompaction(engineContext.getTaskContextSupplier,
options.asJava)
- // Create metaclient on driver to avoid expensive operations on executors
- val metaClient: HoodieTableMetaClient = HoodieTableMetaClient
- .builder().setConf(augmentedStorageConf).setBasePath(tablePath).build
-
- (file: PartitionedFile) => {
- // executor
- val storageConf = new
HadoopStorageConfiguration(broadcastedStorageConf.value.value)
- val iter = file.partitionValues match {
- // Snapshot or incremental queries.
- case fileSliceMapping: HoodiePartitionFileSliceMapping =>
- val fileGroupName = FSUtils.getFileIdFromFilePath(sparkAdapter
- .getSparkPartitionedFileUtils.getPathFromPartitionedFile(file))
- fileSliceMapping.getSlice(fileGroupName) match {
- case Some(fileSlice) if !isCount && (requiredSchema.nonEmpty ||
fileSlice.getLogFiles.findAny().isPresent) =>
- // requiredFilters preserve Spark's row-level filtering
semantics, while instantRangeOpt
- // keeps out-of-range records from participating in the
file-group merge itself.
- val readerContext = new SparkFileFormatInternalRowReaderContext(
- fileGroupBaseFileReader.value, filters, requiredFilters,
storageConf, metaClient.getTableConfig,
- sparkRequiredSchema = Some(requiredSchema), instantRangeOpt =
instantRangeOpt)
-
readerContext.enableLogicalTimestampFieldRepair(storageConf.getBoolean(ENABLE_LOGICAL_TIMESTAMP_REPAIR,
true))
- val props = metaClient.getTableConfig.getProps
- options.foreach(kv => props.setProperty(kv._1, kv._2))
- props.put(HoodieMemoryConfig.MAX_MEMORY_FOR_MERGE.key(),
String.valueOf(maxMemoryPerCompaction))
- val baseFileLength = if (fileSlice.getBaseFile.isPresent) {
- fileSlice.getBaseFile.get.getFileSize
- } else {
- 0
- }
- val reader: HoodieRecordReader[InternalRow] =
- if (LsmReaderUtils.shouldUseLsmReader(
- metaClient.getTableConfig,
- ConfigUtils.getStringWithAltKeys(props,
HoodieReaderConfig.MERGE_TYPE, true))) {
- HoodieLsmFileGroupReader.builder[InternalRow]()
- .withReaderContext(readerContext)
- .withHoodieTableMetaClient(metaClient)
- .withLatestCommitTime(queryTimestamp)
- .withBaseFileOption(fileSlice.getBaseFile)
- .withLogFiles(fileSlice.getLogFiles)
- .withPartitionPath(fileSlice.getPartitionPath)
- .withDataSchema(dataSchema)
- .withRequestedSchema(requestedSchema)
- .withInternalSchemaOpt(internalSchemaOpt)
- .withProps(props)
- .withStart(file.start)
- .withLength(baseFileLength)
- .build()
- } else {
- HoodieFileGroupReader.builder[InternalRow]()
- .withReaderContext(readerContext)
- .withHoodieTableMetaClient(metaClient)
- .withLatestCommitTime(queryTimestamp)
- .withBaseFileOption(fileSlice.getBaseFile)
- .withLogFiles(fileSlice.getLogFiles)
- .withPartitionPath(fileSlice.getPartitionPath)
- .withDataSchema(dataSchema)
- .withRequestedSchema(requestedSchema)
- .withInternalSchemaOpt(internalSchemaOpt)
- .withProps(props)
- .withStart(file.start)
- .withLength(baseFileLength)
- .withShouldUseRecordPosition(shouldUseRecordPosition)
- .build()
- }
- // Append partition values to rows and project to output schema
- appendPartitionAndProject(
- reader.getClosableIterator,
- projectionInputSchema,
- remainingPartitionSchema,
- outputSchema,
- fileSliceMapping.getPartitionValues,
- fixedPartitionIndexes)
-
- case _ =>
- readBaseFile(file, baseFileReader.value, requestedStructType,
remainingPartitionSchema, fixedPartitionIndexes,
- readRequiredSchema, partitionSchema, outputSchema, filters ++
requiredFilters, storageConf)
- }
- // CDC queries.
- case hoodiePartitionCDCFileGroupSliceMapping:
HoodiePartitionCDCFileGroupMapping =>
- buildCDCRecordIterator(hoodiePartitionCDCFileGroupSliceMapping,
fileGroupBaseFileReader.value, storageConf, fileIndexProps, requiredSchema,
metaClient)
-
- case _ =>
- readBaseFile(file, baseFileReader.value, requestedStructType,
remainingPartitionSchema, fixedPartitionIndexes,
- readRequiredSchema, partitionSchema, outputSchema, filters ++
requiredFilters, storageConf)
- }
- CloseableIteratorListener.addListener(iter)
+ // The relation's meta client carries the timeline the scan was planned
against; build one only when the
Review Comment:
Worth stating in the PR description and release note: this is a behavior
change on executors, not only a serde one.
`BaseHoodieLogRecordReader` (line ~303) and `LogReaderUtils` (line ~81) call
`getCommitsTimeline()` / `getActiveTimeline()` on the meta client the task
received. Before, the driver built a fresh client with no loaded timeline, so
each task listed the timeline on the executor and could see commits that
completed after planning. Now tasks see the timeline the file index planned
against. That is the better behavior (consistent snapshot, no per-task
listing), but it changes which log blocks count as committed, so it belongs
next to the read-option note.
Could you also confirm `HoodieFileIndex.refresh()` reloads this same meta
client instance, so a long-lived relation (cached DataFrame, temp view) keeps
the file index and the reader timeline in sync?
##########
hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/execution/datasources/parquet/HoodieFileGroupReaderBasedFileFormat.scala:
##########
@@ -356,97 +339,35 @@ class HoodieFileGroupReaderBasedFileFormat(tablePath:
String,
}
val broadcastedStorageConf = spark.sparkContext.broadcast(new
SerializableConfiguration(augmentedStorageConf.unwrap()))
- val fileIndexProps: TypedProperties =
HoodieFileIndex.getConfigProperties(spark, options, null)
+ val cdcProps: TypedProperties = HoodieFileIndex.getConfigProperties(spark,
options, null)
+ cdcProps.setProperty(HoodieTableConfig.HOODIE_TABLE_NAME_KEY, tableName)
val engineContext = new HoodieSparkEngineContext(new
JavaSparkContext(spark.sparkContext))
val maxMemoryPerCompaction =
MergeUtils.getMaxMemoryPerCompaction(engineContext.getTaskContextSupplier,
options.asJava)
- // Create metaclient on driver to avoid expensive operations on executors
- val metaClient: HoodieTableMetaClient = HoodieTableMetaClient
- .builder().setConf(augmentedStorageConf).setBasePath(tablePath).build
-
- (file: PartitionedFile) => {
- // executor
- val storageConf = new
HadoopStorageConfiguration(broadcastedStorageConf.value.value)
- val iter = file.partitionValues match {
- // Snapshot or incremental queries.
- case fileSliceMapping: HoodiePartitionFileSliceMapping =>
- val fileGroupName = FSUtils.getFileIdFromFilePath(sparkAdapter
- .getSparkPartitionedFileUtils.getPathFromPartitionedFile(file))
- fileSliceMapping.getSlice(fileGroupName) match {
- case Some(fileSlice) if !isCount && (requiredSchema.nonEmpty ||
fileSlice.getLogFiles.findAny().isPresent) =>
- // requiredFilters preserve Spark's row-level filtering
semantics, while instantRangeOpt
- // keeps out-of-range records from participating in the
file-group merge itself.
- val readerContext = new SparkFileFormatInternalRowReaderContext(
- fileGroupBaseFileReader.value, filters, requiredFilters,
storageConf, metaClient.getTableConfig,
- sparkRequiredSchema = Some(requiredSchema), instantRangeOpt =
instantRangeOpt)
-
readerContext.enableLogicalTimestampFieldRepair(storageConf.getBoolean(ENABLE_LOGICAL_TIMESTAMP_REPAIR,
true))
- val props = metaClient.getTableConfig.getProps
- options.foreach(kv => props.setProperty(kv._1, kv._2))
- props.put(HoodieMemoryConfig.MAX_MEMORY_FOR_MERGE.key(),
String.valueOf(maxMemoryPerCompaction))
- val baseFileLength = if (fileSlice.getBaseFile.isPresent) {
- fileSlice.getBaseFile.get.getFileSize
- } else {
- 0
- }
- val reader: HoodieRecordReader[InternalRow] =
- if (LsmReaderUtils.shouldUseLsmReader(
- metaClient.getTableConfig,
- ConfigUtils.getStringWithAltKeys(props,
HoodieReaderConfig.MERGE_TYPE, true))) {
- HoodieLsmFileGroupReader.builder[InternalRow]()
- .withReaderContext(readerContext)
- .withHoodieTableMetaClient(metaClient)
- .withLatestCommitTime(queryTimestamp)
- .withBaseFileOption(fileSlice.getBaseFile)
- .withLogFiles(fileSlice.getLogFiles)
- .withPartitionPath(fileSlice.getPartitionPath)
- .withDataSchema(dataSchema)
- .withRequestedSchema(requestedSchema)
- .withInternalSchemaOpt(internalSchemaOpt)
- .withProps(props)
- .withStart(file.start)
- .withLength(baseFileLength)
- .build()
- } else {
- HoodieFileGroupReader.builder[InternalRow]()
- .withReaderContext(readerContext)
- .withHoodieTableMetaClient(metaClient)
- .withLatestCommitTime(queryTimestamp)
- .withBaseFileOption(fileSlice.getBaseFile)
- .withLogFiles(fileSlice.getLogFiles)
- .withPartitionPath(fileSlice.getPartitionPath)
- .withDataSchema(dataSchema)
- .withRequestedSchema(requestedSchema)
- .withInternalSchemaOpt(internalSchemaOpt)
- .withProps(props)
- .withStart(file.start)
- .withLength(baseFileLength)
- .withShouldUseRecordPosition(shouldUseRecordPosition)
- .build()
- }
- // Append partition values to rows and project to output schema
- appendPartitionAndProject(
- reader.getClosableIterator,
- projectionInputSchema,
- remainingPartitionSchema,
- outputSchema,
- fileSliceMapping.getPartitionValues,
- fixedPartitionIndexes)
-
- case _ =>
- readBaseFile(file, baseFileReader.value, requestedStructType,
remainingPartitionSchema, fixedPartitionIndexes,
- readRequiredSchema, partitionSchema, outputSchema, filters ++
requiredFilters, storageConf)
- }
- // CDC queries.
- case hoodiePartitionCDCFileGroupSliceMapping:
HoodiePartitionCDCFileGroupMapping =>
- buildCDCRecordIterator(hoodiePartitionCDCFileGroupSliceMapping,
fileGroupBaseFileReader.value, storageConf, fileIndexProps, requiredSchema,
metaClient)
-
- case _ =>
- readBaseFile(file, baseFileReader.value, requestedStructType,
remainingPartitionSchema, fixedPartitionIndexes,
- readRequiredSchema, partitionSchema, outputSchema, filters ++
requiredFilters, storageConf)
- }
- CloseableIteratorListener.addListener(iter)
+ // The relation's meta client carries the timeline the scan was planned
against; build one only when the
+ // format is used without a relation. The field is null rather than None
on a deserialized format.
+ val metaClient: HoodieTableMetaClient =
Option(tableMetaClient).flatten.getOrElse(HoodieTableMetaClient
+ .builder().setConf(augmentedStorageConf).setBasePath(tablePath).build)
+ val readerProps = TypedProperties.copy(metaClient.getTableConfig.getProps)
+ options.foreach(kv => readerProps.setProperty(kv._1, kv._2))
Review Comment:
Non-blocking: a user who relied on the accidental override (passing e.g.
merge mode or payload class as a read option) now gets different results with
no signal. A one-time driver-side warning here, when an option key is a table
config key and its value differs from the persisted one, would make the
migration visible.
##########
hudi-common/src/main/java/org/apache/hudi/common/table/HoodieTableMetaClient.java:
##########
@@ -388,6 +389,7 @@ private void readObject(java.io.ObjectInputStream in)
throws IOException, ClassN
in.defaultReadObject();
storage = null; // will be lazily initialized
Review Comment:
Related to the blocking thread on `HoodieFileGroupReadState`: `getStorage()`
(line ~480) is the lazy init that every task on an executor will now hit
concurrently on this shared, deserialized instance. Please make it
`synchronized` or volatile double-checked.
##########
hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/execution/datasources/parquet/HoodieFileGroupReaderFunction.scala:
##########
@@ -0,0 +1,376 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.spark.sql.execution.datasources.parquet
+
+import org.apache.hudi.{HoodiePartitionCDCFileGroupMapping,
HoodiePartitionFileSliceMapping, HoodieTableSchema, SparkAdapterSupport,
SparkFileFormatInternalRowReaderContext}
+import org.apache.hudi.cdc.{CDCFileGroupIterator, HoodieCDCFileGroupSplit,
HoodieCDCFileIndex}
+import org.apache.hudi.common.config.{HoodieReaderConfig, TypedProperties}
+import org.apache.hudi.common.fs.FSUtils
+import org.apache.hudi.common.schema.HoodieSchema
+import org.apache.hudi.common.schema.internal.InternalSchema
+import org.apache.hudi.common.table.{HoodieTableMetaClient,
ParquetTableSchemaResolver}
+import org.apache.hudi.common.table.log.InstantRange
+import org.apache.hudi.common.table.read.{HoodieFileGroupReader,
HoodieRecordReader}
+import org.apache.hudi.common.table.read.lsm.{HoodieLsmFileGroupReader,
LsmReaderUtils}
+import org.apache.hudi.common.util.{ConfigUtils, Option => HOption}
+import org.apache.hudi.common.util.collection.ClosableIterator
+import org.apache.hudi.data.CloseableIteratorListener
+import
org.apache.hudi.io.storage.HoodieSparkParquetReader.ENABLE_LOGICAL_TIMESTAMP_REPAIR
+import org.apache.hudi.io.storage.VectorConversionUtils
+import org.apache.hudi.storage.StorageConfiguration
+import org.apache.hudi.storage.hadoop.HadoopStorageConfiguration
+
+import org.apache.hadoop.conf.Configuration
+import org.apache.parquet.schema.MessageType
+import org.apache.spark.SparkEnv
+import org.apache.spark.broadcast.Broadcast
+import
org.apache.spark.sql.HoodieCatalystExpressionUtils.generateUnsafeProjection
+import org.apache.spark.sql.catalyst.InternalRow
+import org.apache.spark.sql.catalyst.expressions.{JoinedRow, UnsafeProjection}
+import org.apache.spark.sql.execution.datasources.{PartitionedFile,
SparkColumnarFileReader, SparkSchemaTransformUtils}
+import org.apache.spark.sql.sources.Filter
+import org.apache.spark.sql.types.StructType
+import org.apache.spark.sql.vectorized.{ColumnarBatch, ColumnarBatchUtils}
+import org.apache.spark.util.{SerializableConfiguration, Utils}
+
+import java.io.Closeable
+import java.nio.ByteBuffer
+
+import scala.collection.JavaConverters.mapAsJavaMapConverter
+import scala.reflect.ClassTag
+
+/**
+ * Read-only state of one scan of [[HoodieFileGroupReaderBasedFileFormat]],
built on the driver and shared by every
+ * task of an executor through [[HoodieFileGroupReaderFunction]]. Executors
only fill thread-safe lazy caches in it;
Review Comment:
🚨 **Blocking: this claim does not hold yet for everything the state
reaches.**
Before this PR each task deserialized its own copy of the format and meta
client, so these lazy initializers were effectively single-threaded. Now every
concurrent task on an executor hits the same instances through `state`, and two
of them are not written for that:
- `HoodieTableMetaClient.getStorage()` is a plain `if (storage == null)
storage = ...` with no `volatile` or `synchronized`. `HoodieFileGroupReader`
(line 160) and `HoodieLsmFileGroupReader` (line 137) call it for every file.
Threads will race to build and overwrite the instance.
- `InternalSchema` builds `idToName`, `nameToId`, `idToField` and
`nameToPosition` lazily (lines ~99-131 and ~267) and the fields are not
`volatile`, so the `HashMap` is published unsafely to other threads.
`state.internalSchemaOpt` is what schema-on-read column resolution runs
through, so a torn read there is a data correctness problem rather than a perf
one. `HoodieSchema` already handles its equivalent caches with a documented
volatile benign race; `InternalSchema` should match.
Suggested fix, both small and local: `synchronized` (or volatile
double-checked init) on `getStorage()`, and `volatile` on the four
`InternalSchema` caches. With that, this doc comment becomes true.
--
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]