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]

Reply via email to