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]
