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


##########
spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/source/AsyncSparkMicroBatchPlanner.java:
##########
@@ -252,7 +252,9 @@ public synchronized StreamingOffset 
latestOffset(StreamingOffset startOffset, Re
       return StreamingOffset.START_OFFSET;
     }
 
-    if (table().currentSnapshot().timestampMillis() < 
readConf().streamFromTimestamp()) {
+    // Only a new stream starts from the timestamp. A resumed stream continues 
from its offset.
+    if (startOffset.equals(StreamingOffset.START_OFFSET)
+        && table().currentSnapshot().timestampMillis() < 
readConf().streamFromTimestamp()) {

Review Comment:
   Yes. Jobs often pass the current time as `stream-from-timestamp` on every 
launch. After a restart, the newest snapshot was then usually older than that 
timestamp, so this check returned -1 for a resumed stream. Spark used -1 as the 
end offset, which skipped everything committed while the job was down and wrote 
-1 into the offset log. Now a resumed stream continues from its offset. 
`resumeIgnoresLaterStreamFromTimestamp` tests this.
   



##########
spark/v4.2/spark/src/main/java/org/apache/iceberg/spark/source/SparkMicroBatchStream.java:
##########


Review Comment:
   Spark doesn't call it for this source. Since it implements 
`SupportsAdmissionControl`, `MicroBatchExecution` calls 
`latestOffset(startOffset, limit)` instead. `MicroBatchStream` still requires 
the no-arg version and we should probably fail if the no-arg version is called, 
but not sure this should be handled in this PR.



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