Welcome to another article in Batch Processing. We have covered Object Stores and Distributed File Systems in our previous article. It is recommended to read that before this one.
In this read, we will cover distributed job orchestration and a couple of Batch Processing frameworks.
Distributed Job Orchestration
Fortunately, my current job involves ownership of a huge amount of distributed job orchestration. Trust me, there are a lot of problems that occur the moment distributed systems come into the picture. We cannot cover everything and every problem that we face in such scenarios, but it is always possible to have fundamental tools and basic concepts in hand that form the basis of almost all problem solutions in distributed systems architecture.
Whenever we run a program dealing with multiple steps in action on a single host, for example, a single Python program or a Linux command containing multiple commands (like we studied in the last article), programs like aws, sort, etc. in the Linux command we dealt with in the last article are orchestrated by the kernel of the OS. Here, the process aws was allocated CPU, then the data and related information was transmitted to sort, and so on. This is basically orchestration. In the case of distributed jobs where multiple hosts are involved, a job orchestrator does that.
So if we talk about the flow, a job orchestrator scheduler receives a request to run a job. The request is sent by any batch-processing framework (which we will see later). A request can contain information like the number of tasks to execute, memory, CPU, disk needed, secrets and credentials, and other job parameters, etc.
Popular executors are Kubernetes and Hadoop YARN. They execute the job as per the request.
There are 3 common components inside an executor:
Task executors: They execute the tasks on nodes and also send heartbeats to signal their health. Task executors are responsible for tracking the task status and updating them as needed, as well as allocating resources to the task. In short, they execute the command in the task. Sometimes, they are also responsible for providing security and performance isolation. For example, Kubernetes uses Linux cgroups. This ensures a task is able to access only the respective data and no other task can access it.
Resource Manager: It stores data about the resources or nodes, CPU, memory, disk, GPU, etc. Thus, it globally provides a view of the cluster's current state. Popular examples are ZooKeeper and etcd in Kubernetes. The resource manager is centralized, so it is difficult to scale and availability is also difficult.
Scheduler: This is also a centralized subsystem. It receives requests to start, stop, or get the status of a job. This is the one which uses the request to get the state of the resource manager and then determines which task to run on which nodes. It then informs the task executors to begin the execution.
How is Resource Allocated for Executors?
Just as CPU scheduling is very challenging for dispatchers, a similar challenge is faced by a scheduler when allocating resources to competing jobs.
Imagine a scenario where a scheduler has 2 waiting jobs to be scheduled, each wanting 100 cores of CPU spread across 5 nodes. The total available cores in the cluster are 160. The following are different ways to schedule.
- The scheduler could decide to run 80 tasks for each job and then queue the rest.
- One job is scheduled after the other completes. Here, a node can stay idle, which is not efficient. On the other hand, starvation can also happen here.
- Reserve some extra cores for jobs in the future.
There can be many different ways.
Another way can be preempting a job to give room for another important task. But preempting decreases cluster efficiency as preempted/killed tasks need to be restarted again.
Scheduling is often an NP-Hard problem. So, mostly, this is done using a heuristic approach. Keep it non-optimal but reasonable. Several other algorithms similar to what we deal with in CPU scheduling, like FIFO, priority queue, and Round Robin, are also used.
In my current organisation, I recently designed a smart storage processor scheduler that considerably improved our cloud costing. Maybe we can discuss that in detail in the future.
What is a Workflow?
Workflows are Directed Acyclic Graphs (DAGs) of jobs. In other words, for distributed batch processes, the output of one job often becomes the input for another job. A job can also have multiple inputs, each coming from different jobs. The same thing in Linux can be assumed as a pipe being used to transfer the data output of one process into another.
There can be several reasons why a workflow is needed. Maybe different teams are working on different jobs of the workflow, maybe batch processing of the output of several jobs is needed in a final job, or maybe the data pipeline itself requires stages to function.
How is one job in a workflow decoupled from another?
A most common way is using a DFS (Distributed File System) or an object store for storing the data to be shared from one job to another. A dependent job waits until the other jobs whose data it is going to consume as input have completed before it runs.
Fault Tolerance
How are failures handled? Before this question, we should think about how failures occur. There can be many ways, but quite possible ways are simply network failures, hardware faults, and even software faults sometimes. Preemptions of low-priority tasks can also cause faults. Imagine a low-priority task never being able to finish due to frequent preemptions. Such tasks can be scheduled on spot instances, a jargon term.
Fault tolerance becomes harder when one task’s output is used as another task’s input.
MapReduce handles this by:
- Saving intermediate results to the distributed filesystem.
- Waiting until the write is successfully completed.
- Only then allowing other tasks to read those results.
This makes the system reliable even when tasks can be interrupted or preempted. However, it can be inefficient because intermediate data is repeatedly written to the distributed filesystem, causing a lot of disk I/O.
We will discuss MapReduce in detail.
Similarly, Spark keeps intermediate data in memory, “spilling” it to the local disk if it won’t fit, and writes only the final result to the DFS. It also keeps track of how the intermediate data was computed, allowing Spark to recompute it in case it is lost.
Conclusion
In the last article as well as this one, we understood what batch processing is, how it can be useful, and the science behind distributed processing jobs and workflows. We have ramped up the gear to now jump into understanding the popular batch processing models that have formed the basis of modern data-driven technologies and revolution, be it AI or what not.
In the next article, we will cover MapReduce and other common models.
Here is the link to the last article if you have missed it.
https://dev.to/ujjwall-r/batch-processing-from-unix-tools-to-distributed-systems-dbh
Stay tuned for another weekend.
Not an interesting fact: The cover image is sketched by me on the bazaars of Hyderabad. The image uses one-point perspective and is inspired by the streets of Hyderabad. The classiness of the Charminar is adorable. Do visit the city once you get a chance.

Top comments (0)