mxm commented on code in PR #18348:
URL: https://github.com/apache/iceberg/pull/18348#discussion_r4185146224
##########
spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java:
##########
@@ -247,7 +267,7 @@ public void prepareForTriggerAvailableNow() {
lastOffsetForTriggerAvailableNow =
Review Comment:
This part is gone now that the PR is smaller. For a stream without an offset
the cap behaves as before, so the async planner still stops in your scenario.
The only change left is that it no longer returns a -1 cap as the end offset of
a resumed stream.
##########
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:
Gone as well. The start offset lookup is back in the planners, so the
refresh rate is unchanged.
--
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]