jiangzzhu opened a new issue, #40232:
URL: https://github.com/apache/beam/issues/40232

   ### What happened?
   
   When reading from an unbounded Spanner Change Stream (or a Mutable Change 
Stream with a capped query end timestamp), `QueryChangeStreamAction` issues 
change stream queries with an artificial `changeStreamQueryEndTimestamp` (`now 
+ 2 minutes`).
   
   Previously, when the query completed and needed to resume 
(`!stopAfterQuerySucceeds`), `QueryChangeStreamAction` called 
`tracker.tryClaim(changeStreamQueryEndTimestamp)` before returning 
`ProcessContinuation.resume()`:
   
   
https://github.com/apache/beam/blob/ac4acbb6282/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/QueryChangeStreamAction.java#L360-L365
   
   If a change stream query finishes early without returning the 
`ChildPartitionsRecord` at the end of a partition's range, claiming 
`changeStreamQueryEndTimestamp` advances the restriction tracker past the 
partition's actual end timestamp. On the next continuation, the query resumes 
from `changeStreamQueryEndTimestamp + 1ns`, which falls outside the partition's 
valid timestamp range and returns an out-of-range `start_timestamp` error. 
`QueryChangeStreamAction` treats this out-of-range error as the end of the 
partition and marks the partition `FINISHED` without scheduling its child 
partitions, resulting in unscheduled child partitions and missing data from 
those child partitions.
   
   ### Expected Behavior
   
   When `!stopAfterQuerySucceeds`, `QueryChangeStreamAction` should leave the 
restriction tracker at the last claimed position (from the last processed data 
or heartbeat record) instead of advancing to `changeStreamQueryEndTimestamp`, 
so the subsequent query resumes from `lastClaimedTimestamp + 1ns` and reads any 
remaining records (including `ChildPartitionsRecord`) before the partition ends.
   
   ### Issue Priority
   
   Priority: 1 (data loss / correctness bug)
   
   ### Issue Components
   
   - [x] Component: Java SDK
   - [x] Component: IO Connectors


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

Reply via email to