In this RX-M Cloud Native Short Take, Instructor Christian Lacsina discusses the RX-M Druid Architecture Module. Chrisitan explores each of the components of a Druid deployment, demonstrating what occurs in a Druid deployment when a task is submitted, exploring the interaction between Druid and Zookeeper, and discussing the infrastructure considerations that affect Druid. Additionally, we will explore how Druid scales.
Video Transcript
Welcome to another RX-M Cloud Native Short Take my name is Christian Lacsina and today we short take the RX-M Apache Druid Architecture module. In this module we discuss each of the constituent components of a Apache Druid deployment. We also demonstrate what occurs in a Druid deployment when a task is submitted. We explore the interaction between Druid and Apache Zookeeper and talk about the infrastructure considerations that affect Druid. Finally, in this module we explore how Druid scales.
We walk through Apache Druid’s architecture by launching containers for each of Druid’s components and examine what functionality we receive with each new container, building up to a complete Druid cluster. We start by looking at the groundwork for our Druid deployment: its external dependencies. Druid depends on Apache Zookeeper for cluster component coordination and relies on some kind of relational database–like PostgreSQL in our case–to store metadata on its various jobs and other identifiers. In the demo, we run these dependencies in Docker containers. Lastly Druid relies on some kind of file storage which is usually an object store like S3, MinIO, or OpenStack Swift; however for the demo we use local disk.
With the external dependencies covered we need to configure the Druid cluster; many Druid configurations can be provided using environment variables so for the demo we use a file full of them to configure the cluster. Environment variables make Druid very easy to configure for containerized environments like that in our demo.
The first druid component we bring up is the coordinator; a Druid coordinator is responsible for ensuring appropriate data distribution to each of the data processes that are part of a Druid cluster. Additionally a task distributor called the “Overlord” runs alongside the coordinator to trigger data ingestion movement and management tasks across the cluster. It also provides a GUI!
Data in druid is handled by a couple of other components. One of these components is the “Middle Manager”; the Middle Manager is responsible for creating Druid-formatted events by parsing live data streams like those provided by an Apache Kafka clusteror parsing files in batches which can be from local disk or even an Apache Hadoop integration. Once Middle Manager is running–in a container in the demo–we can load data into the cluster–even though other Druid components are not yet up and running! To load data into the cluster a task specification has to be submitted to the coordinator. In the demo this is done by providing a json file that describes where Druid needs to retrieve the data and how events produced from those data should look like. This file is called an “ingestion spec”. The ingestion spec is sent to the coordinator’s Overlord port to initiate a data ingestion task. This task reads a short json file that has been provided to the containers and produces queryable Druid events. During the ingestion task the Middle Manager parses data from the specified source, converts that data into events stored in a columnar format and then produces files called segments which contain parts of that data. The segments are separated into time-based partitions called chunks. Those segments are then sent to the file storage known as the “Deep Store” which, in the demo, is running on the local disk.
Despite the fact that the Middle Manager has these data they’re not queryable yet! There isn’t any querying functionality yet and the batch data that was fed into the cluster aren’t actually queryable from the Middle Managers. The Middle Managers at will however make the streaming data from a Kafka cluster, for example, available immediately for querying.
By checking the log we can see the coordinator is saying that it has no servers that can actually serve these batch data for querying. Now batch data like what has been ingested in the demo is only made queryable once Druid takes these segments and makes them available through a Historical process. Historical processes are responsible for holding batch data for user queries. After launching a historical process the coordinator should find a place to store that data and we see in the logs that a new Historical process has been detected. Data handling capabilities with Apache Druid are handled exclusively by the Middle Manager and the Historical processes. In order to increase or otherwise scale the data handling capacity of Druid we just need to add additional Middle Managers or Historical processes to the cluster.
Once the cluster has successfully recorded the data it can be queried. Query functionality is provided by the “Broker process”. Brokers parse user queries which will then reach out to the Middle Managers and the Historical processes to collate results based on a user’s request. Multiple brokers can handle user queries thanks to a process called a “Router”. A Router can effectively load balance queries between different Brokers and it can also perform different kinds of routing for Brokers. So let’s say you had a Broker that had more data that was live versus other Broker(s) that would almost exclusively deal in historical data. You could do that with a Druid Router. In addition to Broker handling Routers can provide additional UIs for data handling. So if we head over to our Router port we have a new GUI which is slightly more expanded than the previous coordinator GUI. From the Router GUI we can use the “query” option and look for data based on the Druid version of sql called, you guessed it: Druid SQL, which is actually not quite SQL however it gives a sequel like interface for those of you who are more suited to using SQL in your daily lives. So by adding additional Brokers and Routers to a cluster the query capacity can be scaled without affecting data storage capabilities.
And that’s it! You now see how Apache Druid’s deployment architecture can suit a cloud native deployment. Discovering what makes a Druid cluster tick is just one of the many things you will learn in the Druid Architecture module and the druid classes from RX-M. We have many other classes covering databases which you can browse at our course catalog. Under the “Training” menu, click on “Course Catalog” and if you scroll down you’ll see the various categories of courses that we have available. For example, under “Database Courses” you can see that we have Apache Druid among many other offerings for database-focused classes in addition to data science and other types of courses.
That’s our cloud native short take on the RX-M Apache Druid Architecture module!