Not a member of Pastebin yet?
Sign Up,
it unlocks many cool features!
- our team is building a Job Orchestration Platform similar to Temporal, Airflow,
- or AWS Step Functions. The platform is responsible for executing long-running jobs composed of multiple tasks.
- Functional Requirements
- A client should be able to submit a job.
- Each job consists of one or more tasks.
- Tasks can have dependencies, forming a Directed Acyclic Graph (DAG).
- A task can only execute after all of its dependencies have completed successfully.
- The system should support different task types (e.g., API call, database operation, shell script, Spark job).
- If a task fails, it should be retried based on a configurable retry policy.
- Users should be able to:
- Pause a job
- Resume a job
- Cancel a job
- Users should be able to query the current status of a job and its individual tasks.
- Once all tasks complete successfully, the job should be marked as completed.
- // There will be a period where the job is pausing and trying to finish / pause the current running tasks
- // THE CODE I WROTE DURING THE INTERVIEW BEGINS BELOW
- enum RetryFrequency {
- NONE,
- LINEAR,
- EXPONENTIAL_BACKOFF
- }
- abstract class RetryPolicy {
- int retryAfterDuration, int maxNumberOfRetries;
- RetryPolicy retryPolicy = NONE;
- void retry();
- }
- class NoRetry extends RetryPolicy {
- NoRetry() {
- retryAfterDuration = 0;
- maxNumberOfRetries = 0;
- }
- }
- class ConfigurableRetryPolicy extends RetryPolicy {
- ConfigurableRetry(int retryAfterDuration, int maxNumberOfRetries, RetryFrequency retryFrequency) {
- // in actual implementation, add validation that both these values should be >= 0
- super.retryAfterDuration = retryAfterDuration;
- super.maxNumberOfRetries = maxNumberOfRetries;
- super.retryFrequency = retryFrequency;
- }
- }
- enum RetryConfig {
- NONE,
- CONFIG(int retryAfterDuration, int maxNumberOfRetries, RetryFrequency retryFrequency);
- }
- class RetryPolicyFactory {
- public RetryPolicy getRetryPolicy(RetryConfig config) {
- switch (config) {
- case NONE: return new NoRetry();
- case CONFIG(int retryAfterDuration, int maxNumberOfRetries, RetryFrequency retryFrequency):
- return new ConfigurableRetryPolicy(retryAfterDuration, maxNumberOfRetries, retryFrequency);
- }
- }
- }
- enum TaskStatus {
- READY,
- IN_PROGRESS,
- CANCELLED,
- COMPLETED
- }
- enum JobStatus {
- READY,
- IN_PROGRESS,
- PAUSED,
- CANCELLED,
- COMPLETED
- }
- class Task {
- String id;
- String name;
- Callable task;
- TaskStatus status; // maintained internally
- RetryConfig retryConfig;
- }
- class Job {
- String id;
- String name;
- Map<Task, List<Task>> tasks; // for DAG, the key will be the task that needs to be completed first, and the value will be
- // the list of dependent tasks
- JobStatus status;
- }
- class Status {
- JobStatus jobStatus;
- Map<String, taskStatus> tasksStatusMap; // the key = taskId, the value = TaskStatus
- }
- interface JobOrchestrator {
- void submit(Job job);
- void pause(String jobId);
- void resume(String jobId);
- void cancel(String jobId);
- Status getJobStatus(String jobId);
- TaskStatus getTaskStatus(String taskStatus);
- }
Advertisement
Add Comment
Please, Sign In to add comment