fix(spark): propagate incremental instant range to MOR file group reader - #19948
Conversation
Codecov Report❌ Patch coverage is 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
Flags with carried forward coverage won't be shown. Click here to find out more.
🚀 New features to boost your workflow:
|
hudi-agent
left a comment
There was a problem hiding this comment.
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, |
There was a problem hiding this comment.
🤖 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.
|
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, |
There was a problem hiding this comment.
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", |
There was a problem hiding this comment.
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, |
There was a problem hiding this comment.
can we try to extend the existing the constructors instead of adding new ones?
There was a problem hiding this comment.
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())) { |
There was a problem hiding this comment.
isBeforeTimelineStarts means the start insant might be archived, the semantics does not equal to reading from earliest.
There was a problem hiding this comment.
isBeforeTimelineStartsmeans 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
|
cc @cshuo can you enforce this PR to fix the timetravel query issue |
hudi-agent
left a comment
There was a problem hiding this comment.
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, |
There was a problem hiding this comment.
🤖 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?
| trait MergeOnReadIncrementalRelation { | ||
| def listFileSplits(partitionFilters: Seq[Expression], dataFilters: Seq[Expression]): Map[InternalRow, Seq[FileSlice]] | ||
| def getRequiredFilters: Seq[Filter] | ||
| def getInstantRange: HOption[InstantRange] = HOption.empty() |
There was a problem hiding this comment.
🤖 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?
| } | ||
|
|
||
| private boolean isBeforeTimelineStarts() { | ||
| if (START_COMMIT_EARLIEST.equalsIgnoreCase(incrementalQueryStartTime.get())) { |
There was a problem hiding this comment.
🤖 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.
hudi-agent
left a comment
There was a problem hiding this comment.
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.
| * 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 |
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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 |
There was a problem hiding this comment.
🤖 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.
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
InstantRangeto the MOR file-group reader before record merging.InstantRangefromMergeOnReadIncrementalRelationV2.HoodieMergeOnReadRDDV2path.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
InstantRangeand 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.Documentation Update
None. This is an internal correctness fix with no new user-facing configuration or API.
Contributor's checklist