DhruvSaraswat

Hevo_Data_Round_2_Interview_LLD_6th_Oct_2026

Oct 6th, 2026
11
0
Never
Not a member of Pastebin yet? Sign Up, it unlocks many cool features!
text 3.25 KB | Source Code | 0 0
  1. our team is building a Job Orchestration Platform similar to Temporal, Airflow,
  2. or AWS Step Functions. The platform is responsible for executing long-running jobs composed of multiple tasks.
  3. Functional Requirements
  4.  
  5. A client should be able to submit a job.
  6. Each job consists of one or more tasks.
  7. Tasks can have dependencies, forming a Directed Acyclic Graph (DAG).
  8. A task can only execute after all of its dependencies have completed successfully.
  9. The system should support different task types (e.g., API call, database operation, shell script, Spark job).
  10. If a task fails, it should be retried based on a configurable retry policy.
  11. Users should be able to:
  12. Pause a job
  13. Resume a job
  14. Cancel a job
  15.  
  16. Users should be able to query the current status of a job and its individual tasks.
  17. Once all tasks complete successfully, the job should be marked as completed.
  18.  
  19.  
  20. // There will be a period where the job is pausing and trying to finish / pause the current running tasks
  21.  
  22. // THE CODE I WROTE DURING THE INTERVIEW BEGINS BELOW
  23. enum RetryFrequency {
  24. NONE,
  25. LINEAR,
  26. EXPONENTIAL_BACKOFF
  27. }
  28.  
  29. abstract class RetryPolicy {
  30. int retryAfterDuration, int maxNumberOfRetries;
  31. RetryPolicy retryPolicy = NONE;
  32.  
  33. void retry();
  34. }
  35.  
  36. class NoRetry extends RetryPolicy {
  37. NoRetry() {
  38. retryAfterDuration = 0;
  39. maxNumberOfRetries = 0;
  40. }
  41. }
  42.  
  43. class ConfigurableRetryPolicy extends RetryPolicy {
  44. ConfigurableRetry(int retryAfterDuration, int maxNumberOfRetries, RetryFrequency retryFrequency) {
  45. // in actual implementation, add validation that both these values should be >= 0
  46. super.retryAfterDuration = retryAfterDuration;
  47. super.maxNumberOfRetries = maxNumberOfRetries;
  48. super.retryFrequency = retryFrequency;
  49. }
  50. }
  51.  
  52.  
  53. enum RetryConfig {
  54. NONE,
  55. CONFIG(int retryAfterDuration, int maxNumberOfRetries, RetryFrequency retryFrequency);
  56. }
  57.  
  58.  
  59. class RetryPolicyFactory {
  60. public RetryPolicy getRetryPolicy(RetryConfig config) {
  61. switch (config) {
  62. case NONE: return new NoRetry();
  63. case CONFIG(int retryAfterDuration, int maxNumberOfRetries, RetryFrequency retryFrequency):
  64. return new ConfigurableRetryPolicy(retryAfterDuration, maxNumberOfRetries, retryFrequency);
  65. }
  66. }
  67. }
  68.  
  69.  
  70. enum TaskStatus {
  71. READY,
  72. IN_PROGRESS,
  73. CANCELLED,
  74. COMPLETED
  75. }
  76.  
  77. enum JobStatus {
  78. READY,
  79. IN_PROGRESS,
  80. PAUSED,
  81. CANCELLED,
  82. COMPLETED
  83. }
  84.  
  85. class Task {
  86. String id;
  87. String name;
  88. Callable task;
  89. TaskStatus status; // maintained internally
  90. RetryConfig retryConfig;
  91. }
  92.  
  93. class Job {
  94. String id;
  95. String name;
  96. Map<Task, List<Task>> tasks; // for DAG, the key will be the task that needs to be completed first, and the value will be
  97. // the list of dependent tasks
  98. JobStatus status;
  99. }
  100.  
  101. class Status {
  102. JobStatus jobStatus;
  103. Map<String, taskStatus> tasksStatusMap; // the key = taskId, the value = TaskStatus
  104. }
  105.  
  106. interface JobOrchestrator {
  107. void submit(Job job);
  108. void pause(String jobId);
  109. void resume(String jobId);
  110. void cancel(String jobId);
  111. Status getJobStatus(String jobId);
  112. TaskStatus getTaskStatus(String taskStatus);
  113. }
Advertisement
Add Comment
Please, Sign In to add comment