damccorm opened a new pull request, #40251:
URL: https://github.com/apache/beam/pull/40251

   ### Why This Is Needed
   
   In #40205, we updated 
`RestrictionTrackers.RestrictionTrackerObserverWithProgress.getProgress()` to 
use `lock.tryLock(60, TimeUnit.SECONDS)` and return the cached `lastProgress` 
on lock timeout so that a long-running `tryClaim()` (such as a filtered scan in 
`BoundedSourceAsSDFWrapperFn`) does not block progress reporting indefinitely.
   
   However, in runner harnesses (such as Dataflow's worker progress/split 
update loop), `getProgress()` and `trySplit()` execute sequentially within the 
**same** control-plane progress update cycle and share the same 10-minute lease 
renewal deadline:
   1. `getProgress()` is called to gather bundle progress. With #40205, if 
`tryClaim()` is holding `RestrictionTrackers.lock` for an extended period (e.g. 
>10 minutes), `getProgress()` times out after 60 seconds and returns 
`lastProgress`.
   2. That unchanged `lastProgress` is sent to the runner backend. Because the 
work item's reported progress is not advancing while other shards complete, 
dynamic work rebalancing (liquid sharding) flags the shard as a straggler and 
replies with a split request (`ProcessBundleSplitRequest`).
   3. The worker control loop then immediately invokes 
`RestrictionTrackerObserver.trySplit()` as part of handling the update 
response—which calls `lock.lock()` on the **exact same `ReentrantLock`** still 
held by the long-running `tryClaim()`.
   4. `trySplit()` blocks indefinitely on `lock.lock()`, stalling the worker's 
progress/lease-renewal loop until the 10-minute deadline expires and fails the 
work item lease.
   
   ### Solution & Why It Is Safe
   
   This change updates `RestrictionTrackerObserver.trySplit()` to use 
`lock.tryLock(300, TimeUnit.SECONDS)` (5 minutes) instead of an unbounded 
`lock.lock()`, returning `null` if the lock cannot be acquired within the 
timeout:
   
   - **Follows the SDF `trySplit` contract**: In Apache Beam's Splittable 
`DoFn` specification, `RestrictionTracker.trySplit()` is always advisory and is 
explicitly permitted to return `null` whenever a restriction cannot be split at 
the current moment.
   - **No state mutation on timeout**: When `lock.tryLock` times out, 
`delegate.trySplit()` is never invoked, so the active restriction, 
`ClaimObserver` state, and cached progress remain untouched.
   - **Automatic retry on subsequent update cycles**: Returning `null` causes 
`FnApiDoFnRunner` to return an empty `ProcessBundleSplitResponse`, allowing the 
runner's progress/update loop to complete normally and keep renewing the work 
item lease. On subsequent progress update cycles, the runner can continue 
requesting splits, which will succeed as soon as the slow `tryClaim()` finishes 
and releases `lock`.
   - **5-minute timeout**: Using a longer 5-minute timeout (300s) ensures 
splits still wait out moderately slow claims while guaranteeing that 
`getProgress()` (60s) + `trySplit()` (300s) always complete well within a 
10-minute runner lease renewal window.


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