1 of 14

Apache Drill Resource Management Overview

  • Gautam Parai
  • Hanumath Maduri
  • Karthikeyan Manivannan
  • Sorabh Hamirwasia

2 of 14

Agenda

  • Current State of Resource Management
  • Proposed Enhancements
  • Discussions on proposed design

3 of 14

Resource Management - Current State

  • Assumes entire cluster is available for each query
    • Doesn’t have knowledge of actual cluster state
  • Memory Limit set only for Blocking operators
    • Without queueing: Assumes all buffered operators in query plan are running across max parallelization number of minor fragments on a node. BufferedOpMemory = QueryMemoryOnNode / (count of buffered operators in plan * max width)
    • With queueing: Get the maximum of the buffered operators per node for a query. Uses it to calculate memory for each buffered operator in the query. BufferedOpMemory = (QueryMemory per node in queue / (max (buffered operator per node))
  • Support for only 2 cost based queues
    • Relies on statically configured threshold
    • Multi-tenancy is not supported
  • No memory limit for Non-Blocking & Exchange operators
    • Default limit of 10G
    • Possibility of OOM’s
  • Batch Sizing support across multiple operators
  • Spilling support for blocking operators

4 of 14

Resource Management - Enhancements

  • New configuration file for RM
  • Support for YARN like multiple named resource pools
    • Tree like configuration for Resource Pools
    • Provide multi tenancy capability
    • Multiple selectors for Resource Pool -
      • Tag Based
      • ACL Based and
      • Complex selectors using AND, OR, NOT_EQUAL operations
  • Queue throttling parameters:
    • Max admissible queries
    • Max number of waiting queries
    • Max memory per query per node
    • Timeout for waiting queries
  • Connection/ Session level parameter
    • for tags
    • wait_for_preferred_nodes
  • Queries are planned using current cluster state

5 of 14

Resource Management - High Level Workflow

Zookeeper

Drillbit_1 (Foreman)

1. Core Planning with Parallelization info and memory requirement

2. Select queue for the query

3. Check if query can be admitted in queue or has to wait.

4. If admitted then update global state in Zookeeper

5. Execute the query. Send minor fragments to assigned Drillbit

6. On query completion Foreman will update the global states.

Submit Minor Fragments

Submit Minor Fragments

Read/Update Global Cluster State

Submit Query

Drillbit_2

Drillbit_3

Drill Client

Query Response

6 of 14

Resource Management - Planner Changes

  • Generating Major and Minor fragments for the physical plan
    • Parallelization based on slice target
  • Estimation of the memory for the plan
    • Estimates are based on stats.
    • Query resource requirement calculated using:
      • Operators are divided into categories to handle resource requirements.
      • Sum of resources used by all Major Fragments
      • Resources for each Major Fragment = (Resource for each minor fragment) X Parallelization
      • Resources for each Minor Fragment = Sum of resources across all operators in a Minor Fragment
  • Selection of the queue based on ACL or tags (best fit / random in case of multiple queues)
  • Scale down the memory for buffered operators using max_query_memory_per_node.

7 of 14

Resource Management - Planner Changes ...

  • Memory calculation is done per each Major fragment.
  • Scan, Filter, Sender belong to one Major fragment.
    • Scan Mem calculation - 5 MB
    • Filter Mem calculation - 4 MB
    • Sender Mem calculation - XB * 5 * 1 = 1MB
  • Receiver, HashAgg and Screen belong to root frag.
    • Receiver mem calculation - XB * 5 * 3 = 3MB+ 1MB (outgoing batch)
    • HashAgg mem calculation - 8MB (build table) + 1MB (outgoing batch)
  • Total mem required for the query : 10 + 13 = 23MB.
  • Memory is calculated using avg_row_size * #rows
  • Memory is adjusted based on max_query_memory_per_node.

Cluster Information.

  • Num Nodes : 2
  • BatchSize : 1MB
  • Xchg BatchSize (XB) : 200KB.

Query Information.

  • #Rows = 1000
  • Scan Parallelism = 5

UnionXg

Screen

HashAgg

Filter

Scan

3 2

3 2

3 2

1 0

1 0

19 4

Total Memory per node

8 of 14

Resource Management - Operators

  • Operators honoring memory limit
    • Blocking Operators:
      • Hash Join
      • Hash Agg
      • Sort
    • Non Blocking Operators: Supports Batch Sizing
      • Project, Merge Join, etc
  • Operators not honoring memory limit
    • Exchange Operators - Will be improved as part of RM
    • TopN, Scan readers (like Complex parquet, JSON) - RM doesn’t fix this.

*Assuming operators with batch-sizing and spilling support lives within the assigned memory.

9 of 14

Resource Management - Exchange

  • Enabling batch sizing at senders and receivers
    • Move from batch count (3 outstanding batches) based heuristics to batch size (In KB) based heuristics.
  • Exchange RM parameters
    • Sending Semaphore limit
    • Receiver SoftLimit
    • Exchange Batch Size
  • Memory planning for exchange operators
    • #Sender, #Receiver , ExchangeRM Params-> Exchange Memory Requirement
  • No special handling for Mux Exchange since Planner calculates Exchange requirements only post assignment

10 of 14

Resource Management - Queue Scheduler

  • Queue Scheduler distributed across Drillbits
    • A Drillbit acts as a leader to a queue
    • Determines query admission for its queue
    • Removes dependency from distributed semaphore in ZK.
  • Queue state is local to each leader
  • Cluster state is maintained in Zookeeper
  • RPC protocol between Foreman and Queue Scheduler (Leader)

11 of 14

Leader Election

  • Each Drillbit generates unique id for each configured queue
  • Registers queue id with cluster membership information in Zookeeper (i.e. DrillbitEndpoint)
  • Upon receiving membership information from Zookeeper starts leader election
  • Each Drillbit sorts unique id for each queue locally
  • Choose the first in the sorted list as leader
  • Chosen leader Drillbit of the queue will validate it’s cluster information with Zookeeper knowledge
    • If both are same then update QueueLeader blob with chosen leader for queue
    • Otherwise no update is done and process continues

12 of 14

Resource Management - Cluster State

  • 4 blobs stored in Zookeeper to maintain current state
    • ClusterState
    • QueueLeader
    • QueueToForemanQueriesCount
    • QueueToForemanUsedResources
  • Blobs helps to recover in the following failure scenarios
    • Query failure
    • Foreman failure
    • Leader failure

13 of 14

Resource Management - New WorkFlow

Zookeeper

  • QueueLeader
  • Cluster State
  • QueueToForemanUsedResources
  • QueueToForemanQueriesCount

Non-Foreman Drillbit_3

Leader of queue Q2

Drillbit_2 (Foreman)

2. Core Planning with parallelization info and cost. Select queue Q1

5. If admit then contact zookeeper to update blobs

8. Query W1 is completed

Drillbit_1 - Leader of queue Q1

3.2 Insert W1 in Q1

10.2 Remove W1 from Q1

1. Submit Query W1

3.1 Ask leader if W1 can be admitted ?

4. Replies with admit /

wait / server busy

6. Update blob as transactions

7. Schedule minor fragments

7. Schedule minor fragments

9. Update blob as transactions

10.1 Query W1 completed

W1

11. Response to Client

Drill Client

Slot for Waiting Queries

Slot for Running Queries

Non-Foreman Drillbit_4

14 of 14

Q&A