The dask-yarn documentation has the following example.
from dask_yarn import YarnCluster
from dask.distributed import Client
# Create a cluster where each worker has two cores and eight GiB of memory
cluster = YarnCluster(environment='environment.tar.gz',
worker_vcores=2,
worker_memory="8GiB")
# Scale out to ten such workers
cluster.scale(10)
# Connect to the cluster
client = Client(cluster)
In the example, what is the definition of "worker"? Is it a node (i.e., piece of hardware)? Is it a process? What do I pass to YarnCluster and cluster.scale?
To be more concrete, say I have a cluster with 20 nodes, each of which has 8 CPUs and 32 GB of RAM, and I want to maximize resource usage. Do I have to set worker_vcores = 8, worker_memory = 32, and cluster.scale(20)? Or do different settings work; e.g., worker_vcores = 4, worker_memory = 16, and cluster.scale(40)?