This is an automated email from the ASF dual-hosted git repository.

jonvex pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/hudi.git


The following commit(s) were added to refs/heads/master by this push:
     new 7ae0555f1218 [HUDI-8290] Use Filegroup Reader in Spark Structured 
Streaming Read (#13503)
7ae0555f1218 is described below

commit 7ae0555f1218a859da6dc9b0bb3e30627aff0bd7
Author: Jon Vexler <[email protected]>
AuthorDate: Thu Jul 31 19:38:33 2025 -0400

    [HUDI-8290] Use Filegroup Reader in Spark Structured Streaming Read (#13503)
---
 .../main/scala/org/apache/hudi/DefaultSource.scala |  55 +++++-----
 .../scala/org/apache/hudi/HoodieCDCFileIndex.scala |   7 +-
 .../hudi/HoodieHadoopFsRelationFactory.scala       |  15 +--
 .../sql/FileFormatUtilsForFileGroupReader.scala    |   9 +-
 .../sql/hudi/streaming/HoodieStreamSourceV1.scala  |  67 ++++++++----
 .../sql/hudi/streaming/HoodieStreamSourceV2.scala  |  67 ++++++++----
 ...TestStreamSourceReadByStateTransitionTime.scala |  25 ++++-
 .../hudi/functional/TestStreamingSource.scala      | 116 ++++++++++++++++-----
 8 files changed, 262 insertions(+), 99 deletions(-)

diff --git 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/DefaultSource.scala
 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/DefaultSource.scala
index 0973471be41d..9269315de982 100644
--- 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/DefaultSource.scala
+++ 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/DefaultSource.scala
@@ -234,6 +234,16 @@ class DefaultSource extends RelationProvider
         HoodieSpecifiedOffsetRangeLimit(instantTime)
     }
 
+    // NOTE: In cases when Hive Metastore is used as catalog and the table is 
partitioned, schema in the HMS might contain
+    //       Hive-specific partitioning columns created specifically for HMS 
to handle partitioning appropriately. In that
+    //       case  we opt in to not be providing catalog's schema, and instead 
force Hudi relations to fetch the schema
+    //       from the table itself
+    val userSchema = if (isUsingHiveCatalog(sqlContext.sparkSession)) {
+      None
+    } else {
+      schema
+    }
+
     val storageConf = 
HadoopFSUtils.getStorageConf(sqlContext.sparkSession.sessionState.newHadoopConf())
     val tablePath: StoragePath = {
       val path = new StoragePath(parameters.getOrElse("path", "Missing 'path' 
option"))
@@ -247,24 +257,17 @@ class DefaultSource extends RelationProvider
     // does the checkpoint management based on the version.
     val targetTableVersion: Integer = Integer.parseInt(
       parameters.getOrElse(WRITE_TABLE_VERSION.key, 
WRITE_TABLE_VERSION.defaultValue.toString))
-    if (SparkConfigUtils.containsConfigProperty(parameters, 
STREAMING_READ_TABLE_VERSION)) {
-      val sourceTableVersion = 
Integer.parseInt(parameters(STREAMING_READ_TABLE_VERSION.key))
-      if (sourceTableVersion >= HoodieTableVersion.EIGHT.versionCode()) {
-        new HoodieStreamSourceV2(
-          sqlContext, metaClient, metadataPath, schema, parameters, 
offsetRangeLimit, HoodieTableVersion.fromVersionCode(targetTableVersion))
-      } else {
-        new HoodieStreamSourceV1(
-          sqlContext, metaClient, metadataPath, schema, parameters, 
offsetRangeLimit, HoodieTableVersion.fromVersionCode(targetTableVersion))
-      }
+    val readTableVersion = if 
(SparkConfigUtils.containsConfigProperty(parameters, 
STREAMING_READ_TABLE_VERSION)) {
+      Integer.parseInt(parameters(STREAMING_READ_TABLE_VERSION.key))
     } else {
-      val sourceTableVersion = 
metaClient.getTableConfig.getTableVersion.versionCode()
-      if (sourceTableVersion >= HoodieTableVersion.EIGHT.versionCode()) {
-        new HoodieStreamSourceV2(
-          sqlContext, metaClient, metadataPath, schema, parameters, 
offsetRangeLimit, HoodieTableVersion.fromVersionCode(targetTableVersion))
-      } else {
-        new HoodieStreamSourceV1(
-          sqlContext, metaClient, metadataPath, schema, parameters, 
offsetRangeLimit, HoodieTableVersion.fromVersionCode(targetTableVersion))
-      }
+      metaClient.getTableConfig.getTableVersion.versionCode()
+    }
+    if (readTableVersion >= HoodieTableVersion.EIGHT.versionCode()) {
+      new HoodieStreamSourceV2(
+        sqlContext, metaClient, metadataPath, userSchema, parameters, 
offsetRangeLimit, HoodieTableVersion.fromVersionCode(targetTableVersion))
+    } else {
+      new HoodieStreamSourceV1(
+        sqlContext, metaClient, metadataPath, userSchema, parameters, 
offsetRangeLimit, HoodieTableVersion.fromVersionCode(targetTableVersion))
     }
   }
 }
@@ -304,8 +307,8 @@ object DefaultSource {
         Option(schema)
       }
 
-      val useNewParquetFileFormat = 
parameters.getOrElse(HoodieReaderConfig.FILE_GROUP_READER_ENABLED.key(),
-        
HoodieReaderConfig.FILE_GROUP_READER_ENABLED.defaultValue().toString).toBoolean 
&&
+      lazy val enableFileGroupReader = SparkConfigUtils
+        .getStringWithAltKeys(parameters, 
HoodieReaderConfig.FILE_GROUP_READER_ENABLED).toBoolean &&
         !metaClient.isMetadataTable && (globPaths == null || globPaths.isEmpty)
       lazy val tableVersion = if 
(SparkConfigUtils.containsConfigProperty(parameters, 
INCREMENTAL_READ_TABLE_VERSION)) {
         Integer.parseInt(parameters(INCREMENTAL_READ_TABLE_VERSION.key))
@@ -316,7 +319,7 @@ object DefaultSource {
       if 
(metaClient.getCommitsTimeline.filterCompletedInstants.countInstants() == 0) {
         new EmptyRelation(sqlContext, resolveSchema(metaClient, parameters, 
Some(schema)))
       } else if (isCdcQuery) {
-        if (useNewParquetFileFormat) {
+        if (enableFileGroupReader) {
           if (tableType == COPY_ON_WRITE) {
             new HoodieCopyOnWriteCDCHadoopFsRelationFactory(
               sqlContext, metaClient, parameters, userSchema, isBootstrap = 
false).build()
@@ -333,14 +336,14 @@ object DefaultSource {
           case (COPY_ON_WRITE, QUERY_TYPE_SNAPSHOT_OPT_VAL, false) |
                (COPY_ON_WRITE, QUERY_TYPE_READ_OPTIMIZED_OPT_VAL, false) |
                (MERGE_ON_READ, QUERY_TYPE_READ_OPTIMIZED_OPT_VAL, false) =>
-            if (useNewParquetFileFormat) {
+            if (enableFileGroupReader) {
               new HoodieCopyOnWriteSnapshotHadoopFsRelationFactory(
                 sqlContext, metaClient, parameters, userSchema, isBootstrap = 
false).build()
             } else {
               resolveBaseFileOnlyRelation(sqlContext, globPaths, userSchema, 
metaClient, parameters)
             }
           case (COPY_ON_WRITE, QUERY_TYPE_INCREMENTAL_OPT_VAL, _) =>
-            (hoodieTableSupportsCompletionTime, useNewParquetFileFormat) match 
{
+            (hoodieTableSupportsCompletionTime, enableFileGroupReader) match {
               case (true, true) => new 
HoodieCopyOnWriteIncrementalHadoopFsRelationFactoryV2(
                 sqlContext, metaClient, parameters, userSchema, 
isBootstrappedTable, RangeType.CLOSED_CLOSED).build()
               case (true, false) => new IncrementalRelationV2(sqlContext, 
parameters, userSchema, metaClient, RangeType.CLOSED_CLOSED)
@@ -350,7 +353,7 @@ object DefaultSource {
             }
 
           case (MERGE_ON_READ, QUERY_TYPE_SNAPSHOT_OPT_VAL, false) =>
-            if (useNewParquetFileFormat) {
+            if (enableFileGroupReader) {
               new HoodieMergeOnReadSnapshotHadoopFsRelationFactory(
                 sqlContext, metaClient, parameters, userSchema, isBootstrap = 
false).build()
             } else {
@@ -358,7 +361,7 @@ object DefaultSource {
             }
 
           case (MERGE_ON_READ, QUERY_TYPE_SNAPSHOT_OPT_VAL, true) =>
-            if (useNewParquetFileFormat) {
+            if (enableFileGroupReader) {
               new HoodieMergeOnReadSnapshotHadoopFsRelationFactory(
                 sqlContext, metaClient, parameters, userSchema, isBootstrap = 
true).build()
             } else {
@@ -366,7 +369,7 @@ object DefaultSource {
             }
 
           case (MERGE_ON_READ, QUERY_TYPE_INCREMENTAL_OPT_VAL, _) =>
-            (hoodieTableSupportsCompletionTime, useNewParquetFileFormat) match 
{
+            (hoodieTableSupportsCompletionTime, enableFileGroupReader) match {
               case (true, true) => new 
HoodieMergeOnReadIncrementalHadoopFsRelationFactoryV2(
                 sqlContext, metaClient, parameters, userSchema, 
isBootstrappedTable).build()
               case (true, false) => 
MergeOnReadIncrementalRelationV2(sqlContext, parameters, metaClient, userSchema)
@@ -376,7 +379,7 @@ object DefaultSource {
             }
 
           case (_, _, true) =>
-            if (useNewParquetFileFormat) {
+            if (enableFileGroupReader) {
               new HoodieCopyOnWriteSnapshotHadoopFsRelationFactory(
                 sqlContext, metaClient, parameters, userSchema, isBootstrap = 
true).build()
             } else {
diff --git 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieCDCFileIndex.scala
 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieCDCFileIndex.scala
index c2413f33d427..ae0082ff0cea 100644
--- 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieCDCFileIndex.scala
+++ 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieCDCFileIndex.scala
@@ -22,13 +22,13 @@ package org.apache.hudi
 import org.apache.hudi.cdc.CDCRelation
 import org.apache.hudi.common.table.HoodieTableMetaClient
 import org.apache.hudi.common.table.cdc.HoodieCDCExtractor
+import org.apache.hudi.common.table.log.InstantRange.RangeType
 
 import org.apache.hadoop.fs.{FileStatus, Path}
 import org.apache.spark.sql.SparkSession
 import org.apache.spark.sql.catalyst.InternalRow
 import org.apache.spark.sql.catalyst.expressions.{Expression, 
GenericInternalRow}
 import org.apache.spark.sql.execution.datasources.{FileIndex, FileStatusCache, 
NoopCache, PartitionDirectory}
-import org.apache.spark.sql.sources.Filter
 import org.apache.spark.sql.types.StructType
 
 import scala.collection.JavaConverters._
@@ -38,12 +38,13 @@ class HoodieCDCFileIndex(override val spark: SparkSession,
                          override val schemaSpec: Option[StructType],
                          override val options: Map[String, String],
                          @transient override val fileStatusCache: 
FileStatusCache = NoopCache,
-                         override val includeLogFiles: Boolean)
+                         override val includeLogFiles: Boolean,
+                         val rangeType: RangeType)
   extends HoodieFileIndex(
     spark, metaClient, schemaSpec, options, fileStatusCache, includeLogFiles, 
shouldEmbedFileSlices = true
   ) with FileIndex  {
   private val emptyPartitionPath: String = "empty_partition_path";
-  val cdcRelation: CDCRelation = CDCRelation.getCDCRelation(spark.sqlContext, 
metaClient, options)
+  val cdcRelation: CDCRelation = CDCRelation.getCDCRelation(spark.sqlContext, 
metaClient, options, rangeType)
   val cdcExtractor: HoodieCDCExtractor = cdcRelation.cdcExtractor
 
   override def listFiles(partitionFilters: Seq[Expression], dataFilters: 
Seq[Expression]): Seq[PartitionDirectory] = {
diff --git 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieHadoopFsRelationFactory.scala
 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieHadoopFsRelationFactory.scala
index 3bbd61075203..960822b16fd9 100644
--- 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieHadoopFsRelationFactory.scala
+++ 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/hudi/HoodieHadoopFsRelationFactory.scala
@@ -377,18 +377,20 @@ class 
HoodieMergeOnReadIncrementalHadoopFsRelationFactoryV2(override val sqlCont
                                                             override val 
metaClient: HoodieTableMetaClient,
                                                             override val 
options: Map[String, String],
                                                             override val 
schemaSpec: Option[StructType],
-                                                            isBootstrap: 
Boolean)
+                                                            isBootstrap: 
Boolean,
+                                                            rangeType: 
RangeType = RangeType.CLOSED_CLOSED)
   extends HoodieMergeOnReadIncrementalHadoopFsRelationFactory(sqlContext, 
metaClient, options, schemaSpec, isBootstrap,
-    MergeOnReadIncrementalRelationV2(sqlContext, options, metaClient, 
schemaSpec))
+    MergeOnReadIncrementalRelationV2(sqlContext, options, metaClient, 
schemaSpec, None, rangeType))
 
 class HoodieMergeOnReadCDCHadoopFsRelationFactory(override val sqlContext: 
SQLContext,
                                                   override val metaClient: 
HoodieTableMetaClient,
                                                   override val options: 
Map[String, String],
                                                   override val schemaSpec: 
Option[StructType],
-                                                  isBootstrap: Boolean)
+                                                  isBootstrap: Boolean,
+                                                  rangeType: RangeType = 
RangeType.OPEN_CLOSED)
   extends HoodieBaseMergeOnReadIncrementalHadoopFsRelationFactory(sqlContext, 
metaClient, options, schemaSpec, isBootstrap) {
   private val hoodieCDCFileIndex = new HoodieCDCFileIndex(
-    sparkSession, metaClient, schemaSpec, options, fileStatusCache, true)
+    sparkSession, metaClient, schemaSpec, options, fileStatusCache, true, 
rangeType)
 
   override def buildFileIndex(): HoodieFileIndex = hoodieCDCFileIndex
 
@@ -484,11 +486,12 @@ class 
HoodieCopyOnWriteCDCHadoopFsRelationFactory(override val sqlContext: SQLCo
                                                   override val metaClient: 
HoodieTableMetaClient,
                                                   override val options: 
Map[String, String],
                                                   override val schemaSpec: 
Option[StructType],
-                                                  isBootstrap: Boolean)
+                                                  isBootstrap: Boolean,
+                                                  rangeType: RangeType = 
RangeType.OPEN_CLOSED)
   extends HoodieBaseCopyOnWriteIncrementalHadoopFsRelationFactory(sqlContext, 
metaClient, options, schemaSpec, isBootstrap) {
 
   private val hoodieCDCFileIndex = new HoodieCDCFileIndex(
-    sparkSession, metaClient, schemaSpec, options, fileStatusCache, false)
+    sparkSession, metaClient, schemaSpec, options, fileStatusCache, false, 
rangeType)
   override def buildFileIndex(): HoodieFileIndex = hoodieCDCFileIndex
 
   override def buildDataSchema(): StructType = 
hoodieCDCFileIndex.cdcRelation.schema
diff --git 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/FileFormatUtilsForFileGroupReader.scala
 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/FileFormatUtilsForFileGroupReader.scala
index 4050cd2865f4..0eb0d683398b 100644
--- 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/FileFormatUtilsForFileGroupReader.scala
+++ 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/FileFormatUtilsForFileGroupReader.scala
@@ -22,7 +22,7 @@ import org.apache.hudi.{HoodieCDCFileIndex, 
SparkAdapterSupport, SparkHoodieTabl
 
 import org.apache.spark.sql.catalyst.expressions.{And, Attribute, Contains, 
EndsWith, EqualNullSafe, EqualTo, Expression, GreaterThan, GreaterThanOrEqual, 
In, IsNotNull, IsNull, LessThan, LessThanOrEqual, Literal, NamedExpression, 
Not, Or, StartsWith}
 import org.apache.spark.sql.catalyst.plans.logical.{Filter, LogicalPlan, 
Project}
-import org.apache.spark.sql.execution.datasources.HadoopFsRelation
+import org.apache.spark.sql.execution.datasources.{HadoopFsRelation, 
LogicalRelation}
 import org.apache.spark.sql.execution.datasources.parquet.{HoodieFormatTrait, 
ParquetFileFormat}
 import org.apache.spark.sql.types.{BooleanType, StructType}
 
@@ -124,4 +124,11 @@ object FileFormatUtilsForFileGroupReader extends 
SparkAdapterSupport {
       plan
     }
   }
+
+  def createStreamingDataFrame(sqlContext: SQLContext, relation: 
HadoopFsRelation, requiredSchema: StructType): DataFrame = {
+    val logicalRelation = LogicalRelation(relation, isStreaming = true)
+    val resolvedSchema = logicalRelation.resolve(requiredSchema, 
sqlContext.sparkSession.sessionState.analyzer.resolver)
+    Dataset.ofRows(sqlContext.sparkSession, 
applyFiltersToPlan(logicalRelation, requiredSchema, resolvedSchema,
+      relation.fileFormat.asInstanceOf[HoodieFormatTrait].getRequiredFilters))
+  }
 }
diff --git 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/hudi/streaming/HoodieStreamSourceV1.scala
 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/hudi/streaming/HoodieStreamSourceV1.scala
index 902f57c72591..a00cae06cb8d 100644
--- 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/hudi/streaming/HoodieStreamSourceV1.scala
+++ 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/hudi/streaming/HoodieStreamSourceV1.scala
@@ -19,19 +19,21 @@
 
 package org.apache.spark.sql.hudi.streaming
 
-import org.apache.hudi.{AvroConversionUtils, DataSourceReadOptions, 
IncrementalRelationV1, MergeOnReadIncrementalRelationV1, SparkAdapterSupport}
+import org.apache.hudi.{AvroConversionUtils, DataSourceReadOptions, 
HoodieCopyOnWriteCDCHadoopFsRelationFactory, 
HoodieCopyOnWriteIncrementalHadoopFsRelationFactoryV1, 
HoodieMergeOnReadCDCHadoopFsRelationFactory, 
HoodieMergeOnReadIncrementalHadoopFsRelationFactoryV1, IncrementalRelationV1, 
MergeOnReadIncrementalRelationV1, SparkAdapterSupport}
 import 
org.apache.hudi.DataSourceReadOptions.INCREMENTAL_READ_HANDLE_HOLLOW_COMMIT
 import org.apache.hudi.cdc.CDCRelation
+import org.apache.hudi.common.config.HoodieReaderConfig
 import org.apache.hudi.common.model.HoodieTableType
 import org.apache.hudi.common.table.{HoodieTableMetaClient, 
HoodieTableVersion, TableSchemaResolver}
 import org.apache.hudi.common.table.cdc.HoodieCDCUtils
 import org.apache.hudi.common.table.checkpoint.{CheckpointUtils, 
StreamerCheckpointV1}
 import 
org.apache.hudi.common.table.timeline.TimelineUtils.{handleHollowCommitIfNeeded,
 HollowCommitHandling}
 import 
org.apache.hudi.common.table.timeline.TimelineUtils.HollowCommitHandling._
+import org.apache.hudi.util.SparkConfigUtils
 
 import org.apache.spark.internal.Logging
 import org.apache.spark.rdd.RDD
-import org.apache.spark.sql.{DataFrame, SQLContext}
+import org.apache.spark.sql.{DataFrame, FileFormatUtilsForFileGroupReader, 
SQLContext}
 import org.apache.spark.sql.catalyst.InternalRow
 import org.apache.spark.sql.execution.streaming.{Offset, Source}
 import org.apache.spark.sql.hudi.streaming.HoodieSourceOffset.INIT_OFFSET
@@ -54,8 +56,14 @@ class HoodieStreamSourceV1(sqlContext: SQLContext,
                            writeTableVersion: HoodieTableVersion)
   extends Source with Logging with Serializable with SparkAdapterSupport {
 
+  private lazy val enableFileGroupReader = SparkConfigUtils
+    .getStringWithAltKeys(parameters, 
HoodieReaderConfig.FILE_GROUP_READER_ENABLED).toBoolean
+
+
   private lazy val tableType = metaClient.getTableType
 
+  private lazy val isBootstrappedTable = 
metaClient.getTableConfig.getBootstrapBasePath.isPresent
+
   private val isCDCQuery = CDCRelation.isCDCEnabled(metaClient) &&
     
parameters.get(DataSourceReadOptions.QUERY_TYPE.key).contains(DataSourceReadOptions.QUERY_TYPE_INCREMENTAL_OPT_VAL)
 &&
     
parameters.get(DataSourceReadOptions.INCREMENTAL_FORMAT.key).contains(DataSourceReadOptions.INCREMENTAL_FORMAT_CDC_VAL)
@@ -146,10 +154,21 @@ class HoodieStreamSourceV1(sqlContext: SQLContext,
           DataSourceReadOptions.START_COMMIT.key()-> 
startCommitTime(startOffset),
           DataSourceReadOptions.END_COMMIT.key() -> endOffset.offsetCommitTime
         )
-        val rdd = CDCRelation.getCDCRelation(sqlContext, metaClient, 
cdcOptions)
-          .buildScan0(HoodieCDCUtils.CDC_COLUMNS, Array.empty)
+        if (enableFileGroupReader) {
+          val relation = if (tableType == HoodieTableType.COPY_ON_WRITE) {
+            new HoodieCopyOnWriteCDCHadoopFsRelationFactory(
+              sqlContext, metaClient, parameters ++ cdcOptions, schemaOption, 
isBootstrappedTable).build()
+          } else {
+            new HoodieMergeOnReadCDCHadoopFsRelationFactory(
+              sqlContext, metaClient, parameters ++ cdcOptions, schemaOption, 
isBootstrappedTable).build()
+          }
+          
FileFormatUtilsForFileGroupReader.createStreamingDataFrame(sqlContext, 
relation, CDCRelation.FULL_CDC_SPARK_SCHEMA)
+        } else {
+          val rdd = CDCRelation.getCDCRelation(sqlContext, metaClient, 
cdcOptions)
+            .buildScan0(HoodieCDCUtils.CDC_COLUMNS, Array.empty)
 
-        sqlContext.sparkSession.internalCreateDataFrame(rdd, 
CDCRelation.FULL_CDC_SPARK_SCHEMA, isStreaming = true)
+          sqlContext.sparkSession.internalCreateDataFrame(rdd, 
CDCRelation.FULL_CDC_SPARK_SCHEMA, isStreaming = true)
+        }
       } else {
         // Consume the data between (startCommitTime, endCommitTime]
         val incParams = parameters ++ Map(
@@ -158,21 +177,31 @@ class HoodieStreamSourceV1(sqlContext: SQLContext,
           DataSourceReadOptions.END_COMMIT.key -> endOffset.offsetCommitTime,
           INCREMENTAL_READ_HANDLE_HOLLOW_COMMIT.key -> 
hollowCommitHandlingMode.name
         )
-
-        val rdd = tableType match {
-          case HoodieTableType.COPY_ON_WRITE =>
-            val serDe = sparkAdapter.createSparkRowSerDe(schema)
-            new IncrementalRelationV1(sqlContext, incParams, Some(schema), 
metaClient)
-              .buildScan()
-              .map(serDe.serializeRow)
-          case HoodieTableType.MERGE_ON_READ =>
-            val requiredColumns = schema.fields.map(_.name)
-            new MergeOnReadIncrementalRelationV1(sqlContext, incParams, 
metaClient, Some(schema))
-              .buildScan(requiredColumns, Array.empty[Filter])
-              .asInstanceOf[RDD[InternalRow]]
-          case _ => throw new IllegalArgumentException(s"UnSupport tableType: 
$tableType")
+        if (enableFileGroupReader) {
+          val relation = if (tableType == HoodieTableType.COPY_ON_WRITE) {
+            new 
HoodieCopyOnWriteIncrementalHadoopFsRelationFactoryV1(sqlContext, metaClient, 
incParams, Option(schema), isBootstrappedTable)
+              .build()
+          } else {
+            new 
HoodieMergeOnReadIncrementalHadoopFsRelationFactoryV1(sqlContext, metaClient, 
incParams, Option(schema), isBootstrappedTable)
+              .build()
+          }
+          
FileFormatUtilsForFileGroupReader.createStreamingDataFrame(sqlContext, 
relation, schema)
+        } else {
+          val rdd = tableType match {
+            case HoodieTableType.COPY_ON_WRITE =>
+              val serDe = sparkAdapter.createSparkRowSerDe(schema)
+              new IncrementalRelationV1(sqlContext, incParams, Some(schema), 
metaClient)
+                .buildScan()
+                .map(serDe.serializeRow)
+            case HoodieTableType.MERGE_ON_READ =>
+              val requiredColumns = schema.fields.map(_.name)
+              new MergeOnReadIncrementalRelationV1(sqlContext, incParams, 
metaClient, Some(schema))
+                .buildScan(requiredColumns, Array.empty[Filter])
+                .asInstanceOf[RDD[InternalRow]]
+            case _ => throw new IllegalArgumentException(s"UnSupport 
tableType: $tableType")
+          }
+          sqlContext.internalCreateDataFrame(rdd, schema, isStreaming = true)
         }
-        sqlContext.internalCreateDataFrame(rdd, schema, isStreaming = true)
       }
     }
   }
diff --git 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/hudi/streaming/HoodieStreamSourceV2.scala
 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/hudi/streaming/HoodieStreamSourceV2.scala
index 1abb9e264529..d5f984b9b044 100644
--- 
a/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/hudi/streaming/HoodieStreamSourceV2.scala
+++ 
b/hudi-spark-datasource/hudi-spark-common/src/main/scala/org/apache/spark/sql/hudi/streaming/HoodieStreamSourceV2.scala
@@ -17,17 +17,19 @@
 
 package org.apache.spark.sql.hudi.streaming
 
-import org.apache.hudi.{AvroConversionUtils, DataSourceReadOptions, 
IncrementalRelationV2, MergeOnReadIncrementalRelationV2, SparkAdapterSupport}
+import org.apache.hudi.{AvroConversionUtils, DataSourceReadOptions, 
HoodieCopyOnWriteCDCHadoopFsRelationFactory, 
HoodieCopyOnWriteIncrementalHadoopFsRelationFactoryV2, 
HoodieMergeOnReadCDCHadoopFsRelationFactory, 
HoodieMergeOnReadIncrementalHadoopFsRelationFactoryV2, IncrementalRelationV2, 
MergeOnReadIncrementalRelationV2, SparkAdapterSupport}
 import org.apache.hudi.cdc.CDCRelation
+import org.apache.hudi.common.config.HoodieReaderConfig
 import org.apache.hudi.common.model.HoodieTableType
 import org.apache.hudi.common.table.{HoodieTableMetaClient, 
HoodieTableVersion, TableSchemaResolver}
 import org.apache.hudi.common.table.cdc.HoodieCDCUtils
 import org.apache.hudi.common.table.checkpoint.{CheckpointUtils, 
StreamerCheckpointV2}
 import org.apache.hudi.common.table.log.InstantRange.RangeType
+import org.apache.hudi.util.SparkConfigUtils
 
 import org.apache.spark.internal.Logging
 import org.apache.spark.rdd.RDD
-import org.apache.spark.sql.{DataFrame, SQLContext}
+import org.apache.spark.sql.{DataFrame, FileFormatUtilsForFileGroupReader, 
SQLContext}
 import org.apache.spark.sql.catalyst.InternalRow
 import org.apache.spark.sql.execution.streaming.{Offset, Source}
 import org.apache.spark.sql.hudi.streaming.HoodieSourceOffset.INIT_OFFSET
@@ -52,6 +54,11 @@ class HoodieStreamSourceV2(sqlContext: SQLContext,
 
   private lazy val tableType = metaClient.getTableType
 
+  private lazy val isBootstrappedTable = 
metaClient.getTableConfig.getBootstrapBasePath.isPresent
+
+  private lazy val enableFileGroupReader = SparkConfigUtils
+    .getStringWithAltKeys(parameters, 
HoodieReaderConfig.FILE_GROUP_READER_ENABLED).toBoolean
+
   private val isCDCQuery = CDCRelation.isCDCEnabled(metaClient) &&
     
parameters.get(DataSourceReadOptions.QUERY_TYPE.key).contains(DataSourceReadOptions.QUERY_TYPE_INCREMENTAL_OPT_VAL)
 &&
     
parameters.get(DataSourceReadOptions.INCREMENTAL_FORMAT.key).contains(DataSourceReadOptions.INCREMENTAL_FORMAT_CDC_VAL)
@@ -123,10 +130,21 @@ class HoodieStreamSourceV2(sqlContext: SQLContext,
           DataSourceReadOptions.START_COMMIT.key() -> startCompletionTime,
           DataSourceReadOptions.END_COMMIT.key() -> endOffset.offsetCommitTime
         )
-        val rdd = CDCRelation.getCDCRelation(sqlContext, metaClient, 
cdcOptions, rangeType)
-          .buildScan0(HoodieCDCUtils.CDC_COLUMNS, Array.empty)
-
-        sqlContext.sparkSession.internalCreateDataFrame(rdd, 
CDCRelation.FULL_CDC_SPARK_SCHEMA, isStreaming = true)
+        if (enableFileGroupReader) {
+          val relation = if (tableType == HoodieTableType.COPY_ON_WRITE) {
+            new HoodieCopyOnWriteCDCHadoopFsRelationFactory(
+              sqlContext, metaClient, parameters ++ cdcOptions, schemaOption, 
isBootstrappedTable, rangeType).build()
+          } else {
+            new HoodieMergeOnReadCDCHadoopFsRelationFactory(
+              sqlContext, metaClient, parameters ++ cdcOptions, schemaOption, 
isBootstrappedTable, rangeType).build()
+          }
+          
FileFormatUtilsForFileGroupReader.createStreamingDataFrame(sqlContext, 
relation, CDCRelation.FULL_CDC_SPARK_SCHEMA)
+        } else {
+          val rdd = CDCRelation.getCDCRelation(sqlContext, metaClient, 
cdcOptions, rangeType)
+            .buildScan0(HoodieCDCUtils.CDC_COLUMNS, Array.empty)
+
+          sqlContext.sparkSession.internalCreateDataFrame(rdd, 
CDCRelation.FULL_CDC_SPARK_SCHEMA, isStreaming = true)
+        }
       } else {
         // Consume the data between (startCommitTime, endCommitTime]
         val incParams = parameters ++ Map(
@@ -135,20 +153,31 @@ class HoodieStreamSourceV2(sqlContext: SQLContext,
           DataSourceReadOptions.END_COMMIT.key -> endOffset.offsetCommitTime
         )
 
-        val rdd = tableType match {
-          case HoodieTableType.COPY_ON_WRITE =>
-            val serDe = sparkAdapter.createSparkRowSerDe(schema)
-            new IncrementalRelationV2(sqlContext, incParams, Some(schema), 
metaClient, rangeType)
-              .buildScan()
-              .map(serDe.serializeRow)
-          case HoodieTableType.MERGE_ON_READ =>
-            val requiredColumns = schema.fields.map(_.name)
-            new MergeOnReadIncrementalRelationV2(sqlContext, incParams, 
metaClient, Some(schema), rangeType = rangeType)
-              .buildScan(requiredColumns, Array.empty[Filter])
-              .asInstanceOf[RDD[InternalRow]]
-          case _ => throw new IllegalArgumentException(s"UnSupport tableType: 
$tableType")
+        if (enableFileGroupReader) {
+          val relation = if (tableType == HoodieTableType.COPY_ON_WRITE) {
+            new 
HoodieCopyOnWriteIncrementalHadoopFsRelationFactoryV2(sqlContext, metaClient, 
incParams, Option(schema), isBootstrappedTable, rangeType)
+              .build()
+          } else {
+            new 
HoodieMergeOnReadIncrementalHadoopFsRelationFactoryV2(sqlContext, metaClient, 
incParams, Option(schema), isBootstrappedTable, rangeType)
+              .build()
+          }
+          
FileFormatUtilsForFileGroupReader.createStreamingDataFrame(sqlContext, 
relation, schema)
+        } else {
+          val rdd = tableType match {
+            case HoodieTableType.COPY_ON_WRITE =>
+              val serDe = sparkAdapter.createSparkRowSerDe(schema)
+              new IncrementalRelationV2(sqlContext, incParams, Some(schema), 
metaClient, rangeType)
+                .buildScan()
+                .map(serDe.serializeRow)
+            case HoodieTableType.MERGE_ON_READ =>
+              val requiredColumns = schema.fields.map(_.name)
+              new MergeOnReadIncrementalRelationV2(sqlContext, incParams, 
metaClient, Some(schema), rangeType = rangeType)
+                .buildScan(requiredColumns, Array.empty[Filter])
+                .asInstanceOf[RDD[InternalRow]]
+            case _ => throw new IllegalArgumentException(s"UnSupport 
tableType: $tableType")
+          }
+          sqlContext.internalCreateDataFrame(rdd, schema, isStreaming = true)
         }
-        sqlContext.internalCreateDataFrame(rdd, schema, isStreaming = true)
       }
     }
   }
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestStreamSourceReadByStateTransitionTime.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestStreamSourceReadByStateTransitionTime.scala
index d5dc9c301728..d61963f365b8 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestStreamSourceReadByStateTransitionTime.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestStreamSourceReadByStateTransitionTime.scala
@@ -17,6 +17,8 @@
 
 package org.apache.hudi.functional
 
+import org.apache.hudi.DataSourceWriteOptions
+import org.apache.hudi.DataSourceWriteOptions.{PRECOMBINE_FIELD, 
RECORDKEY_FIELD}
 import org.apache.hudi.client.{SparkRDDWriteClient, WriteClientTestUtils}
 import org.apache.hudi.client.common.HoodieSparkEngineContext
 import org.apache.hudi.common.engine.EngineType
@@ -25,13 +27,33 @@ import org.apache.hudi.common.table.HoodieTableMetaClient
 import org.apache.hudi.common.testutils.HoodieTestDataGenerator
 import org.apache.hudi.common.testutils.HoodieTestTable.makeNewCommitTime
 import org.apache.hudi.config.{HoodieCleanConfig, HoodieWriteConfig}
+import org.apache.hudi.config.HoodieWriteConfig.{DELETE_PARALLELISM_VALUE, 
INSERT_PARALLELISM_VALUE, UPSERT_PARALLELISM_VALUE}
 import org.apache.hudi.hadoop.fs.HadoopFSUtils
 
 import org.apache.spark.api.java.JavaRDD
+import org.apache.spark.sql.streaming.StreamTest
 
 import scala.collection.JavaConverters._
 
-class TestStreamSourceReadByStateTransitionTime extends TestStreamingSource {
+class TestStreamSourceReadByStateTransitionTime extends StreamTest  {
+
+  protected val commonOptions: Map[String, String] = Map(
+    RECORDKEY_FIELD.key -> "id",
+    PRECOMBINE_FIELD.key -> "ts",
+    INSERT_PARALLELISM_VALUE.key -> "4",
+    UPSERT_PARALLELISM_VALUE.key -> "4",
+    DELETE_PARALLELISM_VALUE.key -> "4",
+    DataSourceWriteOptions.PARTITIONPATH_FIELD.key() -> "partition_path"
+  )
+
+  org.apache.log4j.Logger.getRootLogger.setLevel(org.apache.log4j.Level.WARN)
+
+  override protected def sparkConf = {
+    super.sparkConf
+      .set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
+      .set("spark.kryo.registrator", 
"org.apache.spark.HoodieSparkKryoRegistrar")
+      .set("spark.sql.extensions", 
"org.apache.spark.sql.hudi.HoodieSparkSessionExtension")
+  }
 
   private val dataGen = new HoodieTestDataGenerator(System.currentTimeMillis())
 
@@ -43,6 +65,7 @@ class TestStreamSourceReadByStateTransitionTime extends 
TestStreamingSource {
           .setTableType(tableType)
           .setTableName(s"test_stream_${tableType.name()}")
           .setPreCombineFields("timestamp")
+          .setPartitionFields("partition_path")
           
.initTable(HadoopFSUtils.getStorageConf(spark.sessionState.newHadoopConf()), 
tablePath)
 
         val writeConfig = HoodieWriteConfig.newBuilder()
diff --git 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestStreamingSource.scala
 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestStreamingSource.scala
index 57f39f75595a..b45f7b5cb037 100644
--- 
a/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestStreamingSource.scala
+++ 
b/hudi-spark-datasource/hudi-spark/src/test/scala/org/apache/hudi/functional/TestStreamingSource.scala
@@ -20,6 +20,7 @@ package org.apache.hudi.functional
 import org.apache.hudi.DataSourceReadOptions
 import org.apache.hudi.DataSourceReadOptions.{START_OFFSET, 
STREAMING_READ_TABLE_VERSION}
 import org.apache.hudi.DataSourceWriteOptions.{PRECOMBINE_FIELD, 
RECORDKEY_FIELD}
+import org.apache.hudi.common.model.HoodieTableType
 import org.apache.hudi.common.model.HoodieTableType.{COPY_ON_WRITE, 
MERGE_ON_READ}
 import org.apache.hudi.common.table.{HoodieTableMetaClient, HoodieTableVersion}
 import org.apache.hudi.common.table.timeline.HoodieTimeline
@@ -215,7 +216,7 @@ class TestStreamingSource extends StreamTest {
   test("Test mor streaming source with clustering") {
     Array("true", "false").foreach(skipCluster => {
       withTempDir { inputDir =>
-        val tablePath = s"${inputDir.getCanonicalPath}/test_mor_stream_cluster"
+        val tablePath = 
s"${inputDir.getCanonicalPath}/test_mor_stream_cluster_$skipCluster"
         val metaClient = HoodieTableMetaClient.newTableBuilder()
           .setTableType(MERGE_ON_READ)
           .setTableName(getTableName(tablePath))
@@ -242,12 +243,20 @@ class TestStreamingSource extends StreamTest {
 
         testStream(df)(
           AssertOnQuery { q => q.processAllAvailable(); true },
-          // Start after the first commit
-          CheckAnswerRows(Seq(
-            Row("2", "a1", "11", "001"),
-            Row("3", "a1", "12", "002"),
-            Row("4", "a1", "13", "003"),
-            Row("5", "a1", "14", "004")), lastOnly = true, isSorted = false)
+          if (skipCluster.toBoolean) {
+            // Start after the first commit
+            CheckAnswerRows(Seq(Row("5", "a1", "14", "004")), lastOnly = true, 
isSorted = false)
+          } else {
+            // Start after the first commit
+            CheckAnswerRows(Seq(
+              Row("2", "a1", "11", "001"),
+              Row("3", "a1", "12", "002"),
+              Row("4", "a1", "13", "003"),
+              Row("5", "a1", "14", "004")), lastOnly = true, isSorted = false)
+          }
+
+
+
         )
         assertTrue(metaClient.reloadActiveTimeline
           .filter(JavaConversions.getPredicate(
@@ -297,61 +306,120 @@ class TestStreamingSource extends StreamTest {
     })
   }
 
-  test("Test checkpoint translation") {
+  private def testCheckpointTranslation(tableName: String,
+                                        tableType: HoodieTableType,
+                                        writeTableVersion: HoodieTableVersion,
+                                        streamingReadVersions: List[Int]): 
Unit = {
     withTempDir { inputDir =>
-      val tablePath = s"${inputDir.getCanonicalPath}/test_cow_stream_ckpt"
+      val tablePath = s"${inputDir.getCanonicalPath}/$tableName"
       val metaClient = HoodieTableMetaClient.newTableBuilder()
-        .setTableType(COPY_ON_WRITE)
+        .setTableType(tableType)
         .setTableName(getTableName(tablePath))
+        .setTableVersion(writeTableVersion)
         .setRecordKeyFields("id")
         .setPreCombineFields("ts")
         
.initTable(HadoopFSUtils.getStorageConf(spark.sessionState.newHadoopConf()), 
tablePath)
 
-      addData(tablePath, Seq(("1", "a1", "10", "000")))
-      addData(tablePath, Seq(("2", "a1", "11", "001")))
-      addData(tablePath, Seq(("3", "a1", "12", "002")))
+      // Add initial data
+      addData(tablePath, Seq(("1", "a1", "10", "000")), tableVersion = 
writeTableVersion)
+      addData(tablePath, Seq(("2", "a1", "11", "001")), tableVersion = 
writeTableVersion)
+      addData(tablePath, Seq(("3", "a1", "12", "002")), tableVersion = 
writeTableVersion)
+
+      // Add update for MOR tests
+      if (tableType == MERGE_ON_READ) {
+        addData(tablePath, Seq(("2", "a2_updated", "16", "003")), tableVersion 
= writeTableVersion)
+      }
 
       val instants = 
metaClient.getActiveTimeline.getCommitsTimeline.filterCompletedInstants.getInstants
-      assertEquals(3, instants.size())
+      val expectedInstantCount = if (tableType == MERGE_ON_READ) 4 else 3
+      assertEquals(expectedInstantCount, instants.size())
+
+      val startTimestampIndex = if (tableType == MERGE_ON_READ) 2 else 1
+      val startTimestamp = instants.get(startTimestampIndex).requestedTime
 
-      // If the request time is used, i.e., V1, then the second record is 
included in the output.
-      // Otherwise, only third record in the output.
-      val startTimestamp = instants.get(1).requestedTime
-      for (streamingReadTableVersion <- 
List(HoodieTableVersion.SIX.versionCode(), 
HoodieTableVersion.EIGHT.versionCode())) {
+      for (streamingReadTableVersion <- streamingReadVersions) {
         val df = spark.readStream
           .format("org.apache.hudi")
           .option(START_OFFSET.key, startTimestamp)
-          .option(WRITE_TABLE_VERSION.key, 
HoodieTableVersion.current().versionCode().toString)
+          .option(WRITE_TABLE_VERSION.key, 
writeTableVersion.versionCode().toString)
           .option(STREAMING_READ_TABLE_VERSION.key, 
streamingReadTableVersion.toString)
           .load(tablePath)
           .select("id", "name", "price", "ts")
-        val expectedRows = if (streamingReadTableVersion == 
HoodieTableVersion.EIGHT.versionCode()) {
-          Seq(Row("2", "a1", "11", "001"), Row("3", "a1", "12", "002"))
+
+        val expectedRows = if (tableType == MERGE_ON_READ) {
+          if (streamingReadTableVersion == 
HoodieTableVersion.current().versionCode()) {
+            Seq(Row("3", "a1", "12", "002"), Row("2", "a2_updated", "16", 
"003"))
+          } else {
+            Seq(Row("2", "a2_updated", "16", "003"))
+          }
         } else {
-          Seq(Row("3", "a1", "12", "002"))
+          if (streamingReadTableVersion == 
HoodieTableVersion.current().versionCode()) {
+            Seq(Row("2", "a1", "11", "001"), Row("3", "a1", "12", "002"))
+          } else {
+            Seq(Row("3", "a1", "12", "002"))
+          }
         }
+
         testStream(df)(
           AssertOnQuery { q => q.processAllAvailable(); true },
-          // Start after the first commit
           CheckAnswerRows(expectedRows, lastOnly = true, isSorted = false)
         )
       }
     }
   }
 
+  test("Test checkpoint translation on COW table") {
+    testCheckpointTranslation(
+      "test_cow_stream_ckpt",
+      COPY_ON_WRITE,
+      HoodieTableVersion.current(),
+      List(HoodieTableVersion.SIX.versionCode(), 
HoodieTableVersion.current().versionCode())
+    )
+  }
+
+  test("Test checkpoint translation on MOR table") {
+    testCheckpointTranslation(
+      "test_mor_stream_ckpt",
+      MERGE_ON_READ,
+      HoodieTableVersion.current(),
+      List(HoodieTableVersion.SIX.versionCode(), 
HoodieTableVersion.current().versionCode())
+    )
+  }
+
+  test("Test checkpoint translation on COW table with table version 6") {
+    testCheckpointTranslation(
+      "test_cow_stream_ckpt_v6",
+      COPY_ON_WRITE,
+      HoodieTableVersion.SIX,
+      List(HoodieTableVersion.SIX.versionCode())
+    )
+  }
+
+  test("Test checkpoint translation on MOR table with table version 6") {
+    testCheckpointTranslation(
+      "test_mor_stream_ckpt_v6",
+      MERGE_ON_READ,
+      HoodieTableVersion.SIX,
+      List(HoodieTableVersion.SIX.versionCode())
+    )
+  }
+
   private def addData(inputPath: String,
                       rows: Seq[(String, String, String, String)],
                       enableInlineCompaction: Boolean = false,
-                      enableInlineCluster: Boolean = false) : Unit = {
+                      enableInlineCluster: Boolean = false,
+                      tableVersion: HoodieTableVersion = 
HoodieTableVersion.current) : Unit = {
     rows.toDF(columns: _*)
       .write
       .format("org.apache.hudi")
       .options(commonOptions)
       .option(TBL_NAME.key, getTableName(inputPath))
+      .option(WRITE_TABLE_VERSION.key, tableVersion.versionCode().toString)
       .option(HoodieCompactionConfig.INLINE_COMPACT.key(), 
enableInlineCompaction.toString)
       .option(HoodieCompactionConfig.INLINE_COMPACT_NUM_DELTA_COMMITS.key(), 
"2")
       .option(HoodieClusteringConfig.INLINE_CLUSTERING.key(), 
enableInlineCluster.toString)
       .option(HoodieClusteringConfig.INLINE_CLUSTERING_MAX_COMMITS.key(), "2")
+      .option(HoodieCompactionConfig.PARQUET_SMALL_FILE_LIMIT.key, "0")
       .mode(SaveMode.Append)
       .save(inputPath)
   }


Reply via email to