# Scaling XGBoost with Dask and Coiled

This notebook walks through training a distributed [XGBoost](https://xgboost.readthedocs.io/en/latest/) model locally on a small dataset using [Dask](https://dask.org/) and then using Dask and [Coiled](https://coiled.io/) to scale out to the cloud to run XGBoost on a larger-than-memory dataset.

In [None]:
# coiled.create_software_environment(
#     name='coiled-xgboost-test',
#     account='coiled-examples',
#     conda="/Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/xgboost-test.yml"
# )

## 1. Importing Libraries

We'll start by importing all the libraries we'll need to run this notebook.

In [1]:
import coiled
import dask.dataframe as  dd
from dask.distributed import Client, LocalCluster
from dask_ml.preprocessing import Categorizer
from dask_ml.model_selection import train_test_split
import xgboost as xgb

## 2. Local Distributed XGBoost Model using Dask

Next, let's instantiate a local version of the Dask distributed scheduler using the **LocalCluster** object. 

This object will handle parallelism for us on our local machine.

In [5]:
# local dask cluster
cluster = LocalCluster(n_workers=8)
client = Client(cluster)
client

0,1
Client  Scheduler: tcp://127.0.0.1:58322  Dashboard: http://127.0.0.1:8787/status,Cluster  Workers: 8  Cores: 8  Memory: 16.00 GiB


In [5]:
client.restart()

0,1
Connection method: Cluster object,Cluster type: LocalCluster
Dashboard: http://127.0.0.1:8787/status,

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

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

0,1
Comm: tcp://127.0.0.1:56414,Total threads: 1
Dashboard: http://127.0.0.1:56415/status,Memory: 2.00 GiB
Nanny: tcp://127.0.0.1:56280,
Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-qaab67cf,Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-qaab67cf

0,1
Comm: tcp://127.0.0.1:56402,Total threads: 1
Dashboard: http://127.0.0.1:56403/status,Memory: 2.00 GiB
Nanny: tcp://127.0.0.1:56281,
Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-hn7__frm,Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-hn7__frm

0,1
Comm: tcp://127.0.0.1:56411,Total threads: 1
Dashboard: http://127.0.0.1:56412/status,Memory: 2.00 GiB
Nanny: tcp://127.0.0.1:56287,
Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-pqlq_b5n,Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-pqlq_b5n

0,1
Comm: tcp://127.0.0.1:56405,Total threads: 1
Dashboard: http://127.0.0.1:56406/status,Memory: 2.00 GiB
Nanny: tcp://127.0.0.1:56286,
Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-rlxg8abu,Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-rlxg8abu

0,1
Comm: tcp://127.0.0.1:56399,Total threads: 1
Dashboard: http://127.0.0.1:56400/status,Memory: 2.00 GiB
Nanny: tcp://127.0.0.1:56285,
Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-gggel_rb,Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-gggel_rb

0,1
Comm: tcp://127.0.0.1:56396,Total threads: 1
Dashboard: http://127.0.0.1:56397/status,Memory: 2.00 GiB
Nanny: tcp://127.0.0.1:56284,
Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-75pzb_5a,Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-75pzb_5a

0,1
Comm: tcp://127.0.0.1:56408,Total threads: 1
Dashboard: http://127.0.0.1:56409/status,Memory: 2.00 GiB
Nanny: tcp://127.0.0.1:56283,
Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-upj0b7gk,Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-upj0b7gk

0,1
Comm: tcp://127.0.0.1:56393,Total threads: 1
Dashboard: http://127.0.0.1:56394/status,Memory: 2.00 GiB
Nanny: tcp://127.0.0.1:56282,
Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-kniumbdr,Local directory: /Users/rpelgrim/Documents/git/coiled-resources/xgboost-with-coiled/dask-worker-space/worker-kniumbdr


In [6]:
# Specify the columns we want to download
columns = [
    "interest_rate", "loan_age", "num_borrowers", 
    "borrower_credit_score", "num_units"
]

categorical = [
    "orig_channel", "occupancy_status", "property_state",
    "first_home_buyer", "loan_purpose", "property_type",
    "zip", "relocation_mortgage_indicator", "delinquency_12"
]

In [7]:
# Download data from S3
mortgage_data_local = dd.read_parquet(
    "s3://coiled-data/mortgage-2000.parq/part.*.parquet", 
    columns=columns + categorical, 
    storage_options={"anon": True}
)

# Cache the data on Cluster workers
mortgage_data_local = mortgage_data_local.persist()

In [None]:
# inspect the first 5 entries
mortgage_data_local.head()



This is looking good.

Before we can start training our XGBoost model, however, we'll have to conduct two preprocessing steps:
1. Cast our categorical columns to the correct types (XGBoost only accepts float, integer and boolean dtypes)
2. Create our train and test splits

*Note: we're using the **dask_ml** library for this, which mimics the familiar scikit-learn API*

In [None]:
# Cast categorical columns to the correct type
ce = Categorizer(columns=categorical)
mortgage_data_local = ce.fit_transform(mortgage_data_local)
for col in categorical:
    mortgage_data_local[col] = mortgage_data_local[col].cat.codes

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

X_train = X_train.persist()
X_test = X_test.persist()
y_train = y_train.persist()
y_test = y_test.persist()

Great, now we're all set to start training our XGBoost model.

First, we'll create the XGBoost DMatrix and set the model parameters.

In [None]:
# Create the XGBoost DMatrix

dtrain = xgb.dask.DaskDMatrix(client, X_train, y_train)

# Set parameters
params = {
    "max_depth": 8,
    "max_leaves": 2 ** 8,
    "gamma": 0.1,
    "eta": 0.1,
    "min_child_weight": 30,
    "objective": "binary:logistic",
    "grow_policy": "lossguide"
}


Then let's go ahead and train the model.

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

And see the results:

In [None]:
# 'booster' is the trained model
booster = output['booster']  

# 'history' is a dictionary containing evaluation metrics
history = output['history']  

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

## 3. Cloud-Based Distributed XGBoost using Dask and Coiled

Let's now expand this workflow to process the entire dataset (~200GB). We'll run almost exactly the same code 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 download the entire dataset, instead of a single partition.

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 adjust the cell that downloads the data as well, of course.

### Instantiate Coiled Cluster
Let's create our Coiled cluster in the cloud. We'll specify a cluster of 20 workers, with 4 CPU cores and 16GB of RAM each. That should allow the entire dataset to fit into the cluster's memory comfortably.

In [2]:
# Create Coiled Cloud cluster
cluster = coiled.Cluster(
    name='xgboost-2',
    n_workers=10,
    worker_cpu=4,
    worker_memory='24GiB',
    account='coiled-examples',
    software='rrpelgrim/coiled-xgboost',
    shutdown_on_close=False,
    scheduler_options={'idle_timeout':'2hours'}
)

# Connect Dask client to the Coiled cluster
client = Client(cluster)
client

Output()

Found software environment build


0,1
Connection method: Cluster object,Cluster type: Cluster
Dashboard: http://ec2-3-80-18-232.compute-1.amazonaws.com:8787,

0,1
Dashboard: http://ec2-3-80-18-232.compute-1.amazonaws.com:8787,Workers: 5
Total threads:  20,Total memory:  120.00 GiB

0,1
Comm: tls://10.3.190.54:8786,Workers: 5
Dashboard: http://10.3.190.54:8787/status,Total threads:  20
Started:  Just now,Total memory:  120.00 GiB

0,1
Comm: tls://10.3.183.203:40233,Total threads: 4
Dashboard: http://10.3.183.203:37555/status,Memory: 24.00 GiB
Nanny: tls://10.3.183.203:44463,
Local directory: /dask-worker-space/worker-g_u2v5n0,Local directory: /dask-worker-space/worker-g_u2v5n0

0,1
Comm: tls://10.3.170.107:33391,Total threads: 4
Dashboard: http://10.3.170.107:43579/status,Memory: 24.00 GiB
Nanny: tls://10.3.170.107:41903,
Local directory: /dask-worker-space/worker-5u89sirr,Local directory: /dask-worker-space/worker-5u89sirr

0,1
Comm: tls://10.3.130.22:45077,Total threads: 4
Dashboard: http://10.3.130.22:35437/status,Memory: 24.00 GiB
Nanny: tls://10.3.130.22:38953,
Local directory: /dask-worker-space/worker-qbgz7e80,Local directory: /dask-worker-space/worker-qbgz7e80

0,1
Comm: tls://10.3.126.71:38269,Total threads: 4
Dashboard: http://10.3.126.71:46189/status,Memory: 24.00 GiB
Nanny: tls://10.3.126.71:34219,
Local directory: /dask-worker-space/worker-nip35mgf,Local directory: /dask-worker-space/worker-nip35mgf

0,1
Comm: tls://10.3.120.105:38171,Total threads: 4
Dashboard: http://10.3.120.105:39781/status,Memory: 24.00 GiB
Nanny: tls://10.3.120.105:39035,
Local directory: /dask-worker-space/worker-l1ov8oao,Local directory: /dask-worker-space/worker-l1ov8oao


### Download the Data

In [26]:
client.restart()

0,1
Connection method: Cluster object,Cluster type: Cluster
Dashboard: http://3.88.2.193:8787,

0,1
Dashboard: http://3.88.2.193:8787,Workers: 10
Total threads:  40,Total memory:  240.00 GiB

0,1
Comm: tls://10.4.0.16:8786,Workers: 10
Dashboard: http://10.4.0.16:8787/status,Total threads:  40
Started:  4 minutes ago,Total memory:  240.00 GiB

0,1
Comm: tls://10.4.1.73:43089,Total threads: 4
Dashboard: http://10.4.1.73:40847/status,Memory: 24.00 GiB
Nanny: tls://10.4.1.73:45875,
Local directory: /dask-worker-space/worker-ubkp5wcs,Local directory: /dask-worker-space/worker-ubkp5wcs

0,1
Comm: tls://10.4.1.46:42887,Total threads: 4
Dashboard: http://10.4.1.46:36255/status,Memory: 24.00 GiB
Nanny: tls://10.4.1.46:42025,
Local directory: /dask-worker-space/worker-o_3t2mt1,Local directory: /dask-worker-space/worker-o_3t2mt1

0,1
Comm: tls://10.4.1.246:46711,Total threads: 4
Dashboard: http://10.4.1.246:41331/status,Memory: 24.00 GiB
Nanny: tls://10.4.1.246:34275,
Local directory: /dask-worker-space/worker-9hp4_iwv,Local directory: /dask-worker-space/worker-9hp4_iwv

0,1
Comm: tls://10.4.1.47:37263,Total threads: 4
Dashboard: http://10.4.1.47:44917/status,Memory: 24.00 GiB
Nanny: tls://10.4.1.47:44489,
Local directory: /dask-worker-space/worker-43wxrajj,Local directory: /dask-worker-space/worker-43wxrajj

0,1
Comm: tls://10.4.1.41:43913,Total threads: 4
Dashboard: http://10.4.1.41:35569/status,Memory: 24.00 GiB
Nanny: tls://10.4.1.41:35629,
Local directory: /dask-worker-space/worker-ebzd4l4f,Local directory: /dask-worker-space/worker-ebzd4l4f

0,1
Comm: tls://10.4.1.147:46801,Total threads: 4
Dashboard: http://10.4.1.147:38693/status,Memory: 24.00 GiB
Nanny: tls://10.4.1.147:45955,
Local directory: /dask-worker-space/worker-q5aw5psl,Local directory: /dask-worker-space/worker-q5aw5psl

0,1
Comm: tls://10.4.1.183:37243,Total threads: 4
Dashboard: http://10.4.1.183:38877/status,Memory: 24.00 GiB
Nanny: tls://10.4.1.183:45445,
Local directory: /dask-worker-space/worker-4gsh4ib2,Local directory: /dask-worker-space/worker-4gsh4ib2

0,1
Comm: tls://10.4.1.247:46663,Total threads: 4
Dashboard: http://10.4.1.247:46655/status,Memory: 24.00 GiB
Nanny: tls://10.4.1.247:40877,
Local directory: /dask-worker-space/worker-xd18z1gj,Local directory: /dask-worker-space/worker-xd18z1gj

0,1
Comm: tls://10.4.1.43:39481,Total threads: 4
Dashboard: http://10.4.1.43:35193/status,Memory: 24.00 GiB
Nanny: tls://10.4.1.43:45919,
Local directory: /dask-worker-space/worker-l2mp69ez,Local directory: /dask-worker-space/worker-l2mp69ez

0,1
Comm: tls://10.4.1.214:32957,Total threads: 4
Dashboard: http://10.4.1.214:37895/status,Memory: 24.00 GiB
Nanny: tls://10.4.1.214:33921,
Local directory: /dask-worker-space/worker-768cceqp,Local directory: /dask-worker-space/worker-768cceqp


In [3]:
# Specify the columns we want to download
columns = [
    "interest_rate", "loan_age", "num_borrowers", 
    "borrower_credit_score", "num_units"
]

categorical = [
    "orig_channel", "occupancy_status", "property_state",
    "first_home_buyer", "loan_purpose", "property_type",
    "zip", "relocation_mortgage_indicator", "delinquency_12"
]

In [5]:
# Download data from S3
mortgage_data_all = dd.read_parquet(
    "s3://coiled-data/mortgage-2000.parq/*", 
    storage_options={"anon": True},
    columns = columns + categorical,
)

# Cache the data on Cluster workers
mortgage_data_all = mortgage_data_all.repartition(partition_size='50MB').persist()

In [6]:
mortgage_data_all

Unnamed: 0_level_0,interest_rate,loan_age,num_borrowers,borrower_credit_score,num_units,orig_channel,occupancy_status,property_state,first_home_buyer,loan_purpose,property_type,zip,relocation_mortgage_indicator,delinquency_12
npartitions=386,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
,float64,float64,float64,float64,int32,object,object,object,object,object,object,int32,object,object
,...,...,...,...,...,...,...,...,...,...,...,...,...,...
...,...,...,...,...,...,...,...,...,...,...,...,...,...,...
,...,...,...,...,...,...,...,...,...,...,...,...,...,...
,...,...,...,...,...,...,...,...,...,...,...,...,...,...


In [7]:
# inspect the first 5 entries
mortgage_data_all.head()

Unnamed: 0_level_0,interest_rate,loan_age,num_borrowers,borrower_credit_score,num_units,orig_channel,occupancy_status,property_state,first_home_buyer,loan_purpose,property_type,zip,relocation_mortgage_indicator,delinquency_12
loan_id,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
100000174660,7.875,18.0,2.0,673.0,1,B,P,MA,N,C,SF,26,N,False
100000174660,7.875,5.0,2.0,673.0,1,B,P,MA,N,C,SF,26,N,False
100000174660,7.875,17.0,2.0,673.0,1,B,P,MA,N,C,SF,26,N,False
100000174660,7.875,6.0,2.0,673.0,1,B,P,MA,N,C,SF,26,N,False
100000174660,7.875,7.0,2.0,673.0,1,B,P,MA,N,C,SF,26,N,False


### Preprocessing

In [8]:
# Cast categorical columns to the correct type
ce = Categorizer(columns=categorical)
mortgage_data_all = ce.fit_transform(mortgage_data_all)
for col in categorical:
    mortgage_data_all[col] = mortgage_data_all[col].cat.codes

In [9]:
# Create the train-test split
X, y = mortgage_data_all.iloc[:, :-1], mortgage_data_all["delinquency_12"]
X_train, X_test, y_train, y_test = train_test_split(
    X, y, test_size=0.2, shuffle=True, random_state=2
)

X_train = X_train.persist()
y_train = y_train.persist()
X_test = X_test.persist()
y_test = y_test.persist()

### Training Model

In [10]:
# Create the XGBoost DMatrix
dtrain = xgb.dask.DaskDMatrix(client, X_train, y_train)

# Set model parameters
params = {
    "max_depth": 8,
    "max_leaves": 2 ** 8,
    "gamma": 0.1,
    "eta": 0.1,
    "min_child_weight": 30,
    "objective": "binary:logistic",
    "grow_policy": "lossguide"
}


AssertionError: 

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

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

CPU times: user 224 ms, sys: 21.4 ms, total: 246 ms
Wall time: 48.2 s


48.2seconds with 5 workers, 4 cpu, 24GB memory


In [None]:
# 'booster' is the trained model
booster = output['booster']  

# 'history' is a dictionary containing evaluation metrics
history = output['history']  

### Shutting down the cluster

In [None]:
# Stop the cluster and close the client
coiled.delete_cluster(name='xgboost')
client.close()

## 4. Recap

In this notebook, we:
- trained a distributed XGBoost model on a portion of the XXX dataset using all of the cores of our machine in parallel by instantiating a Dask LocalCluster,
- expanded the distributed XGBoost model to train on the entire dataset using a Coiled Cluster of XX machines and XX total memory in the cloud.

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 own 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 at us.