Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

[query] Improve error handling and job tracking in ServiceBackend #14751

Open
wants to merge 1 commit into
base: qob-fast-cancel
Choose a base branch
from

Conversation

grohli
Copy link
Contributor

@grohli grohli commented Nov 5, 2024

Change Description

  • Added ability to check job status in individual jobs in a JobGroup
  • Cancel JobGroup if any jobs in the partition fail
  • Added test functionality for detecting cancelled, failed jobs
  • Query batch for which jobs have failed within a JobGroup, rather than go through every job in the group.

Security Assessment

Delete all except the correct answer:

  • This change has no security impact

Impact Description

Mainly just code for testing, nothing security related

(Reviewers: please confirm the security impact before approving)

Copy link
Contributor Author

grohli commented Nov 5, 2024

Warning

This pull request is not mergeable via GitHub because a downstack PR is open. Once all requirements are satisfied, merge this PR as a stack on Graphite.
Learn more

This stack of pull requests is managed by Graphite. Learn more about stacking.

@grohli grohli changed the title Added ability to check job status of individual job, made associated tests [batch] Improve error handling and job tracking in ServiceBackend Nov 5, 2024
@grohli grohli changed the title [batch] Improve error handling and job tracking in ServiceBackend [query] Improve error handling and job tracking in ServiceBackend Nov 5, 2024
@grohli grohli marked this pull request as ready for review November 5, 2024 19:30
Comment on lines -315 to +350
new CancellationException(
s"Job group ${jobGroup.job_group_id} for batch ${batchConfig.batchId} was cancelled"
)
}

r = (Some(error), results)
}
def streamSuccessfulJobResults: Stream[(Array[Byte], Int)] =
for {
successes <- batchClient.getJobGroupJobs(
jobGroup.batch_id,
jobGroup.job_group_id,
Some(JobStates.Success),
)
job <- successes
partIdx = job.job_id - startJobId
} yield (readPartitionResult(root, partIdx), partIdx)

val r @ (_, results) =
jobGroup.state match {
case Success =>
runAllKeepFirstError(executor) {
(partIdxs, parts.indices).zipped.map { (partIdx, jobIndex) =>
(() => readPartitionResult(root, jobIndex), partIdx)
}
}
case Failure =>
val failedEntries = batchClient.getJobGroupJobs(
jobGroup.batch_id,
jobGroup.job_group_id,
Some(JobStates.Failed),
)
assert(
failedEntries.nonEmpty,
s"Job group ${jobGroup.job_group_id} failed, but no failed jobs found.",
)
val error = readPartitionError(root, failedEntries.head.head.job_id - startJobId)

(Some(error), streamSuccessfulJobResults.toIndexedSeq)
case Cancelled =>
val error =
new CancellationException(s"Job Group ${jobGroup.job_group_id} was cancelled.")
Copy link
Member

Choose a reason for hiding this comment

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

Revert this - there was more context in the last message

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