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


##########
spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java:
##########
@@ -215,12 +213,34 @@ public Offset latestOffset(Offset startOffset, ReadLimit 
limit) {
         "Invalid start offset: %s is not a StreamingOffset",
         startOffset);
 
+    StreamingOffset start = (StreamingOffset) startOffset;
+    // Resolve start offset once, if it hasn't been resolved already.
+    if (StreamingOffset.START_OFFSET.equals(start)) {
+      start = startingOffsetFromTimestamp();
+      if (StreamingOffset.START_OFFSET.equals(start)) {
+        return null;
+      }
+
+      // Spark starts the first batch at initialOffset(), also when it replays 
the batch after a
+      // restart. Store the resolved offset before Spark records the first 
batch's end offset in the
+      // checkpoint, so the replay starts from the same offset.
+      if (StreamingOffset.START_OFFSET.equals(initialOffset)) {
+        initialOffsetStore.storeInitialOffset(start);
+        this.initialOffset = start;
+      }
+    }
+
     // Initialize planner if not already done
     if (planner == null) {
-      initializePlanner((StreamingOffset) startOffset, null);
+      initializePlanner(start, null);
     }
 
-    return planner.latestOffset((StreamingOffset) startOffset, limit);
+    return planner.latestOffset(start, limit);
+  }
+
+  private StreamingOffset startingOffsetFromTimestamp() {
+    table.refresh();

Review Comment:
   In async mode, this adds a synchronous `table.refresh()` to every 
`latestOffset` call until a snapshot after the timestamp appears. With the 
default trigger, Spark waits only `spark.sql.streaming.pollingDelay` (10 ms) 
between empty triggers. A job started with the timestamp set to "now" can then 
load table metadata many times a second until the next commit, and again after 
each restart. Before this change, the async planner's background thread did 
this every `streaming-snapshot-polling-interval-ms`, which defaults to 30 s. 
Could this lookup follow the polling interval in async mode, or stay with the 
async planner?



##########
spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/source/StreamingInitialOffsetStore.java:
##########
@@ -50,16 +50,28 @@ class StreamingInitialOffsetStore {
   StreamingOffset initialOffset() {
     InputFile inputFile = io.newInputFile(initialOffsetLocation);
     if (inputFile.exists()) {
-      return readOffset(inputFile);
+      StreamingOffset offset = readOffset(inputFile);
+      if (!StreamingOffset.START_OFFSET.equals(offset)) {
+        return offset;
+      }
     }
 
+    // START_OFFSET has no position yet, so it is derived again instead of 
stored
     StreamingOffset offset = offsetSupplier.get();

Review Comment:
   Thanks for handling checkpoints that already have `-1` stored. I'm curious 
about the stuck case from the description. The old code stored `-1`, Spark 
recorded batch 0 with end `E`, and the job stopped before committing it. After 
upgrading, Spark replays `planInputPartitions(initialOffset(), E)`:
   
   - If nothing exists after the new timestamp, `initialOffset()` is still 
`START_OFFSET`. `SyncSparkMicroBatchPlanner.planFiles` then derives the start 
from the timestamp again and fails with the same "Cannot load current offset at 
snapshot -1".
   - If a snapshot does exist after the new timestamp, this overwrites 
`offsets/0` with a start that comes after `E`. The replay can't reach `E`. And 
because the new start is now persisted, restarting with the original timestamp 
no longer recovers the query.
   
   What do you think about returning a stored `START_OFFSET` unchanged here, 
and letting `latestOffset` resolve and store it? A batch 0 that Spark hasn't 
recorded yet already takes that path, and it never replaces the start under a 
batch Spark has recorded.



##########
spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java:
##########
@@ -247,7 +267,7 @@ public void prepareForTriggerAvailableNow() {
     lastOffsetForTriggerAvailableNow =

Review Comment:
   When no snapshot at or after `stream-from-timestamp` exists yet, 
`latestOffset` now returns `null` here, so `lastOffsetForTriggerAvailableNow` 
ends up `null`. Both planners treat a `null` cap as no cap. If a snapshot is 
committed between `prepareForTriggerAvailableNow` and the first `latestOffset` 
call, the run reads it and keeps picking up new commits until a batch comes 
back empty. Before this change, the async planner got `START_OFFSET` as its cap 
and stopped. The sync planner was already uncapped in this case.
   
   I checked this locally with the test below. On this branch it fails for both 
`async=true` and `async=false`, because `latestOffset` returns the new snapshot 
at position 1.
   
   ```java
     @TestTemplate
     void availableNowIgnoresSnapshotsCommittedAfterStart() {
       appendData(List.of(new SimpleRecord(1, "one")));
       SparkMicroBatchStream stream =
           newMicroBatchStream(
               Map.of(
                   SparkReadOptions.STREAM_FROM_TIMESTAMP,
                   Long.toString(timestampAfterCurrentSnapshot())),
               "available-now-unresolved-start-checkpoint");
   
       try {
         stream.prepareForTriggerAvailableNow();
         appendData(List.of(new SimpleRecord(2, "two")));
   
         assertThat(stream.latestOffset(stream.initialOffset(), 
stream.getDefaultReadLimit()))
             .isNull();
       } finally {
         stream.stop();
       }
     }
   ```
   
   Would it work to have `prepareForTriggerAvailableNow` remember that nothing 
was available, and then return `null` from `latestOffset` for the rest of the 
run?



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