# Scaling XGBoost with Dask and Coiled

This notebook shows you how to solve the common **MemoryError** issue that is thrown whenever you try to train an XGBoost model that doesn't fit into your memory. 

You'll learn how to leverage **distributed [XGBoost](https://xgboost.readthedocs.io/en/latest/) training** for effective modelling on datasets that exceed the hardware limitations of your local machine.

Specifically, you will write code to:
1. Train a distributed XGBoost model locally on a small dataset using [Dask](https://dask.org/), 
2. Scale your distributed XGBoost model to the cloud using Dask and [Coiled](https://coiled.io/) to train on a larger-than-memory dataset,
3. Speed up your training with Pro tips from the Dask core team.

### About the Dataset
We'll be using **a 100GB dataset** containing synthetic data, generated using the `dask-ml make_regression` API. The dataset is stored in the public `coiled-datasets` S3 bucket.

In [1]:
import warnings
warnings.filterwarnings('ignore')

import logging
logger = logging.getLogger("distributed.utils_perf")
logger.setLevel(logging.ERROR)

## 1. Local Distributed XGBoost Model using Dask

By default, XGBoost trains models sequentially. This is fine for smaller projects, but when the size of your dataset and/or ML model exceeds the limitations of your local machine, you will want to leverage the potential of distributed computing.

Starting from version 1.0, XGBoost comes with a native Dask integration that makes this possible. 

It only requires two changes to your regular XGBoostcode:
1. substitute `dtrain = xgb.DMatrix(X_train, y_train)` with `dtrain = xgb.dask.DaskDMatrix(X_train, y_train)`, and
2. substitute `xgb.train(params, dtrain, ...)` with `xgb.dask.train(client, params, dtrain, ...)`

Let's see this in action.

### Instantiate Dask Cluster

We'll begin by instantiating a local version of the Dask distributed scheduler, which will orchestrate the distributed training of our model. Read more about the Dask schedulers [here](https://distributed.dask.org/en/latest/).

In [2]:
from dask.distributed import Client, LocalCluster

# local dask cluster
cluster = LocalCluster(n_workers=4)
client = Client(cluster)
client

0,1
Connection method: Cluster object,Cluster type: distributed.LocalCluster
Dashboard: http://127.0.0.1:51647/status,

0,1
Dashboard: http://127.0.0.1:51647/status,Workers: 4
Total threads: 8,Total memory: 16.00 GiB
Status: running,Using processes: True

0,1
Comm: tcp://127.0.0.1:51648,Workers: 4
Dashboard: http://127.0.0.1:51647/status,Total threads: 8
Started: Just now,Total memory: 16.00 GiB

0,1
Comm: tcp://127.0.0.1:51668,Total threads: 2
Dashboard: http://127.0.0.1:51673/status,Memory: 4.00 GiB
Nanny: tcp://127.0.0.1:51653,
Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-wnmgn8ge,Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-wnmgn8ge

0,1
Comm: tcp://127.0.0.1:51663,Total threads: 2
Dashboard: http://127.0.0.1:51664/status,Memory: 4.00 GiB
Nanny: tcp://127.0.0.1:51651,
Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-lmmugaby,Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-lmmugaby

0,1
Comm: tcp://127.0.0.1:51666,Total threads: 2
Dashboard: http://127.0.0.1:51667/status,Memory: 4.00 GiB
Nanny: tcp://127.0.0.1:51654,
Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-kv5uh7xp,Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-kv5uh7xp

0,1
Comm: tcp://127.0.0.1:51670,Total threads: 2
Dashboard: http://127.0.0.1:51671/status,Memory: 4.00 GiB
Nanny: tcp://127.0.0.1:51652,
Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-1cegiq8q,Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-1cegiq8q


### Import the Data
Since we are working with a synthetic dataset, we can import the data and start training right away. No preprocessing needed.

For an example notebook with real-world data that does include some preprocessing work, check out [this notebook](https://github.com/coiled/coiled-resources/blob/main/xgboost-with-coiled/coiled-xgboost-arcos-20GB.ipynb) that trains an XGBoost model on a 20GB subset of the ARCOS dataset.

In [3]:
import dask.dataframe as dd

# download data from S3
data = dd.read_parquet(
    "s3://coiled-datasets/synthetic-data/synth-reg-104GB.parquet/", 
    compression="lz4",
    storage_options={"anon": True, 'use_ssl': True},
)

In [4]:
data

Unnamed: 0_level_0,0,1,2,3,4,5,6,7,8,9,10,11,12,13,14,15,16,17,18,19,20,21,22,23,24,25,26,27,28,29,30,31,32,33,34,35,36,37,38,39,40,41,42,43,44,45,46,47,48,49,target
npartitions=2750,Unnamed: 1_level_1,Unnamed: 2_level_1,Unnamed: 3_level_1,Unnamed: 4_level_1,Unnamed: 5_level_1,Unnamed: 6_level_1,Unnamed: 7_level_1,Unnamed: 8_level_1,Unnamed: 9_level_1,Unnamed: 10_level_1,Unnamed: 11_level_1,Unnamed: 12_level_1,Unnamed: 13_level_1,Unnamed: 14_level_1,Unnamed: 15_level_1,Unnamed: 16_level_1,Unnamed: 17_level_1,Unnamed: 18_level_1,Unnamed: 19_level_1,Unnamed: 20_level_1,Unnamed: 21_level_1,Unnamed: 22_level_1,Unnamed: 23_level_1,Unnamed: 24_level_1,Unnamed: 25_level_1,Unnamed: 26_level_1,Unnamed: 27_level_1,Unnamed: 28_level_1,Unnamed: 29_level_1,Unnamed: 30_level_1,Unnamed: 31_level_1,Unnamed: 32_level_1,Unnamed: 33_level_1,Unnamed: 34_level_1,Unnamed: 35_level_1,Unnamed: 36_level_1,Unnamed: 37_level_1,Unnamed: 38_level_1,Unnamed: 39_level_1,Unnamed: 40_level_1,Unnamed: 41_level_1,Unnamed: 42_level_1,Unnamed: 43_level_1,Unnamed: 44_level_1,Unnamed: 45_level_1,Unnamed: 46_level_1,Unnamed: 47_level_1,Unnamed: 48_level_1,Unnamed: 49_level_1,Unnamed: 50_level_1,Unnamed: 51_level_1
,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64
,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...
...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...
,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...
,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...


Since we're working locally to begin with, we won't be able to process the entire 100GB dataset.

We'll subset the first 10 partitions and persist them to our Dask cluster memory for quicker access.

In [5]:
# select the first 10 partitions
data_local = data.partitions[0:10]
data_local = data_local.persist()

In [6]:
# inspect the first 5 entries
data_local.head()

Unnamed: 0,0,1,2,3,4,5,6,7,8,9,...,41,42,43,44,45,46,47,48,49,target
0,-1.083516,0.173372,-0.973546,-1.465443,1.973955,-0.922526,1.058072,0.302878,1.160762,-0.690999,...,0.478698,-1.286906,0.037474,-0.448159,-0.652509,-1.205982,0.166634,2.526275,-0.890744,223.602485
1,2.077819,-0.507675,1.188347,-0.958974,0.666332,0.699718,0.416365,-0.006916,-0.561665,-0.535323,...,-0.406144,-0.122424,1.623143,0.438106,-1.510411,-0.909098,-0.416044,0.16966,-1.343285,-63.876627
2,-1.545396,-1.001309,-0.185548,-0.507883,1.223005,0.405486,-0.838138,-0.521867,1.16429,0.566665,...,1.341402,-0.206474,-1.203585,0.7965,-2.083753,0.670345,1.243194,-0.513658,-1.388109,182.856379
3,-0.548436,-0.754629,1.62849,0.954295,0.190117,-0.359459,1.901831,-0.137075,-0.005027,0.918249,...,1.214883,-0.115838,0.287735,-0.115192,-0.49933,0.349165,-1.618127,1.421938,-0.43924,-211.527657
4,-0.981102,0.993449,-0.173022,0.503123,0.823864,0.083351,0.242027,0.661806,0.463781,-0.799858,...,-0.98889,-0.541225,-0.298992,0.306095,0.351885,2.269911,0.465673,0.909917,0.513545,-165.464021


This is looking good.

### Train-Test Splits

The next step is to define our train and test splits. The target feature in this synthetic dataset is the last column, conveniently named "target".

In [7]:
from dask_ml.model_selection import train_test_split

In [8]:
# Create the train-test split
X, y = data_local.iloc[:, :-1], data_local["target"]
X_train, X_test, y_train, y_test = train_test_split(
    X, y, test_size=0.3, shuffle=True, random_state=21
)

### Train XGBoost Model

Now we're all set to start training our XGBoost model.

First, we'll create the XGBoost DMatrix and set the model parameters. We'll use the default parameters for this example.

For more information on training XGBoost models and setting model parameter, have a look at the [XGBoost documentation](https://xgboost.readthedocs.io/en/latest/get_started.html).

In [9]:
import xgboost as xgb

In [10]:
# Create the XGBoost DMatrix for our training and testing splits
dtrain = xgb.dask.DaskDMatrix(client, X_train, y_train)
dtest = xgb.dask.DaskDMatrix(client, X_test, y_test)

# Set model parameters (XGBoost defaults)
params = {
    "max_depth": 6,
    "gamma": 0,
    "eta": 0.3,
    "min_child_weight": 30,
    "objective": "reg:squarederror",
    "grow_policy": "depthwise"
}

Then let's go ahead and train the model.

In [11]:
%%time 
# train the model
output = xgb.dask.train(
    client, params, dtrain, num_boost_round=5,
    evals=[(dtrain, 'train')]
)

  from pandas import MultiIndex, Int64Index
  from pandas import MultiIndex, Int64Index
  from pandas import MultiIndex, Int64Index
  from pandas import MultiIndex, Int64Index
[11:36:41] task [xgboost.dask]:tcp://127.0.0.1:51670 got new rank 0
[11:36:41] task [xgboost.dask]:tcp://127.0.0.1:51668 got new rank 1
[11:36:41] task [xgboost.dask]:tcp://127.0.0.1:51663 got new rank 2
[11:36:41] task [xgboost.dask]:tcp://127.0.0.1:51666 got new rank 3
  elif isinstance(data.columns, (pd.Int64Index, pd.RangeIndex)):
  elif isinstance(data.columns, (pd.Int64Index, pd.RangeIndex)):
  elif isinstance(data.columns, (pd.Int64Index, pd.RangeIndex)):
  elif isinstance(data.columns, (pd.Int64Index, pd.RangeIndex)):


[0]	train-rmse:191.10689
[1]	train-rmse:166.94313
[2]	train-rmse:147.69171
[3]	train-rmse:132.75594
[4]	train-rmse:120.05722
CPU times: user 121 ms, sys: 38.8 ms, total: 160 ms
Wall time: 3.9 s


And use our trained model together with our testing split to make predictions.

In [12]:
# make predictions
y_pred = xgb.dask.predict(client, output, dtest)

And finally, let's evaluate our results by getting the accuracy score.

In [13]:
from sklearn.metrics import mean_absolute_error

In [14]:
mae = mean_absolute_error(y_test, y_pred)
print(f"Mean Absolute Error: {mae}")

Mean Absolute Error: 94.71184525241688


### Try Locally with Entire Dataset... if you dare...

Unless you're running this on a supercomputer, uncommenting and running the cell below will likely not complete.

But don't just take our word for it, of course ;)

In [None]:
# # Create the train-test split
# X, y = data.iloc[:, :-1], data["target"]
# X_train, X_test, y_train, y_test = train_test_split(
#     X, y, test_size=0.3, shuffle=True, random_state=2
# )

# # Create DaskDMatrices
# dtrain = xgb.dask.DaskDMatrix(client, X_train, y_train)
# dtest = xgb.dask.DaskDMatrix(client, X_test, y_test)

```MemoryError
distributed.batched - ERROR - Error in batched write
```
```
MemoryError
```
```
distributed.worker - WARNING - Worker is at 80% memory usage. Pausing worker.  Process memory: 1.49 GiB -- Worker memory limit: 1.86 GiB
```

## 2. Distributed XGBoost in the Cloud using Dask and Coiled

Let's now expand this workflow to process the entire dataset (>100 GB). 

We'll run the exact same code as above except for **2 changes**:
1. We'll connect Dask to a Coiled cluster in the cloud, instead of to our local CPU cores,
2. We'll work with the entire 100GB dataset, instead of the first 10 partitions.

In the section below we've copied and pasted the cells from above so that you can run this notebook from top to bottom in one go. Alternatively, you could run the cell below (where we instantiate the Coiled Cluster) and then simply re-run the cells above -- making sure to work with the entire dataset, of course.

### Instantiate Coiled Cluster
Let's create our Coiled cluster in the cloud. 

We'll specify a cluster of 50 workers, with 4 CPU cores and 16GB of RAM each. That will allow the entire dataset to fit into the cluster's memory comfortably and should make for quick training.

> *Note: if you're running this using the Coiled Free Tier, you'll want to reduce your **n_workers** to 25 to stay within the Total Core limit.*

In [1]:
import coiled

cluster = coiled.Cluster(
    name="xgboost",
    software="coiled-examples/xgboost",
    n_workers=50,
    worker_cpu=4,
    worker_memory='16Gib',
    shutdown_on_close=False,
)



Found software environment build
Created fw rule: inbound [8786-8787] [0.0.0.0/0] []
Created FW rules: coiled-dask-rrpelgr71-80652-firewall
Created fw rule: cluster [0-65535] [None] [coiled-dask-rrpelgr71-80652-firewall -> coiled-dask-rrpelgr71-80652-firewall]
Created FW rules: coiled-dask-rrpelgr71-80652-cluster-firewall
Created fw rule: cluster [0-65535] [None] [coiled-dask-rrpelgr71-80652-cluster-firewall -> coiled-dask-rrpelgr71-80652-cluster-firewall]
Created scheduler VM: coiled-dask-rrpelgr71-80652-scheduler (type: t3a.medium, ip: ['44.196.47.64'])


In [3]:
from distributed import Client

client = Client(cluster)
client

0,1
Connection method: Cluster object,Cluster type: coiled.Cluster
Dashboard: http://44.196.47.64:8787,

0,1
Dashboard: http://44.196.47.64:8787,Workers: 38
Total threads: 152,Total memory: 587.02 GiB

0,1
Comm: tls://10.4.3.21:8786,Workers: 38
Dashboard: http://10.4.3.21:8787/status,Total threads: 152
Started: 1 minute ago,Total memory: 587.02 GiB

0,1
Comm: tls://10.4.10.160:44261,Total threads: 4
Dashboard: http://10.4.10.160:42177/status,Memory: 15.45 GiB
Nanny: tls://10.4.10.160:34503,
Local directory: /dask-worker-space/worker-l3unmyhv,Local directory: /dask-worker-space/worker-l3unmyhv

0,1
Comm: tls://10.4.11.95:46471,Total threads: 4
Dashboard: http://10.4.11.95:44063/status,Memory: 15.45 GiB
Nanny: tls://10.4.11.95:35969,
Local directory: /dask-worker-space/worker-1eyt3_1k,Local directory: /dask-worker-space/worker-1eyt3_1k

0,1
Comm: tls://10.4.14.3:45679,Total threads: 4
Dashboard: http://10.4.14.3:39797/status,Memory: 15.45 GiB
Nanny: tls://10.4.14.3:41713,
Local directory: /dask-worker-space/worker-usjoa5xi,Local directory: /dask-worker-space/worker-usjoa5xi

0,1
Comm: tls://10.4.9.213:41655,Total threads: 4
Dashboard: http://10.4.9.213:46159/status,Memory: 15.45 GiB
Nanny: tls://10.4.9.213:36173,
Local directory: /dask-worker-space/worker-fs_zh10d,Local directory: /dask-worker-space/worker-fs_zh10d

0,1
Comm: tls://10.4.10.196:35857,Total threads: 4
Dashboard: http://10.4.10.196:45587/status,Memory: 15.45 GiB
Nanny: tls://10.4.10.196:32911,
Local directory: /dask-worker-space/worker-y9wnzjtg,Local directory: /dask-worker-space/worker-y9wnzjtg

0,1
Comm: tls://10.4.0.63:43785,Total threads: 4
Dashboard: http://10.4.0.63:36609/status,Memory: 15.45 GiB
Nanny: tls://10.4.0.63:39887,
Local directory: /dask-worker-space/worker-g1nfftw8,Local directory: /dask-worker-space/worker-g1nfftw8

0,1
Comm: tls://10.4.5.227:37369,Total threads: 4
Dashboard: http://10.4.5.227:33423/status,Memory: 15.45 GiB
Nanny: tls://10.4.5.227:34585,
Local directory: /dask-worker-space/worker-ucpd57cr,Local directory: /dask-worker-space/worker-ucpd57cr

0,1
Comm: tls://10.4.1.120:37339,Total threads: 4
Dashboard: http://10.4.1.120:38143/status,Memory: 15.45 GiB
Nanny: tls://10.4.1.120:39929,
Local directory: /dask-worker-space/worker-evkux69r,Local directory: /dask-worker-space/worker-evkux69r

0,1
Comm: tls://10.4.8.112:35651,Total threads: 4
Dashboard: http://10.4.8.112:34643/status,Memory: 15.45 GiB
Nanny: tls://10.4.8.112:34723,
Local directory: /dask-worker-space/worker-gmyit6lv,Local directory: /dask-worker-space/worker-gmyit6lv

0,1
Comm: tls://10.4.7.240:37825,Total threads: 4
Dashboard: http://10.4.7.240:39669/status,Memory: 15.45 GiB
Nanny: tls://10.4.7.240:40589,
Local directory: /dask-worker-space/worker-03jw7z1m,Local directory: /dask-worker-space/worker-03jw7z1m

0,1
Comm: tls://10.4.0.163:46119,Total threads: 4
Dashboard: http://10.4.0.163:35429/status,Memory: 15.45 GiB
Nanny: tls://10.4.0.163:43255,
Local directory: /dask-worker-space/worker-et7ml9g3,Local directory: /dask-worker-space/worker-et7ml9g3

0,1
Comm: tls://10.4.0.42:43771,Total threads: 4
Dashboard: http://10.4.0.42:39315/status,Memory: 15.45 GiB
Nanny: tls://10.4.0.42:42267,
Local directory: /dask-worker-space/worker-58xcgmq3,Local directory: /dask-worker-space/worker-58xcgmq3

0,1
Comm: tls://10.4.15.12:46217,Total threads: 4
Dashboard: http://10.4.15.12:40009/status,Memory: 15.45 GiB
Nanny: tls://10.4.15.12:33087,
Local directory: /dask-worker-space/worker-4av8j8vg,Local directory: /dask-worker-space/worker-4av8j8vg

0,1
Comm: tls://10.4.11.117:42267,Total threads: 4
Dashboard: http://10.4.11.117:39013/status,Memory: 15.45 GiB
Nanny: tls://10.4.11.117:39563,
Local directory: /dask-worker-space/worker-uupti2pa,Local directory: /dask-worker-space/worker-uupti2pa

0,1
Comm: tls://10.4.5.211:46587,Total threads: 4
Dashboard: http://10.4.5.211:41057/status,Memory: 15.45 GiB
Nanny: tls://10.4.5.211:44041,
Local directory: /dask-worker-space/worker-ufhwh7wn,Local directory: /dask-worker-space/worker-ufhwh7wn

0,1
Comm: tls://10.4.13.45:40735,Total threads: 4
Dashboard: http://10.4.13.45:40233/status,Memory: 15.45 GiB
Nanny: tls://10.4.13.45:39371,
Local directory: /dask-worker-space/worker-piyitx6z,Local directory: /dask-worker-space/worker-piyitx6z

0,1
Comm: tls://10.4.2.207:40725,Total threads: 4
Dashboard: http://10.4.2.207:39665/status,Memory: 15.45 GiB
Nanny: tls://10.4.2.207:40871,
Local directory: /dask-worker-space/worker-iybzijhc,Local directory: /dask-worker-space/worker-iybzijhc

0,1
Comm: tls://10.4.1.142:33811,Total threads: 4
Dashboard: http://10.4.1.142:37329/status,Memory: 15.45 GiB
Nanny: tls://10.4.1.142:33619,
Local directory: /dask-worker-space/worker-9xr2r6rp,Local directory: /dask-worker-space/worker-9xr2r6rp

0,1
Comm: tls://10.4.8.156:43101,Total threads: 4
Dashboard: http://10.4.8.156:37169/status,Memory: 15.45 GiB
Nanny: tls://10.4.8.156:35633,
Local directory: /dask-worker-space/worker-ayfkbmzj,Local directory: /dask-worker-space/worker-ayfkbmzj

0,1
Comm: tls://10.4.11.173:42661,Total threads: 4
Dashboard: http://10.4.11.173:46413/status,Memory: 15.45 GiB
Nanny: tls://10.4.11.173:44993,
Local directory: /dask-worker-space/worker-3s3jk988,Local directory: /dask-worker-space/worker-3s3jk988

0,1
Comm: tls://10.4.7.212:34085,Total threads: 4
Dashboard: http://10.4.7.212:42197/status,Memory: 15.45 GiB
Nanny: tls://10.4.7.212:32775,
Local directory: /dask-worker-space/worker-p_fj21yu,Local directory: /dask-worker-space/worker-p_fj21yu

0,1
Comm: tls://10.4.13.225:39301,Total threads: 4
Dashboard: http://10.4.13.225:37559/status,Memory: 15.45 GiB
Nanny: tls://10.4.13.225:38955,
Local directory: /dask-worker-space/worker-gcdp1rjv,Local directory: /dask-worker-space/worker-gcdp1rjv

0,1
Comm: tls://10.4.2.253:39385,Total threads: 4
Dashboard: http://10.4.2.253:40775/status,Memory: 15.45 GiB
Nanny: tls://10.4.2.253:35889,
Local directory: /dask-worker-space/worker-h03xtua2,Local directory: /dask-worker-space/worker-h03xtua2

0,1
Comm: tls://10.4.5.23:32775,Total threads: 4
Dashboard: http://10.4.5.23:34893/status,Memory: 15.45 GiB
Nanny: tls://10.4.5.23:36807,
Local directory: /dask-worker-space/worker-pfttt54m,Local directory: /dask-worker-space/worker-pfttt54m

0,1
Comm: tls://10.4.15.216:45527,Total threads: 4
Dashboard: http://10.4.15.216:39531/status,Memory: 15.45 GiB
Nanny: tls://10.4.15.216:36961,
Local directory: /dask-worker-space/worker-6qtoi64u,Local directory: /dask-worker-space/worker-6qtoi64u

0,1
Comm: tls://10.4.2.67:46271,Total threads: 4
Dashboard: http://10.4.2.67:37367/status,Memory: 15.45 GiB
Nanny: tls://10.4.2.67:39955,
Local directory: /dask-worker-space/worker-0qg91e29,Local directory: /dask-worker-space/worker-0qg91e29

0,1
Comm: tls://10.4.4.189:42661,Total threads: 4
Dashboard: http://10.4.4.189:33933/status,Memory: 15.45 GiB
Nanny: tls://10.4.4.189:33121,
Local directory: /dask-worker-space/worker-ic3d9d8r,Local directory: /dask-worker-space/worker-ic3d9d8r

0,1
Comm: tls://10.4.8.228:35831,Total threads: 4
Dashboard: http://10.4.8.228:41941/status,Memory: 15.45 GiB
Nanny: tls://10.4.8.228:44775,
Local directory: /dask-worker-space/worker-58ttbcoy,Local directory: /dask-worker-space/worker-58ttbcoy

0,1
Comm: tls://10.4.12.111:41059,Total threads: 4
Dashboard: http://10.4.12.111:40071/status,Memory: 15.45 GiB
Nanny: tls://10.4.12.111:42587,
Local directory: /dask-worker-space/worker-yux3feo7,Local directory: /dask-worker-space/worker-yux3feo7

0,1
Comm: tls://10.4.3.160:41751,Total threads: 4
Dashboard: http://10.4.3.160:40305/status,Memory: 15.45 GiB
Nanny: tls://10.4.3.160:37963,
Local directory: /dask-worker-space/worker-btirrgy_,Local directory: /dask-worker-space/worker-btirrgy_

0,1
Comm: tls://10.4.1.67:45725,Total threads: 4
Dashboard: http://10.4.1.67:42717/status,Memory: 15.45 GiB
Nanny: tls://10.4.1.67:46001,
Local directory: /dask-worker-space/worker-mlk9lo9_,Local directory: /dask-worker-space/worker-mlk9lo9_

0,1
Comm: tls://10.4.11.122:35027,Total threads: 4
Dashboard: http://10.4.11.122:32825/status,Memory: 15.45 GiB
Nanny: tls://10.4.11.122:33661,
Local directory: /dask-worker-space/worker-1qvsk9_9,Local directory: /dask-worker-space/worker-1qvsk9_9

0,1
Comm: tls://10.4.1.170:41761,Total threads: 4
Dashboard: http://10.4.1.170:38021/status,Memory: 15.45 GiB
Nanny: tls://10.4.1.170:40861,
Local directory: /dask-worker-space/worker-afr_k5s3,Local directory: /dask-worker-space/worker-afr_k5s3

0,1
Comm: tls://10.4.10.12:36857,Total threads: 4
Dashboard: http://10.4.10.12:46363/status,Memory: 15.45 GiB
Nanny: tls://10.4.10.12:41035,
Local directory: /dask-worker-space/worker-0km5w0fi,Local directory: /dask-worker-space/worker-0km5w0fi

0,1
Comm: tls://10.4.6.203:36915,Total threads: 4
Dashboard: http://10.4.6.203:45331/status,Memory: 15.45 GiB
Nanny: tls://10.4.6.203:32807,
Local directory: /dask-worker-space/worker-o9u5b48y,Local directory: /dask-worker-space/worker-o9u5b48y

0,1
Comm: tls://10.4.9.167:33857,Total threads: 4
Dashboard: http://10.4.9.167:33519/status,Memory: 15.45 GiB
Nanny: tls://10.4.9.167:39243,
Local directory: /dask-worker-space/worker-csfm31nb,Local directory: /dask-worker-space/worker-csfm31nb

0,1
Comm: tls://10.4.10.33:36041,Total threads: 4
Dashboard: http://10.4.10.33:33051/status,Memory: 15.45 GiB
Nanny: tls://10.4.10.33:41529,
Local directory: /dask-worker-space/worker-my8kzk5r,Local directory: /dask-worker-space/worker-my8kzk5r

0,1
Comm: tls://10.4.9.190:43913,Total threads: 4
Dashboard: http://10.4.9.190:43709/status,Memory: 15.45 GiB
Nanny: tls://10.4.9.190:36569,
Local directory: /dask-worker-space/worker-b8fv6v2s,Local directory: /dask-worker-space/worker-b8fv6v2s


### Inspecting Entire Dataset

Let's load the entire dataset into our Dask dataframe **data**.

As you can see below, it consists of 2750 partitions.

In [4]:
import dask.dataframe as dd

In [5]:
data = dd.read_parquet(
    "s3://coiled-datasets/synthetic-data/synth-reg-104GB.parquet/", 
    compression="lz4",
    storage_options={"anon": True, 'use_ssl': True},
)


In [6]:
data.head()

Unnamed: 0,0,1,2,3,4,5,6,7,8,9,...,41,42,43,44,45,46,47,48,49,target
0,-1.083516,0.173372,-0.973546,-1.465443,1.973955,-0.922526,1.058072,0.302878,1.160762,-0.690999,...,0.478698,-1.286906,0.037474,-0.448159,-0.652509,-1.205982,0.166634,2.526275,-0.890744,223.602485
1,2.077819,-0.507675,1.188347,-0.958974,0.666332,0.699718,0.416365,-0.006916,-0.561665,-0.535323,...,-0.406144,-0.122424,1.623143,0.438106,-1.510411,-0.909098,-0.416044,0.16966,-1.343285,-63.876627
2,-1.545396,-1.001309,-0.185548,-0.507883,1.223005,0.405486,-0.838138,-0.521867,1.16429,0.566665,...,1.341402,-0.206474,-1.203585,0.7965,-2.083753,0.670345,1.243194,-0.513658,-1.388109,182.856379
3,-0.548436,-0.754629,1.62849,0.954295,0.190117,-0.359459,1.901831,-0.137075,-0.005027,0.918249,...,1.214883,-0.115838,0.287735,-0.115192,-0.49933,0.349165,-1.618127,1.421938,-0.43924,-211.527657
4,-0.981102,0.993449,-0.173022,0.503123,0.823864,0.083351,0.242027,0.661806,0.463781,-0.799858,...,-0.98889,-0.541225,-0.298992,0.306095,0.351885,2.269911,0.465673,0.909917,0.513545,-165.464021


In [7]:
data

Unnamed: 0_level_0,0,1,2,3,4,5,6,7,8,9,10,11,12,13,14,15,16,17,18,19,20,21,22,23,24,25,26,27,28,29,30,31,32,33,34,35,36,37,38,39,40,41,42,43,44,45,46,47,48,49,target
npartitions=2750,Unnamed: 1_level_1,Unnamed: 2_level_1,Unnamed: 3_level_1,Unnamed: 4_level_1,Unnamed: 5_level_1,Unnamed: 6_level_1,Unnamed: 7_level_1,Unnamed: 8_level_1,Unnamed: 9_level_1,Unnamed: 10_level_1,Unnamed: 11_level_1,Unnamed: 12_level_1,Unnamed: 13_level_1,Unnamed: 14_level_1,Unnamed: 15_level_1,Unnamed: 16_level_1,Unnamed: 17_level_1,Unnamed: 18_level_1,Unnamed: 19_level_1,Unnamed: 20_level_1,Unnamed: 21_level_1,Unnamed: 22_level_1,Unnamed: 23_level_1,Unnamed: 24_level_1,Unnamed: 25_level_1,Unnamed: 26_level_1,Unnamed: 27_level_1,Unnamed: 28_level_1,Unnamed: 29_level_1,Unnamed: 30_level_1,Unnamed: 31_level_1,Unnamed: 32_level_1,Unnamed: 33_level_1,Unnamed: 34_level_1,Unnamed: 35_level_1,Unnamed: 36_level_1,Unnamed: 37_level_1,Unnamed: 38_level_1,Unnamed: 39_level_1,Unnamed: 40_level_1,Unnamed: 41_level_1,Unnamed: 42_level_1,Unnamed: 43_level_1,Unnamed: 44_level_1,Unnamed: 45_level_1,Unnamed: 46_level_1,Unnamed: 47_level_1,Unnamed: 48_level_1,Unnamed: 49_level_1,Unnamed: 50_level_1,Unnamed: 51_level_1
0,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64,float64
100000,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...
...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...
274900000,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...
274999999,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...,...


### Train / Test Splits

Below we apply the same code we used above to create out training and testing splits. 

We also persist the splits to the cluster's memory for faster training.

In [8]:
from dask_ml.model_selection import train_test_split

In [9]:
# Create the train-test split
X, y = data.iloc[:, :-1], data["target"]
X_train, X_test, y_train, y_test = train_test_split(
    X, y, test_size=0.3, shuffle=True, random_state=13
)

# persist the train/test splits to cluster memory to speed up training
import dask
dask.persist(X_train, X_test, y_train, y_test)

(Dask DataFrame Structure:
                         0        1        2        3        4        5        6        7        8        9       10       11       12       13       14       15       16       17       18       19       20       21       22       23       24       25       26       27       28       29       30       31       32       33       34       35       36       37       38       39       40       41       42       43       44       45       46       47       48       49
 npartitions=2750                                                                                                                                                                                                                                                                                                                                                                                                                                                                  
 0                 float64  float64  

### XGBoost Training
Alright, the moment we've all been waiting for!

You're now all set to train your distributed XGBoost model on the entire 500GB dataset.

The cells below will create the DaskDMatrix, set the model parameters (using the XGBoost defaults for now) and train your XGBoost model.

In [10]:
import xgboost as xgb

In [11]:
%%time
# Create the XGBoost DMatrices
dtrain = xgb.dask.DaskDMatrix(client, X_train, y_train)
dtest = xgb.dask.DaskDMatrix(client, X_test, y_test)

CPU times: user 7.09 s, sys: 803 ms, total: 7.9 s
Wall time: 1min 22s


In [12]:
# Set model parameters (XGBoost defaults)
params = {
    "max_depth": 6,
    "gamma": 0,
    "eta": 0.3,
    "min_child_weight": 30,
    "objective": "reg:squarederror",
    "grow_policy": "depthwise"
}

In [13]:
%%time 
# train the model 
output = xgb.dask.train(
    client, params, dtrain, num_boost_round=5,
    evals=[(dtrain, 'train')]
)

CPU times: user 1.23 s, sys: 111 ms, total: 1.35 s
Wall time: 3min 47s


In [14]:
%%time
# make predictions
y_pred = xgb.dask.predict(client, output, dtest)

CPU times: user 2.57 s, sys: 96.1 ms, total: 2.67 s
Wall time: 9.15 s


Great work! You just trained an XGBoost model on 100GB of data in a matter of minutes!

### Shutting down the cluster
After our training is done, we can close down the cluster, releasing the resources. Should you forget to do so for whatever reason, Coiled automatically shuts down clusters after 20 minutes of inactivity, to help avoid unnecessary costs.


In [15]:
# Shut down the cluster
client.close()

distributed.client - ERROR - Failed to reconnect to scheduler after 30.00 seconds, closing client
_GatheringFuture exception was never retrieved
future: <_GatheringFuture finished exception=CancelledError()>
asyncio.exceptions.CancelledError
Traceback (most recent call last):
  File "/Users/rpelgrim/mambaforge/envs/xgboost/lib/python3.9/site-packages/distributed/comm/tcp.py", line 398, in connect
    stream = await self.client.connect(
  File "/Users/rpelgrim/mambaforge/envs/xgboost/lib/python3.9/site-packages/tornado/tcpclient.py", line 275, in connect
    af, addr, stream = await connector.start(connect_timeout=timeout)
asyncio.exceptions.CancelledError

During handling of the above exception, another exception occurred:

Traceback (most recent call last):
  File "/Users/rpelgrim/mambaforge/envs/xgboost/lib/python3.9/asyncio/tasks.py", line 492, in wait_for
    fut.result()
asyncio.exceptions.CancelledError

The above exception was the direct cause of the following exception:

Traceb

## 3. Pro Tips to Speed Up Training
Below we’ve collected some pro tips straight from the Dask core team to help you speed up your XGBoost training:

- Increase the number of workers in your Coiled cluster using the `n_workers` keyword argument.
- Re-cast numerical columns to less memory-intensive dtypes. For example, convert float64 into int16 whenever possible. This will reduce the memory load of your dataframe and thereby speed up training.
- The Dask Dashboard is a great way to spot bottle-necks and identify opportunities for increased performance in your code. Watch the initial author of Dask, Matt Rocklin, explain how to get the most out of the Dask Dashboard [here](https://www.youtube.com/watch?v=N_GqzcuGLCY).
- Read Matthew Power’s blog on setting up the Dask Dashboard in your Jupyter Lab environment [here](https://coiled.io/blog/dask-jupyterlab-workflow/). 
- Read Dask core contributor Guido Imperiale’s blog on how to tackle the specific issue of unmanaged memory in Dask workers [here](https://coiled.io/blog/tackling-unmanaged-memory-with-dask/). 



## 4. Recap

Let’s recap what we’ve discussed in this notebook:
- When training XGBoost with large datasets, running out of local memory can be a challenge. 
- Connecting XGboost to a local Dask cluster allows you to make the most out of the multiple cores in your machine.
- If that’s still not enough, you can connect Dask to Coiled and burst to the cloud as and when needed.
- You can tweak your distributed XGBoost performance by inspecting the Dask Dashboard.

We’d love to see you apply distributed XGBoost to a dataset that’s meaningful to you. If you’d like to try, swap your dataset into this notebook and see how well it does! 

Let us know how you get on in our [Coiled Community Slack channel](https://join.slack.com/t/coiled-users/shared_invite/zt-hx1fnr7k-In~Q8ui3XkQfvQon0yN5WQ) or by [tweeting](https://twitter.com/coiledhq) at us.