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]