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]

Reply via email to