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