CDH 5.14.4集群中Dask-Yarn无法分配预期Worker数量的问题咨询
Let's break down your problems step by step, starting with the core issue of not being able to spin up more than 18 workers, then moving to the adaptive cluster question.
Problem 1: Can't Allocate More Than 18 Workers
First, let's recap your setup: 6 nodes (384 total vcores), requesting 24 workers with 4 vcores and 64GB memory each. Even bumping n_workers to 50 doesn't help, but adjusting worker resources (vcores/memory) works. Here are the most likely causes and fixes:
1. Check YARN Queue Resource Quotas
This is the most common culprit in managed Hadoop clusters like CDH. Even if the cluster has total resources available, your assigned queue might have hard limits on how much vcore/memory it can use.
- Go to the YARN Web UI (typically
http://<resourcemanager-host>:8088), navigate to Cluster > Scheduler to find your application's queue. - Look at the queue's Used Resources vs Max Resources. For your 24-worker request, you'd need
24*4 + 2 = 98 vcores(plus 2 for scheduler/logger) and24*64 = 1536 GBof memory. If your queue's max vcores or memory is lower than these numbers, that's the block. - If you have admin access, check queue configs (e.g., for Capacity Scheduler):
yarn.scheduler.capacity.<your-queue>.maximum-allocation-vcores: Max vcores per container (should be at least 4 in your case)yarn.scheduler.capacity.<your-queue>.maximum-allocation-mb: Max memory per container (should be at least 65536 MB)yarn.scheduler.capacity.<your-queue>.maximum-capacity: The upper limit of cluster resources this queue can consume (needs to cover your total request)
2. Verify Node Resource Availability
Even if the queue has enough quota, individual nodes might be busy or have reserved resources that block worker deployment:
- In the YARN UI's Nodes page, check each node's Available Resources. If some nodes are already using most of their vcores/memory for other jobs (like Hive queries or Spark tasks), Dask can't place workers there.
- Check YARN NodeManager configs (via Cloudera Manager if you're using CDH):
yarn.nodemanager.resource.cpu-vcores: The total vcores YARN can allocate per node (often less than physical cores, reserved for system processes)yarn.nodemanager.resource.memory-mb: Total memory YARN can allocate per node (again, system memory is reserved here)
For your 64GB worker memory, a node with ~300GB total allocatable memory could run 4 workers max. If 3 of your nodes are partially occupied, you'd only get 34 + 32 = 18 workers (which matches your symptom).
3. Check Worker Startup Logs for Failures
Sometimes it looks like workers aren't being allocated, but they're actually starting and crashing immediately.
- In the YARN UI, find your Dask application and click Logs. Check the worker logs (not just scheduler/logger) for errors like missing dependencies from your
env.tar.gz, permission issues, or memory errors. - If your environment package is missing critical libraries (like dask, skein, or pandas), workers will fail to initialize and YARN will clean up the container, making it seem like no allocation happened.
4. Confirm Queue Assignment
If you didn't specify the queue parameter in YarnCluster(), Dask might be using the default queue, which often has strict resource limits. Try explicitly setting your queue:
cluster = YarnCluster( environment='path/to/my/env.tar.gz', n_workers=24, worker_vcores=4, worker_memory='64GB', queue='your-priority-queue' # Replace with your actual queue name )
Problem 2: Adaptive Cluster Shows Only 2 Containers for 10 Workers
You mentioned using cluster.adapt() with 10 workers (10 threads, 100GB memory each), but YARN UI only shows 2 containers. Here's what to check:
1. Adaptive Cluster Isn't Triggering Expansion
Dask's adaptive mode only spins up workers when there's pending work. If your ETL job isn't generating enough tasks to saturate the 2 workers, the adaptive controller won't scale up.
- Check the Dask Dashboard's Tasks tab: if there are no pending tasks (all tasks are either running or completed), the cluster will stay at the minimum worker count (default is 2 if you didn't set
minimum). - Try setting explicit bounds for adaptive mode to test:
cluster.adapt(minimum=5, maximum=10)
Then submit a large enough workload to trigger scaling.
2. Worker Startup Failures
Again, check the YARN application logs. If workers are failing to start, the adaptive controller will keep trying but only maintain the number of successfully running workers (in your case, 2). Common issues here include:
- 100GB worker memory exceeding YARN's per-container memory limit (check your queue's
maximum-allocation-mbconfig) - Missing dependencies in your environment package that are required for the heavier workload
- Permissions issues accessing data sources or writing outputs
3. Misunderstanding Dask Worker vs YARN Container Mapping
You noted that you thought YARN containers = Dask workers, which is correct for dask-yarn (each worker runs in its own container). So if you see 2 containers in YARN UI, that means only 2 workers are actually running (plus scheduler/logger, which should show as separate containers). The "10 workers" you're seeing in Dask Dashboard might be a misinterpretation—double-check the Workers tab to confirm how many are actually active.
内容的提问来源于stack exchange,提问作者skibee

