damccorm commented on code in PR #40205:
URL: https://github.com/apache/beam/pull/40205#discussion_r4076095288
##########
sdks/java/core/src/main/java/org/apache/beam/sdk/fn/splittabledofn/RestrictionTrackers.java:
##########
@@ -139,32 +151,29 @@ protected RestrictionTrackerObserverWithProgress(
@Override
public Progress getProgress() {
- return getProgress(FIRST_PROGRESS_TIMEOUT_SEC);
+ return getProgress(PROGRESS_TIMEOUT_SEC);
}
@VisibleForTesting
Progress getProgress(int timeOutSec) {
- if (!hasInitialProgress) {
- Progress progress = Progress.NONE;
- try {
- // lock can be held long by long-running tryClaim/trySplit. We
tolerate this scenario
- // by returning zero progress when initial progress never evaluated
before due to lock
- // timeout.
- if (lock.tryLock(timeOutSec, TimeUnit.SECONDS)) {
- try {
- progress = getProgressBlocking();
- hasInitialProgress = true;
- } finally {
- lock.unlock();
- }
+ try {
+ // lock can be held long by long-running tryClaim/trySplit. We
tolerate this scenario
+ // by returning the last evaluated progress (or zero progress if never
evaluated before)
+ // when lock timeout occurs.
+ if (lock.tryLock(timeOutSec, TimeUnit.SECONDS)) {
+ try {
+ updateProgressBlocking();
+ } finally {
+ lock.unlock();
Review Comment:
I didn't do this because `updateProgressAndUnlock()` implies that we're
always updating progress, which we're not. But I see the downside, and I think
I prefer your approach.
--
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]