Skip to content

fix(spark): propagate incremental instant range to MOR file group reader - #19948

Merged
danny0405 merged 4 commits into
apache:masterfrom
fhan688:propagate-incremental-instant-range-to-MOR-file-group-reader
Sep 22, 2026
Merged

danny0405 merged 4 commits into
apache:masterfrom
fhan688:propagate-incremental-instant-range-to-MOR-file-group-reader

Conversation

@fhan688

@fhan688 fhan688 commented Sep 14, 2026

Copy link
Copy Markdown
Contributor

Describe the issue this Pull Request addresses

Part of #19286.

For pre-8 tables using a layout V1 timeline, Spark incremental queries can use the 1.x reader with completion-time query boundaries. The query analyzer converts these boundaries into a requested-time InstantRange.

When a MOR incremental query falls back to a full-table scan, the latest file slice can contain log records beyond the requested end instant. The Spark MOR file-group reader paths did not receive the analyzed InstantRange, so an out-of-range update could participate in record merging and mask the latest in-range version of the same key.

Summary and Changelog

This PR propagates the incremental query's InstantRange to the MOR file-group reader before record merging.

  • Exposes the analyzed InstantRange from MergeOnReadIncrementalRelationV2.
  • Propagates the range through the default HadoopFsRelation/file-format path.
  • Propagates the same range through the legacy HoodieMergeOnReadRDDV2 path.
  • Passes the range into the Spark reader context so base and log records are filtered before merging.
  • Keeps snapshot, V1 incremental, and other unaffected paths unchanged by defaulting to an empty range.
  • Adds regression coverage for V6 MOR tables read with the V8 incremental reader, including archived instants, full-table-scan fallback, and an out-of-range log update.
  • Verifies both the default HadoopFsRelation path and the legacy RDD path.

Impact

Fixes incremental read correctness for MOR pre-8 tables when the V8 reader falls back to a full-table scan.

There are no storage format, configuration, or public API changes. Other read paths continue to use an empty InstantRange and retain their existing behavior.

Risk Level

Low.

The change is limited to MOR incremental reads with an analyzed instant range. Regression tests cover both Spark file-group reader entry paths and confirm that a physically present log update after the query end instant is excluded before merging.

Verification:

  • TestIncrementalReadWithFileGroupReader: 9 tests passed, 0 failures, 0 errors.
  • Relevant modules compile successfully with 0 scalastyle errors.

Documentation Update

None. This is an internal correctness fix with no new user-facing configuration or API.

Contributor's checklist

  • Read through contributor's guide
  • Enough context is provided in the sections above
  • Adequate tests were added if applicable

@github-actions github-actions Bot added the size:M PR with lines of changes in (100, 300] label Sep 14, 2026
@codecov-commenter

codecov-commenter commented Sep 14, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 86.95652% with 3 lines in your changes missing coverage. Please review.
✅ Project coverage is 80.35%. Comparing base (648996f) to head (d0cbab7).
⚠️ Report is 45 commits behind head on master.

Files with missing lines Patch % Lines
...hudi/SparkFileFormatInternalRowReaderContext.scala 66.66% 2 Missing ⚠️
...pache/hudi/core/read/BaseHoodieTableFileIndex.java 75.00% 0 Missing and 1 partial ⚠️
Additional details and impacted files
@@             Coverage Diff              @@
##             master   #19948      +/-   ##
============================================
+ Coverage     80.20%   80.35%   +0.14%     
- Complexity    34661    34783     +122     
============================================
  Files          2546     2545       -1     
  Lines        142452   142646     +194     
  Branches      17330    17713     +383     
============================================
+ Hits         114260   114628     +368     
+ Misses        20281    20111     -170     
+ Partials       7911     7907       -4     
Components Coverage Δ
hudi-common 83.85% <75.00%> (+0.01%) ⬆️
hudi-client 83.42% <71.42%> (+0.14%) ⬆️
hudi-flink 85.70% <ø> (+0.12%) ⬆️
hudi-spark-datasource 73.79% <100.00%> (+0.50%) ⬆️
hudi-utilities 78.18% <ø> (+0.01%) ⬆️
hudi-cli 69.99% <ø> (ø)
hudi-hadoop 70.98% <ø> (+0.16%) ⬆️
hudi-sync 76.02% <ø> (+0.02%) ⬆️
hudi-io 81.57% <ø> (-0.04%) ⬇️
hudi-timeline-service 83.34% <ø> (ø)
hudi-cloud 81.00% <ø> (+<0.01%) ⬆️
hudi-kafka-connect 53.20% <ø> (ø)
Flag Coverage Δ
common-and-other-modules 52.23% <47.82%> (+0.24%) ⬆️
flink-integration-tests 49.41% <0.00%> (+0.30%) ⬆️
hadoop-mr-java-client 43.85% <0.00%> (-0.07%) ⬇️
integration-tests 13.45% <0.00%> (-0.01%) ⬇️
spark-client-hadoop-common 38.57% <39.13%> (+0.02%) ⬆️
spark-java-tests 52.32% <86.95%> (+0.16%) ⬆️
spark-scala-tests 46.95% <82.60%> (-0.01%) ⬇️
utilities 36.81% <73.91%> (-0.03%) ⬇️

Flags with carried forward coverage won't be shown. Click here to find out more.

Files with missing lines Coverage Δ
...apache/hudi/BaseSparkInternalRowReaderContext.java 89.47% <100.00%> (ø)
...rg/apache/hudi/HoodieHadoopFsRelationFactory.scala 82.92% <100.00%> (+0.21%) ⬆️
...scala/org/apache/hudi/HoodieMergeOnReadRDDV2.scala 80.23% <100.00%> (+0.70%) ⬆️
...apache/hudi/MergeOnReadIncrementalRelationV2.scala 88.78% <100.00%> (+2.37%) ⬆️
...parquet/HoodieFileGroupReaderBasedFileFormat.scala 84.30% <100.00%> (+0.05%) ⬆️
...pache/hudi/core/read/BaseHoodieTableFileIndex.java 90.00% <75.00%> (+0.12%) ⬆️
...hudi/SparkFileFormatInternalRowReaderContext.scala 73.42% <66.66%> (-1.47%) ⬇️

... and 119 files with indirect coverage changes

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@hudi-agent hudi-agent left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ 🤖 This review was generated by an AI agent and may contain mistakes. Please verify any suggestions before applying.

Thanks for working on this! This PR threads the incremental query's analyzed InstantRange from MergeOnReadIncrementalRelationV2 into the Spark MOR file-group reader context (both the HadoopFsRelation/file-format path and the legacy HoodieMergeOnReadRDDV2 path), so out-of-range log blocks and base records are excluded before merging rather than only being filtered afterward. I traced the range through HoodieReaderContext.applyInstantRangeFilter, the log scanner's block-level check, and the LSM reader, and confirmed the EXACT_MATCH set matches the existing In(_hoodie_commit_time, …) required filters, that command/rollback blocks are not affected, and that non-V2 paths still receive an empty range. No correctness issues found. A few style/readability suggestions in the inline comments. Please take a look, and this should be ready for a Hudi committer or PMC member to take it from here. Changes are consistent with existing patterns in the codebase; only a minor readability note on the growing set of auxiliary constructors.

cc @yihua

this(baseFileReader, filters, requiredFilters, storageConfiguration, tableConfig, None)
this(baseFileReader, filters, requiredFilters, storageConfiguration, tableConfig, None, HOption.empty())

def this(baseFileReader: SparkColumnarFileReader,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🤖 nit: we now have three Java-friendly auxiliary constructors here differing only by which trailing optional params are included. Might be worth a comment above them noting they exist purely to work around Scala default args not generating Java overloads, so future readers don't try to collapse them.

⚠️ AI-generated; verify before applying. React 👍/👎 to flag quality.

@fhan688

fhan688 commented Sep 16, 2026

Copy link
Copy Markdown
Contributor Author

please help review, thanks! @yihua

HoodieWriteConfig.AUTO_UPGRADE_VERSION.key -> "false")

batches.zipWithIndex.foreach { case ((data, operation), i) =>
write(data, "MERGE_ON_READ", 6, operation,

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Cover full-table-scan fallback on V8 tables

Could we parameterize this regression test over table versions 6 and 8? The range propagation applies to all V2 MOR incremental reads, while this fallback test only exercises a V6 table. A V8 case using completion-time query boundaries would also verify the conversion to the requested-time range passed to the reader.


val readOptions = Map(
DataSourceReadOptions.QUERY_TYPE.key -> DataSourceReadOptions.QUERY_TYPE_INCREMENTAL_OPT_VAL,
DataSourceReadOptions.START_COMMIT.key -> "000",

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Exercise a bounded start and the earliest range variant

This case only exercises an EXACT_MATCH range starting from "000". Consider adding (c3, c5] to verify exclusion of older base/log records, and earliest → c5 to exercise the CLOSED_CLOSED range with a null lower boundary. The former should return only k1=5, k2=4, k3=5, and k4=5.

this(storageConfig, tableConfig, Option.empty(), recordContext);
}

protected BaseSparkInternalRowReaderContext(StorageConfiguration<?> storageConfig,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

can we try to extend the existing the constructors instead of adding new ones?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

can we try to extend the existing the constructors instead of adding new ones?

Thanks, makes sense for the protected base constructor. I’ll extend it with instantRangeOpt and update the internal subclass/test call sites. I’ll keep the existing Java-friendly constructor on SparkFileFormatInternalRowReaderContext, since Java callers rely on that overload.

}

private boolean isBeforeTimelineStarts() {
if (START_COMMIT_EARLIEST.equalsIgnoreCase(incrementalQueryStartTime.get())) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

isBeforeTimelineStarts means the start insant might be archived, the semantics does not equal to reading from earliest.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

isBeforeTimelineStarts means the start insant might be archived, the semantics does not equal to reading from earliest.

Good point. earliest is a sentinel rather than an instant before the active timeline. I’ll keep the special handling but move it out of isBeforeTimelineStarts, so the method retains its original semantics

@danny0405

Copy link
Copy Markdown
Contributor

cc @cshuo can you enforce this PR to fix the timetravel query issue

@hudi-agent hudi-agent left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ 🤖 This review was generated by an AI agent and may contain mistakes. Please verify any suggestions before applying.

Thanks for working on this! The PR threads the analyzed incremental InstantRange into the Spark MOR file-group reader (both the HadoopFsRelation file-format path and the legacy HoodieMergeOnReadRDDV2 path) so out-of-range log blocks no longer participate in the merge and mask an in-range version of the same key, and it makes START_COMMIT_EARLIEST fall through to a full partition listing in the file index. I traced the range through applyInstantRangeFilter (commit-time ordinal is guaranteed present for MOR incremental), the log scanner's block skip, positional merging (positions come from the parquet row index, so pre-buffer filtering doesn't desync), and serialization — those all check out. One semantics question worth double-checking in the inline comment: the range is also applied on the non-full-scan legacy RDD path, where base rows were previously merged over the full history and post-filtered, so results can differ for event-time ordering / custom payloads. Please take a look at any inline comments, and this should be ready for a Hudi committer or PMC member to take it from here. A couple of small readability suggestions below; otherwise the change reads cleanly.

mergeType = mergeType,
fileSplits = fileSplits,
includedInstantTimeSet = Option(includedCommits.map(_.requestedTime).toSet),
instantRangeOpt = queryContext.getInstantRange,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🤖 This passes the range on the non-full-scan path too, and on this legacy RDD path requiredFilters is Seq.empty, so before this change base rows and pre-range log blocks were merged first and only then post-filtered by includedInstantTimeSet. With EVENT_TIME_ORDERING (or a partial-update payload), an in-range log record with a lower ordering value than the out-of-range base row used to be suppressed by the merge but will now be emitted on its own — is that the intended incremental semantics here, or should the range be limited to the fullTableScan case on this path?

⚠️ AI-generated; verify before applying. React 👍/👎 to flag quality.

trait MergeOnReadIncrementalRelation {
def listFileSplits(partitionFilters: Seq[Expression], dataFilters: Seq[Expression]): Map[InternalRow, Seq[FileSlice]]
def getRequiredFilters: Seq[Filter]
def getInstantRange: HOption[InstantRange] = HOption.empty()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🤖 nit: giving getInstantRange a default HOption.empty() in the trait means any other implementation (e.g. V1) silently opts out of range filtering. Have you considered leaving it abstract like getRequiredFilters, so each relation makes the choice explicit?

⚠️ AI-generated; verify before applying. React 👍/👎 to flag quality.

}

private boolean isBeforeTimelineStarts() {
if (START_COMMIT_EARLIEST.equalsIgnoreCase(incrementalQueryStartTime.get())) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🤖 nit: a one-line comment on why the earliest marker short-circuits here (it's not a real timestamp, so the completion/requested-time comparisons below would be meaningless) would help — could also be pulled into a small isStartCommitEarliest(...) helper next to START_COMMIT_EARLIEST in IncrementalQueryAnalyzer if that check is repeated elsewhere.

⚠️ AI-generated; verify before applying. React 👍/👎 to flag quality.

@hudi-agent hudi-agent left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ 🤖 This review was generated by an AI agent and may contain mistakes. Please verify any suggestions before applying.

Thanks for the updates! The new change moves the START_COMMIT_EARLIEST sentinel check out of isBeforeTimelineStarts() and into the listPartitionPaths condition, so the helper keeps its original "start instant may be archived" meaning while earliest still falls through to a full partition listing — this addresses @danny0405's semantics concern, and the listing behavior is unchanged (the .get() is still guarded by isPresent()). The earlier PR revision also extended the protected BaseSparkInternalRowReaderContext constructor with instantRangeOpt rather than adding a new one, per the constructor discussion. My two prior nits (abstract getInstantRange in the trait; a short comment on the sentinel short-circuit) and the open question about applying the range on the non-full-scan legacy RDD path under EVENT_TIME_ORDERING / custom payloads haven't been responded to yet, but none of them are blocking — the semantics question is the one I'd still appreciate a quick confirmation on. No new issues flagged from this pass — a Hudi committer or PMC member can take it from here for a final review.

@fhan688 fhan688 closed this Sep 20, 2026
@fhan688 fhan688 reopened this Sep 20, 2026

@cshuo cshuo left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

+1.

* reader needs this in addition to Spark's required filters so that out-of-range log records do
* not participate in record merging and mask an earlier in-range version of the same key.
*/
override def getInstantRange: HOption[InstantRange] = queryContext.getInstantRange

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could we scope the new instant-range propagation to full-table-scan fallback? The reported bug occurs when fallback selects the latest file slices, allowing out-of-range log blocks to participate in merging. This would keep the fix focused on that case and preserve the existing behavior of normal incremental reads.

In MergeOnReadIncrementalRelationV2:

override def getInstantRange: HOption[InstantRange] =
  if (fullTableScan) queryContext.getInstantRange
  else HOption.empty()

Then use the same accessor in composeRDD:

instantRangeOpt = getInstantRange,

The HadoopFsRelation factory already calls this accessor, so both reader paths would follow the same condition. We can keep incrementalSpanRecordFilters unchanged and retain the fallback regression coverage for both paths and table versions 6/8, alongside the normal incremental-read tests.

This is a scope suggestion for the fallback fix, rather than a claim that file selection guarantees every log block is in range on all normal-read paths.

@hudi-agent hudi-agent left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ 🤖 This review was generated by an AI agent and may contain mistakes. Please verify any suggestions before applying.

Thanks for the follow-up on this! The PR threads the analyzed InstantRange from MergeOnReadIncrementalRelationV2 (now scoped to the full-table-scan fallback) into both the HadoopFsRelation file-format path and the legacy HoodieMergeOnReadRDDV2 path so out-of-range log records are excluded before file-group merging. I traced the range through HoodieReaderContext.applyInstantRangeFilter / FileGroupReaderSchemaHandler (commit-time field is added to the required schema whenever a range is present, and both incremental relations already request _hoodie_commit_time) and the LogScanningRecordBufferLoader log-side filter; the non-fallback and snapshot paths still receive an empty range. No new critical correctness issues flagged from this automated pass — a Hudi committer or PMC member can take it from here for a final review. Code is clean overall; one small readability suggestion on the growing partition-listing predicate.

cc @yihua

try {
if (isPartitionedTable()) {
if (queryType == HoodieTableQueryType.INCREMENTAL && incrementalQueryStartTime.isPresent() && !isBeforeTimelineStarts()) {
if (queryType == HoodieTableQueryType.INCREMENTAL

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🤖 nit: this predicate now has four clauses; it might be worth pulling it into a small helper like canListPartitionsFromIncrementalTimeline() so the intent ("we have a concrete start instant that is still on the active timeline") reads at a glance.

⚠️ AI-generated; verify before applying. React 👍/👎 to flag quality.

@hudi-bot

Copy link
Copy Markdown
Collaborator

CI report:

Bot commands @hudi-bot supports the following commands:
  • @hudi-bot run azure re-run the last Azure build

@danny0405
danny0405 merged commit 319afce into apache:master Sep 22, 2026
23 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

size:M PR with lines of changes in (100, 300]

Projects

None yet

Development

Successfully merging this pull request may close these issues.

7 participants