anuragmantri commented on code in PR #18348:
URL: https://github.com/apache/iceberg/pull/18348#discussion_r4186848615


##########
spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/source/AsyncSparkMicroBatchPlanner.java:
##########
@@ -267,10 +269,16 @@ public synchronized StreamingOffset 
latestOffset(StreamingOffset startOffset, Re
 
     // if we want to read all available we don't need to scan files, just 
snapshots
     if (limit instanceof ReadAllAvailable) {
+      // A START_OFFSET cap only means that no snapshot matched 
stream-from-timestamp when the run
+      // started. Returned as the end offset of a resumed stream, it would 
reset the stream.
+      boolean capResetsStream =

Review Comment:
   Thanks for slimming this down. The new condition is much easier to follow. I 
tried one more case on top of 
`availableNowResumesCheckpointWithUnresolvedInitialOffset`, where a snapshot is 
committed after `prepareForTriggerAvailableNow`. For a resumed stream with a 
`START_OFFSET` cap, this branch walks to whatever snapshot is latest when it's 
called, so the run also reads snapshots committed after it started. The sync 
planner does the same, because `latestSnapshotId` becomes `-1` and never 
matches. On a table with steady commits, the AvailableNow run keeps going for 
as long as new commits arrive.
   
   This variant fails on this branch for both `async=true` and `async=false`. 
`latestOffset` returns the snapshot committed after prepare.
   
   ```java
     @TestTemplate
     void availableNowResumeStopsAtSnapshotsAvailableAtStart() {
       appendData(List.of(new SimpleRecord(1, "one")));
       table.refresh();
       Snapshot committed = table.currentSnapshot();
       appendData(List.of(new SimpleRecord(2, "two")));
       table.refresh();
       Snapshot latest = table.currentSnapshot();
       SparkMicroBatchStream stream =
           newMicroBatchStream(
               ImmutableMap.of(
                   SparkReadOptions.STREAM_FROM_TIMESTAMP,
                   Long.toString(timestampAfterCurrentSnapshot())),
               "available-now-resume-cap-checkpoint");
   
       try {
         stream.prepareForTriggerAvailableNow();
         appendData(List.of(new SimpleRecord(3, "three")));
         StreamingOffset committedOffset =
             new StreamingOffset(
                 committed.snapshotId(), MicroBatchUtils.addedFilesCount(table, 
committed), false);
   
         assertThat(stream.latestOffset(committedOffset, 
stream.getDefaultReadLimit()))
             .isEqualTo(
                 new StreamingOffset(
                     latest.snapshotId(), 
MicroBatchUtils.addedFilesCount(table, latest), false));
       } finally {
         stream.stop();
       }
     }
   ```
   
   Could `prepareForTriggerAvailableNow` also record the snapshot that is 
current at that point, and both planners use it as the cap for a resumed 
stream? A new stream would keep the `START_OFFSET` cap it has today. That would 
also give the async planner a real end for its initial preload. Right now 
`initialPreloadEndSnapshot()` returns `table().snapshot(-1)`, which is `null`, 
so `fillQueueInitialBuffer` skips the preload as if the table were empty.



-- 
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]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to