KubernetesTaskRunner for running druid tasks as kubernetes jobs
- Dominant language
- Java
- Stars
- 14.1k
- Forks
- 3.8k
- Avg merge
- 2d 58m
- Merged PRs (30d)
- 233
Description
### Motivation
Below talk by @Jinchul81 outlines a way to autoscale druid Middlemanagers by emitting druid metrics -
https://www.slideshare.net/Hadoop_Summit/apache-druid-auto-scaleoutin-for-streaming-data-ingestion-on-kubernetes
However, there are some limitations with this approach especially for selecting MMs for scaling down discussed towards end slides.
* A middlemanager can only be scaled down when all its tasks have been completed
* Selecting a specific MM to scale down is not trivial
Another requirement with MMs is over-provisioning where `required workerCapacity = 2 * replicas * taskCount` to accomodate one set of tasks publishing while another set is reading
### Proposed changes
This proposal is to implement a KubernetesTaskRunner as part of a new extension druid-kubernetes
* K8sTaskRunner will use K8s API to directly submit jobs to kubernetes cluster using java kubernetes-client https://github.com/kubernetes-client/java
* Each K8s Job will only run Peon process.
* Kubernetes task runner will use Watch ([link](https://github.com/kubernetes-client/java/blob/master/examples/src/main/java/io/kubernetes/client/examples/WatchExample.java)) to watch for the status changes of the submitted jobs
* Task logs can will also be streamed in the console using existing kubernetes APIs. ([link](https://github.com/kubernetes-client/java/blob/master/examples/src/main/java/io/kubernetes/client/examples/LogsExample.java))
* For replica tasks, kubernetes antiaffinity can be used to make sure that replica tasks are assigned to different instances
* Once the job completes, allocated resources will be freed and can be assigned to other tasks.
* Instead of MM pushing the task logs on completion to log storage, Overlord will fetch the logs and push them to the log storage
### Rationale
* KubernetesTaskRunner would help in better resource utilization in the cloud as we need not allocate larger pods which can host multiple tasks that are mostly under-utilized
* Kubernetes [ClusterAutoscaler](https://github.com/kubernetes/autoscaler/tree/master/cluster-autoscaler) can be used to add more nodes whenever needed, If enough resources are not available, tasks will remain in pending state, cluster-autoscaler will create more nodes for running these tasks
* No over-provisioning is needed to allow simultaneously publishing and reading tasks is required
### Operational impact
* Simplified autoscaling in cloud.
* No need to manage extra configuration for MiddleManagers in the cloud
* This is proposed to be done as a kubernetes extension with the hope of adding better cloud support for druid. It will be an Optional opt-in feature and no change in core druid are required.
Contributor guide
Research direction
The proposed entry point is a new KubernetesTaskRunner in a druid-kubernetes extension; start with the linked Kubernetes Java client examples for job submission, Watch status handling, and log streaming. Done would mean the extension can submit and monitor one Peon job per task, stream logs, and release resources after completion.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- java, kubernetes
- Domain
- cloud, infrastructure
- Issue type
- Feature
- Difficulty
- 5/5
- Estimated time
- Over a week
- Activity status
- Stale
- Clarity
- Needs clarification
- Newbie friendliness
- 25/100