Skip to content

SAMZA-2266: Introduce a backoff when there are repeated failures for host-affinity allocations - #1104

Merged
rmatharu-zz merged 9 commits into
apache:masterfrom
dnishimura:samza-2266-host-affinity-retry-backoff
Aug 14, 2019
Merged

rmatharu-zz merged 9 commits into
apache:masterfrom
dnishimura:samza-2266-host-affinity-retry-backoff

Conversation

@dnishimura

Copy link
Copy Markdown
Contributor

Motivation
For host-affinity enabled jobs, a bad physical host may not immediately be marked as invalid by the Resource Manager (RM). As a result, when the HostAwareContainerAllocator requests preferred hosts, the RM generates the onResourceCompleted callback even though the host can't be allocated. The status error in the onResourceCompleted is equivalent to an application error and the retry logic kicks in to restart the failed container. Adding delays in the retry logic will prevent the job from failing prematurely (after 8 retries) before the bad host is marked invalid.

Implementation notes
Added an exponential back-off with a max delay. Container allocation requests are put in a priority queue with the priority determined by type and request timestamp. For retries that have a delay, I set the request timestamp in the future by time X where X is the calculated back-off.

Testing
Unit tests and tested a Samza job on a YARN cluster. I simulated the scenario by forcing an uncaught exception in a few containers to force the containers to fail.

@rmatharu @abhishekshivanna and others please take a look

public static final int DEFAULT_CONTAINER_RETRY_COUNT = 8;

public static final String CLUSTER_MANAGER_CONTAINER_RETRY_MAX_DELAY_MS = "cluster-manager.container.host-affinity-retry.max.delay.ms";
public static final long DEFAULT_CONTAINER_RETRY_MAX_DELAY_MS = Duration.ofSeconds(120).toMillis();

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Debating if this should be a non-configurable constant. Thoughts?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This should be in greater than the max amount of time it takes for the RM to mark a dead NM unhealthy,
which in YARN 2.9.2 is 5 minutes?
So perhaps not configurable and set to a higher value.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

As we discussed, as long as the total time across all retries is greater than 5 minutes, we are good.

@rmatharu-zz rmatharu-zz Jul 18, 2019 •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes maybe add 1 second to account for clock-skew and asynchrony.
(typically clock skew with a dc w/ ntp is 100s of micros).

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sure will add

*/
protected final SamzaResourceRequest peekPendingRequest() {
return resourceRequestState.peekPendingRequest();
protected final SamzaResourceRequest peekReadyPendingRequest() {

@rmatharu-zz rmatharu-zz Jul 15, 2019 •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Code nitpicks:
Perhaps use Optional.
Update to Javadoc:
Retrieves, but does not remove, the next pending request in the queue with request-timestamp >= current-timestamp.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Will change. thanks.


/**
* Called within {@link #onResourceCompleted(SamzaResourceStatus)} for unknown exit statuses. Usually these type of
* exit statuses are due to application errors causing the container resource to fail or for other unknown reasons.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nitpick:
These exit statuses correspond to container completion other than container run-to-completion, abort or preemption, or disk failure (e.g., detected by YARN's NM healthchecks).

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Will edit. Thanks.

@rmatharu-zz rmatharu-zz left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Taking a pass over the change

currentFailCount = 1;
lastFailureTime = Instant.now();
}
if (currentFailCount > maxRetryCount) {

@rmatharu-zz rmatharu-zz Jul 15, 2019 •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should be currentFailCount >= maxRetryCount ?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Max fail count should be 1 greater than the max retries since the initial attempt isn't considered a retry, correct? Examples: 1 max failure for 0 retries, 2 max failures for 1 retry, 3 max failures for 2 retries, etc...

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

got it, makes sense.
1.
The code used to have currentFailCount >= retryCount and we've changed it to
currentFailCount > maxRetryCount which means without any change to the configparam,
now the job would fail at 9 failures (1 initial failure + 8 failed retries). Perhaps we can add a lil comment to say that.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sure will add.

Duration retryDelay = getHostRetryDelay(lastSeenOn, currentFailCount - 1);
long currentRetryWindowMs = Instant.now().toEpochMilli() - lastFailureTime.plus(retryDelay).toEpochMilli();

if (currentRetryWindowMs < maxRetryWindowMs) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should be currentRetryWindowMs <= maxRetryWindowMs

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Will change thanks.

lastFailureTime = Instant.now();
}
if (currentFailCount > maxRetryCount) {
Duration retryDelay = getHostRetryDelay(lastSeenOn, currentFailCount - 1);

@rmatharu-zz rmatharu-zz Jul 15, 2019 •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

For readability, can getHostRetryDelay be made to accept currentFailCount?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Do you mean removing the lastSeenOn parameter or removing the - 1? Here I'm trying to get the delta of the retry delay window from the last failed request.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Removing the -1

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not obvious, but the retryDelay is always zero on the first failure, so the scenario you mentioned wouldn't happen. I'll wrap this logic in a private method and document the logic to make it more clear. Thanks.

@Override
public int compareTo(SamzaResourceRequest o) {
if (!StandbyTaskUtil.isStandbyContainer(this.processorId) && StandbyTaskUtil.isStandbyContainer(o.processorId)) {
if (!isInFuture() && o.isInFuture()) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Couple of things
a. The isInFuture() check here is redundant since there is the timestamp based check below.
b. Having this check here, above the check for standby vs. active is problematic.
Imagine a scenario, where two containers one active (e.g., container 1) and one standby container (e.g., container 2-0) fail concurrently.
If the failure-callback for the standby arrives before the callback of the active, the timestamp for the standby will be less than the timestamp for the active.
And at some point in time, the standby's request will be not in future, while the one for active will be in future.
This scenario will violate the priority-queue assumption that active containers take precedence over standby containers.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'll be adding a separate DelayedRequestQueue, so the "Is in future" logic will no longer be needed. I'll revert this change.

@rmatharu-zz rmatharu-zz left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Took an initial pass.
Mostly code and doc comments, but also a couple of correctness issues.

@dnishimura dnishimura left a comment

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks @rmatharu for the review. Please see my responses. I'll push the changes soon.


/**
* Called within {@link #onResourceCompleted(SamzaResourceStatus)} for unknown exit statuses. Usually these type of
* exit statuses are due to application errors causing the container resource to fail or for other unknown reasons.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Will edit. Thanks.

currentFailCount = 1;
lastFailureTime = Instant.now();
}
if (currentFailCount > maxRetryCount) {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Max fail count should be 1 greater than the max retries since the initial attempt isn't considered a retry, correct? Examples: 1 max failure for 0 retries, 2 max failures for 1 retry, 3 max failures for 2 retries, etc...

lastFailureTime = Instant.now();
}
if (currentFailCount > maxRetryCount) {
Duration retryDelay = getHostRetryDelay(lastSeenOn, currentFailCount - 1);

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Do you mean removing the lastSeenOn parameter or removing the - 1? Here I'm trying to get the delta of the retry delay window from the last failed request.

Duration retryDelay = getHostRetryDelay(lastSeenOn, currentFailCount - 1);
long currentRetryWindowMs = Instant.now().toEpochMilli() - lastFailureTime.plus(retryDelay).toEpochMilli();

if (currentRetryWindowMs < maxRetryWindowMs) {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Will change thanks.

@Override
public int compareTo(SamzaResourceRequest o) {
if (!StandbyTaskUtil.isStandbyContainer(this.processorId) && StandbyTaskUtil.isStandbyContainer(o.processorId)) {
if (!isInFuture() && o.isInFuture()) {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I'll be adding a separate DelayedRequestQueue, so the "Is in future" logic will no longer be needed. I'll revert this change.

public static final int DEFAULT_CONTAINER_RETRY_COUNT = 8;

public static final String CLUSTER_MANAGER_CONTAINER_RETRY_MAX_DELAY_MS = "cluster-manager.container.host-affinity-retry.max.delay.ms";
public static final long DEFAULT_CONTAINER_RETRY_MAX_DELAY_MS = Duration.ofSeconds(120).toMillis();

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

As we discussed, as long as the total time across all retries is greater than 5 minutes, we are good.

*/
protected final SamzaResourceRequest peekPendingRequest() {
return resourceRequestState.peekPendingRequest();
protected final SamzaResourceRequest peekReadyPendingRequest() {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Will change. thanks.

}
if (currentFailCount > maxRetryCount) {
Duration retryDelay = getHostRetryDelay(lastSeenOn, currentFailCount - 1);
long currentRetryWindowMs = Instant.now().toEpochMilli() - lastFailureTime.plus(retryDelay).toEpochMilli();

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

As we discussed, the problem here is that, because now the lastFailureTime is set to currentTime, as opposed to 0 (in the current code).
The currentRetryWindowMs value can now be < 0, as opposed to being the currTime value.

This means the check below on line 532 can erroneously pass, and will fail the job.
In current code, in this case the currentRetryWindowMs (called lastFailureMsDiff)
gets set = currentTime, which fails the check and falls into the else.

One way to fix is to set lastFailureTime to 0 (as in current code), or set a flag (firstFailure) and use that to make a determination on the if branch (since there is no failure-window for the first fail).

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not obvious, but the retryDelay is always zero on the first failure, so the scenario you mentioned wouldn't happen. I'll wrap this logic in a private method and document the logic to make it more clear. Thanks.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Ah i see,
But is it possible that the retryDelay is > 0 on a subsequent failure, and due to scheduling delays on the machine running this code
Instant.now().toEpochMilli() - lastFailureTime.plus(retryDelay).toEpochMilli(); becomes < 0 because lastFailureTime is set to Instant.now() above.
Could this happen even with retryDelay = 0?
Because if it becomes < 0 then the check below will fail.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

lastFailureTime is only set to Instant.now() on the first failure. I'm I overlooking something?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Youre right, I was confusing the current code, which has
currentFailCount >= retryCount with the currentFailCount > maxRetryCount
the equality could lead to a erroneous situation but with the change it shouldnt happen

@rmatharu-zz rmatharu-zz left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Took another pass, one corner case needs a fix, other are mostly minor.

@dnishimura dnishimura left a comment

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks @rmatharu for the review. Please see my responses. New commit will follow later.

currentFailCount = 1;
lastFailureTime = Instant.now();
}
if (currentFailCount > maxRetryCount) {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sure will add.

}
if (currentFailCount > maxRetryCount) {
Duration retryDelay = getHostRetryDelay(lastSeenOn, currentFailCount - 1);
long currentRetryWindowMs = Instant.now().toEpochMilli() - lastFailureTime.plus(retryDelay).toEpochMilli();

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not obvious, but the retryDelay is always zero on the first failure, so the scenario you mentioned wouldn't happen. I'll wrap this logic in a private method and document the logic to make it more clear. Thanks.

public static final int DEFAULT_CONTAINER_RETRY_COUNT = 8;

public static final String CLUSTER_MANAGER_CONTAINER_RETRY_MAX_DELAY_MS = "cluster-manager.container.host-affinity-retry.max.delay.ms";
public static final long DEFAULT_CONTAINER_RETRY_MAX_DELAY_MS = Duration.ofSeconds(120).toMillis();

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sure will add

lastFailureTime = Instant.now();
}
if (currentFailCount > maxRetryCount) {
Duration retryDelay = getHostRetryDelay(lastSeenOn, currentFailCount - 1);

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not obvious, but the retryDelay is always zero on the first failure, so the scenario you mentioned wouldn't happen. I'll wrap this logic in a private method and document the logic to make it more clear. Thanks.

@rmatharu-zz rmatharu-zz left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

clarification

@rmatharu-zz rmatharu-zz left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

thanks

@rmatharu-zz rmatharu-zz left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Looks like we have a bug where the retry-window doesnt get recomputed upon each failure.

lastFailureTime = 0L;
}
if (currentFailCount >= retryCount) {
long lastFailureMsDiff = System.currentTimeMillis() - lastFailureTime;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Looks like we have a bug where the retry-window doesnt get recomputed upon each failure.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yup will fix his existing bug in this PR since it's related. Thanks for pointing it out.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Summarizing the discussion, for anyone looking at this review,
there was a bug where the "fail job if N failures in a M window" was only looking at the "last failure being in the M window"

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in #1108

@dnishimura

Copy link
Copy Markdown
Contributor Author

PR on hold due to a bug that @rmatharu and I discovered. Will work on a PR for that bug and merge that in here afterwards once that bug fix is in master. Details of the bug here SAMZA-2277

@dnishimura

Copy link
Copy Markdown
Contributor Author

@rmatharu please take a look at PR #1108 for the fix with the retry window. I'll merge master here once that PR is committed to master. Thanks.

asfgit pushed a commit that referenced this pull request Jul 29, 2019
…ot reflected in code

In the current code, the window is only applied and checked on the last retry. However, the check should be done at all retries.

This was found during #1104

rmatharu please take a look

Author: Daniel Nishimura <dnishimura@gmail.com>

Reviewers: Ray Matharu <rmatharu@linkedin.com>

Closes #1108 from dnishimura/samza-2277-cluster-manager-retry-window-bug-fix and squashes the following commits:

23db835 [Daniel Nishimura] Address minor comments from @rmatharu
0a7c927 [Daniel Nishimura] Trigger build b/c of flakey test.
0c56499 [Daniel Nishimura] Fix checkstyle
2674f75 [Daniel Nishimura] Address @rmatharu's comments.
ecfe63e [Daniel Nishimura] SAMZA-2277: Semantics for cluster-manager.container.retry.window.ms not reflected in code
@dnishimura

Copy link
Copy Markdown
Contributor Author

@rmatharu please take a look at the latest changes. I merged in master to pull in the changes from PR #1108 that has the fix for the retry window. The rest of the commits where to clean up some things and to address your comments.
I'll be testing on an actual job, so please hold off from merging until I'm done, but feel free to review the code in the meantime. Thanks!

@dnishimura

Copy link
Copy Markdown
Contributor Author

@rmatharu I tested on an actual job and it seems to work well. Per our discussion, it's no longer doing an exponential backoff, but instead it adds a delay at the last retry for the container restart. The last delayed retry honors the cluster-manager.container.retry.window.ms by not including the delay as part of the window. Please take a look when you get a chance.

Comment thread docs/learn/documentation/versioned/jobs/samza-configurations.md Outdated
if (processorFailures.containsKey(processorId)) {
ProcessorFailure failure = processorFailures.get(processorId);
currentFailCount = failure.getCount() + 1;
Duration lastRetryDelay =

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Nit: Could this Duration lastRetryDelay = processorFailures.containsKey(processorId) ? processorFailures.get(processorId).getLastRetryDelay() : Duration.ZERO; be extracted into a method, since it is used below on line 574

}

public Long getLastFailure() {
public Instant getLastFailure() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

+1, thanks for this

* Sends the {@link SamzaResourceRequest}s in the delayed requests queue that have expired.
* @return number of delayed requests sent.
*/
public int sendExpiredDelayedResourceRequests() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Could we use another name, e.g., sendPendingDelayedRequests
In this code "expired resource requests" is also used to refer to resource requests which were never answered with a resource by the RM, so could confuse a new dev looking at this.

private final PriorityQueue<SamzaResourceRequest> requestsQueue = new PriorityQueue<>();

/**
* Represents the queue of delayed resource requests made by the {@link ContainerProcessManager}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

maybe add to documentation: "The difference between requestsQueue and delayedRequestsQueue is that, any request present the requestsQueue has been sent out to the YARN-RM,
while requests in the delayedRequestsQueue will be sent to the YARN-RM only when their delay reaches 0."

synchronized (lock) {
int numReleasedResources = 0;
if (requestsQueue.isEmpty()) {
if (requestsQueue.isEmpty() && delayedRequestsQueue.isEmpty()) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We don't need to check the delayedRequestsQueue here, because a resource-request for anything in the delayedRequestsQueue will be sent out only when the requests are no-longer delayed.

The implication of this is that if requestsQueue is empty and delayedRequestsQueue is not, the CPM will continue to hold onto allocated resources, in the worst case for a period of 5 mins (default).
Aggregated over num_containers, could be significant especially in resource-crunch scenarios.

Alternatively, we could release the resource only looking at the requestsQueue, and when a request in the delayedRequestsQueue "expires", we will send out the request to the YARN-RM and allocation flow shall resume.

@rmatharu-zz rmatharu-zz left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A couple of things, others are mostly minor

@dnishimura

Copy link
Copy Markdown
Contributor Author

@rmatharu - I addressed your latest comments and the build is green. Please merge after your approve. Thanks for reviewing this PR.

@rmatharu-zz

Copy link
Copy Markdown
Contributor

thanks, will merge

@dnishimura

Copy link
Copy Markdown
Contributor Author

thanks, will merge

Great! Thanks again for reviewing.

@rmatharu-zz
rmatharu-zz merged commit 2e17e08 into apache:master Aug 14, 2019
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants