leaves12138 commented on code in PR #10172:
URL: https://github.com/apache/paimon/pull/10172#discussion_r4101491887
##########
paimon-core/src/main/java/org/apache/paimon/table/source/DataEvolutionTableRead.java:
##########
@@ -75,43 +66,42 @@ public DataEvolutionTableRead(
@Override
public RecordReader<InternalRow> createReader(Split split) throws
IOException {
QueryAuthContext queryAuthContext = unwrapQueryAuthSplit(split);
- final Split dataSplit;
- boolean filterOnRead = executeFilter;
- if (queryAuthContext.split() instanceof LazyIndexedSplit) {
- if (fileIO == null) {
- throw new IllegalStateException("FileIO is required for lazy
index evaluation.");
+ if (queryAuthContext.split() instanceof IndexQuerySplit) {
+ return createIndexQueryReader(
+ (IndexQuerySplit) queryAuthContext.split(),
queryAuthContext);
+ }
+ return createSelectedReader(queryAuthContext.split(),
queryAuthContext, executeFilter);
+ }
+
+ private RecordReader<InternalRow> createIndexQueryReader(
+ IndexQuerySplit split, QueryAuthContext queryAuthContext) throws
IOException {
+ final IndexedSplit indexedSplit;
+ try {
+ indexedSplit = split.evaluate(fileIO);
+ } catch (IOException e) {
+ if (!ExceptionUtils.findThrowable(
+ e,
+ cause ->
+ cause instanceof
FileNotFoundException
+ || cause instanceof
NoSuchFileException)
+ .isPresent()
+ || options.scalarIndexSearchMode() ==
CoreOptions.GlobalIndexSearchMode.FAST) {
+ throw e;
}
- LazyIndexedSplit lazySplit = (LazyIndexedSplit)
queryAuthContext.split();
- Split selectedSplit;
- try {
- IndexedSplit indexedSplit = lazySplit.evaluate(fileIO);
- if (indexedSplit.rowRanges().isEmpty()) {
- return new EmptyRecordReader<>();
- }
- selectedSplit = indexedSplit;
- } catch (IOException e) {
- if (!ExceptionUtils.findThrowable(
- e,
- cause ->
- cause instanceof
FileNotFoundException
- || cause instanceof
NoSuchFileException)
- .isPresent()
- || options.scalarIndexSearchMode()
- == CoreOptions.GlobalIndexSearchMode.FAST) {
- throw e;
- }
- if (predicate() == null) {
- throw new IOException(
- "Cannot scan a split without its index and query
filter", e);
- }
- selectedSplit = lazySplit.dataSplit();
- filterOnRead = true;
+ if (predicate() == null) {
+ throw new IOException("Cannot scan a split without its index
and query filter", e);
}
- dataSplit = selectedSplit;
- } else {
- dataSplit = queryAuthContext.split();
+ return createSelectedReader(split.dataSplit(), queryAuthContext,
true);
Review Comment:
[P1] Preserve the replay sequence when falling back after index loss
This forces exact residual filtering on fallback, while the normal indexed
path at line 99 uses `executeFilter`, which is false for the standard Flink
source. Therefore the two paths can emit different candidate sequences even
when physical row ordering is identical. I reproduced this in
`IndexQuerySplitTest`: write rows 0..99, index only `f1`, use `f1 startsWith
'a' AND f2 = 'b50'` in FULL mode, and read without `executeFilter()` as
`FlinkSource.createReader` does. Before index loss, the split emits candidates
0..99; checkpointing after candidate 0 records `recordsToSkip=1`.
Serialize/restore the same split and delete its planned BTree files: this
fallback emits only row 50, so replaying the saved skip consumes the only
matching row and produces `[]` instead of `[50]`.
`FileStoreSourceSplitReader.seek` skips from exactly this reader output.
This confirms the PR's documented recovery blocker, rather than only a
hypothetical ordering issue. Removing the implicit protection tag also makes
index retirement relevant under normal retention. Before merging, ensure normal
and fallback paths preserve the same replay sequence (including residual-filter
placement), use a stable recovery position, or fail safely when that cannot be
guaranteed. Please add a restore-after-partial-consumption test with an
unindexed residual and missing index files; the existing recovery test
explicitly enables `executeFilter()` and does not exercise this transition.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]