Mars 是一个并行和分布式 Python 框架,能轻松把单机大家耳熟能详的的 numpy、pandas、scikit-learn 等库,以及 Python 函数利用多核或者多机加速。这其中,并行和分布式 Python 函数主要利用 Mars Remote API。
启动 Mars 分布式环境可以参考:
- 命令行方式在集群中部署。
- Kubernetes 中部署。
- MaxCompute 开箱即用的环境,购买了 MaxCompute 服务的可以直接使用。
如何使用 Mars Remote API
使用 Mars Remote API 非常简单,只需要对原有的代码做少许改动,就可以分布式执行。
采用蒙特卡洛方法计算 π 为例。代码如下,我们编写了两个函数,calc_chunk 用来计算每个分片内落在圆内的点的个数,calc_pi 用来把多个分片 calc_chunk 计算的结果汇总最后得出 π 值。
1from typing import List 2import numpy as np 3 4def calc_chunk(n: int, i: int): 5 # 计算n个随机点(x和y轴落在-1到1之间)到原点距离小于1的点的个数 6 rs = np.random.RandomState(i) 7 a = rs.uniform(-1, 1, size=(n, 2)) 8 d = np.linalg.norm(a, axis=1) 9 return (d < 1).sum() 10 11def calc_pi(fs: List[int], N: int): 12 # 将若干次 calc_chunk 计算的结果汇总,计算 pi 的值 13 return sum(fs) * 4 / N 14 15N = 200_000_000 16n = 10_000_000 17 18fs = [calc_chunk(n, i) 19 for i in range(N // n)] 20pi = calc_pi(fs, N) 21print(pi)
%%time 下可以看到结果:
13.1416312 2CPU times: user 9.47 s, sys: 2.62 s, total: 12.1 s 3Wall time: 12.3 s
在单机需要 12.3 s。
要让这个计算使用 Mars Remote API 并行起来,我们不需要对函数做任何改动,需要变动的仅仅是最后部分。
1import mars.remote as mr 2 3# 函数调用改成 mars.remote.spawn 4fs = [mr.spawn(calc_chunk, args=(n, i)) 5 for i in range(N // n)] 6# 把 spawn 的列表传入作为参数,再 spawn 新的函数 7pi = mr.spawn(calc_pi, args=(fs, N)) 8# 通过 execute() 触发执行,fetch() 获取结果 9print(pi.execute().fetch())
%%time 下看到结果:
13.1416312 2CPU times: user 29.6 ms, sys: 4.23 ms, total: 33.8 ms 3Wall time: 2.85 s
结果一模一样,但是却有数倍的性能提升。
可以看到,对已有的 Python 代码,Mars remote API 几乎不需要做多少改动,就能有效并行和分布式来加速执行过程。
一个例子
为了让读者理解 Mars Remote API 的作用,我们从另一个例子开始。现在我们有一个数据集,我们希望对它们做一个分类任务。要做分类,我们有很多算法和库可以选择,这里我们用 RandomForest、LogisticRegression,以及 XGBoost。
困难的地方是,除了有多个模型选择,这些模型也会包含多个超参,那哪个超参效果最好呢?对于调参不那么有经验的同学,跑过了才知道。所以,我们希望能生成一堆可选的超参,然后把他们都跑一遍,看看效果。
准备数据
这个例子里我们使用 otto 数据集。
首先,我们准备数据。读取数据后,我们按 2:1 的比例把数据分成训练集和测试集。
1import pandas as pd 2from sklearn.preprocessing import LabelEncoder 3from sklearn.model_selection import train_test_split 4 5def gen_data(): 6 df = pd.read_csv('otto/train.csv') 7 8 X = df.drop(['target', 'id'], axis=1) 9 y = df['target'] 10 11 label_encoder = LabelEncoder() 12 label_encoder.fit(y) 13 y = label_encoder.transform(y) 14 15 return train_test_split(X, y, test_size=0.33, random_state=123) 16 17X_train, X_test, y_train, y_test = gen_data()
模型
接着,我们使用 scikit-learn 的 RandomForest 和 LogisticRegression 来处理分类。
RandomForest:
1from sklearn.ensemble import RandomForestClassifier 2 3def random_forest(X_train: pd.DataFrame, 4 y_train: pd.Series, 5 verbose: bool = False, 6 **kw): 7 model = RandomForestClassifier(verbose=verbose, **kw) 8 model.fit(X_train, y_train) 9 return model
接着,我们生成供 RandomForest 使用的超参,我们用 yield 的方式来迭代返回。
1def gen_random_forest_parameters(): 2 for n_estimators in [50, 100, 600]: 3 for max_depth in [None, 3, 15]: 4 for criterion in ['gini', 'entropy']: 5 yield { 6 'n_estimators': n_estimators, 7 'max_depth': max_depth, 8 'criterion': criterion 9 }
LogisticRegression 也是这个过程。我们先定义模型。
1from sklearn.linear_model import LogisticRegression 2 3def logistic_regression(X_train: pd.DataFrame, 4 y_train: pd.Series, 5 verbose: bool = False, 6 **kw): 7 model = LogisticRegression(verbose=verbose, **kw) 8 model.fit(X_train, y_train) 9 return model
接着生成供 LogisticRegression 使用的超参。
1def gen_lr_parameters(): 2 for penalty in ['l2', 'none']: 3 for tol in [0.1, 0.01, 1e-4]: 4 yield { 5 'penalty': penalty, 6 'tol': tol 7 }
XGBoost 也是一样,我们用 XGBClassifier 来执行分类任务。
1from xgboost import XGBClassifier 2 3def xgb(X_train: pd.DataFrame, 4 y_train: pd.Series, 5 verbose: bool = False, 6 **kw): 7 model = XGBClassifier(verbosity=int(verbose), **kw) 8 model.fit(X_train, y_train) 9 return model
生成一系列超参。
1def gen_xgb_parameters(): 2 for n_estimators in [100, 600]: 3 for criterion in ['gini', 'entropy']: 4 for learning_rate in [0.001, 0.1, 0.5]: 5 yield { 6 'n_estimators': n_estimators, 7 'criterion': criterion, 8 'learning_rate': learning_rate 9 }
验证
接着我们编写验证逻辑,这里我们使用 log_loss 来作为评价函数。
1from sklearn.metrics import log_loss 2 3def metric_model(model, 4 X_test: pd.DataFrame, 5 y_test: pd.Series) -> float: 6 if isinstance(model, bytes): 7 model = pickle.loads(model) 8 y_pred = model.predict_proba(X_test) 9 return log_loss(y_test, y_pred) 10 11 12def train_and_metric(train_func, 13 train_params: dict, 14 X_train: pd.DataFrame, 15 y_train: pd.Series, 16 X_test: pd.DataFrame, 17 y_test: pd.Series, 18 verbose: bool = False 19 ): 20 # 把训练和验证封装到一起 21 model = train_func(X_train, y_train, verbose=verbose, **train_params) 22 metric = metric_model(model, X_test, y_test) 23 return model, metric
找出最好的模型
做好准备工作后,我们就开始来跑模型了。针对每个模型,我们把每次生成的超参们送进去训练,除了这些超参,我们还把 n_jobs 设成 -1,这样能更好利用单机的多核。
1results = [] 2 3# ------------- 4# Random Forest 5# ------------- 6 7for params in gen_random_forest_parameters(): 8 print(f'calculating on {params}') 9 # fixed random_state 10 params['random_state'] = 123 11 # use all CPU cores 12 params['n_jobs'] = -1 13 model, metric = train_and_metric(random_forest, params, 14 X_train, y_train, 15 X_test, y_test) 16 print(f'metric: {metric}') 17 results.append({'model': model, 18 'metric': metric}) 19 20# ------------------- 21# Logistic Regression 22# ------------------- 23 24for params in gen_lr_parameters(): 25 print(f'calculating on {params}') 26 # fixed random_state 27 params['random_state'] = 123 28 # use all CPU cores 29 params['n_jobs'] = -1 30 model, metric = train_and_metric(logistic_regression, params, 31 X_train, y_train, 32 X_test, y_test) 33 print(f'metric: {metric}') 34 results.append({'model': model, 35 'metric': metric}) 36 37# ------- 38# XGBoost 39# ------- 40 41for params in gen_xgb_parameters(): 42 print(f'calculating on {params}') 43 # fixed random_state 44 params['random_state'] = 123 45 # use all CPU cores 46 params['n_jobs'] = -1 47 model, metric = train_and_metric(xgb, params, 48 X_train, y_train, 49 X_test, y_test) 50 print(f'metric: {metric}') 51 results.append({'model': model, 52 'metric': metric})
运行一下,需要相当长时间,我们省略掉一部分输出内容。
1calculating on {'n_estimators': 50, 'max_depth': None, 'criterion': 'gini'} 2metric: 0.6964123781828575 3calculating on {'n_estimators': 50, 'max_depth': None, 'criterion': 'entropy'} 4metric: 0.6912312790832288 5# 省略其他模型的输出结果 6CPU times: user 3h 41min 53s, sys: 2min 34s, total: 3h 44min 28s 7Wall time: 31min 44s
从 CPU 时间和 Wall 时间,能看出来这些训练还是充分利用了多核的性能。但整个过程还是花费了 31 分钟。
使用 Remote API 分布式加速
现在我们尝试使用 Remote API 通过分布式方式加速整个过程。
集群方面,我们使用最开始说的第三种方式,直接在 MaxCompute 上拉起一个集群。大家可以选择其他方式,效果是一样的。
1n_cores = 8 2mem = 2 * n_cores # 16G 3# o 是 MaxCompute 入口,这里创建 10 个 worker 的集群,每个 worker 8核16G 4cluster = o.create_mars_cluster(10, n_cores, mem, image='extended')
为了方便在分布式读取数据,我们对数据处理稍作改动,把数据上传到 MaxCompute 资源。对于其他环境,用户可以考虑 HDFS、Aliyun OSS 或者 Amazon S3 等存储。
1if not o.exist_resource('otto_train.csv'): 2 with open('otto/train.csv') as f: 3 # 上传资源 4 o.create_resource('otto_train.csv', 'file', fileobj=f) 5 6def gen_data(): 7 # 改成从资源读取 8 df = pd.read_csv(o.open_resource('otto_train.csv')) 9 10 X = df.drop(['target', 'id'], axis=1) 11 y = df['target'] 12 13 label_encoder = LabelEncoder() 14 label_encoder.fit(y) 15 y = label_encoder.transform(y) 16 17 return train_test_split(X, y, test_size=0.33, random_state=123)
稍作改动之后,我们使用 mars.remote.spawn 方法来让 gen_data 调度到集群上运行。
1import mars.remote as mr 2 3# n_output 说明是 4 输出 4# execute() 执行后,数据会读取到 Mars 集群内部 5data = mr.ExecutableTuple(mr.spawn(gen_data, n_output=4)).execute() 6# remote_ 开头的都是 Mars 对象,这时候数据在集群内,这些对象只是引用 7remote_X_train, remote_X_test, remote_y_train, remote_y_test = data
目前 Mars 能正确序列化 numpy ndarray、pandas DataFrame 等,还不能序列化模型,所以,我们要对 train_and_metric 稍作改动,把模型 pickle 了之后再返回。
1def distributed_train_and_metric(train_func, 2 train_params: dict, 3 X_train: pd.DataFrame, 4 y_train: pd.Series, 5 X_test: pd.DataFrame, 6 y_test: pd.Series, 7 verbose: bool = False 8 ): 9 model, metric = train_and_metric(train_func, train_params, 10 X_train, y_train, 11 X_test, y_test, verbose=verbose) 12 return pickle.dumps(model), metric
后续 Mars 支持了序列化模型后可以直接 spawn 原本的函数。
接着我们就对前面的执行过程稍作改动,把函数调用全部都用 mars.remote.spawn 来改写。
1import numpy as np 2 3tasks = [] 4models = [] 5metrics = [] 6 7# ------------- 8# Random Forest 9# ------------- 10 11for params in gen_random_forest_parameters(): 12 # fixed random_state 13 params['random_state'] = 123 14 task = mr.spawn(distributed_train_and_metric, 15 args=(random_forest, params, 16 remote_X_train, remote_y_train, 17 remote_X_test, remote_y_test), 18 kwargs={'verbose': 2}, 19 n_output=2 20 ) 21 tasks.extend(task) 22 # 把模型和评价分别存储 23 models.append(task[0]) 24 metrics.append(task[1]) 25 26 27# ------------------- 28# Logistic Regression 29# ------------------- 30 31for params in gen_lr_parameters(): 32 # fixed random_state 33 params['random_state'] = 123 34 task = mr.spawn(distributed_train_and_metric, 35 args=(logistic_regression, params, 36 remote_X_train, remote_y_train, 37 remote_X_test, remote_y_test), 38 kwargs={'verbose': 2}, 39 n_output=2 40 ) 41 tasks.extend(task) 42 # 把模型和评价分别存储 43 models.append(task[0]) 44 metrics.append(task[1]) 45 46# ------- 47# XGBoost 48# ------- 49 50for params in gen_xgb_parameters(): 51 # fixed random_state 52 params['random_state'] = 123 53 # 再指定并发为核的个数 54 params['n_jobs'] = n_cores 55 task = mr.spawn(distributed_train_and_metric, 56 args=(xgb, params, 57 remote_X_train, remote_y_train, 58 remote_X_test, remote_y_test), 59 kwargs={'verbose': 2}, 60 n_output=2 61 ) 62 tasks.extend(task) 63 # 把模型和评价分别存储 64 models.append(task[0]) 65 metrics.append(task[1]) 66 67 68# 把顺序打乱,目的是能分散到 worker 上平均一点 69shuffled_tasks = np.random.permutation(tasks) 70_ = mr.ExecutableTuple(shuffled_tasks).execute()
可以看到代码几乎一致。
运行查看结果:
1CPU times: user 69.1 ms, sys: 10.9 ms, total: 80 ms 2Wall time: 1min 59s
时间一下子从 31 分钟多来到了 2 分钟,提升 15x+。但代码修改的代价可以忽略不计。
细心的读者可能注意到了,分布式运行的代码中,我们把模型的 verbose 给打开了,在分布式环境下,因为这些函数远程执行,打印的内容只会输出到 worker 的标准输出流,我们在客户端不会看到打印的结果,但 Mars 提供了一个非常有用的接口来让我们查看每个模型运行时的输出。
以第0个模型为例,我们可以在 Mars 对象上直接调用 fetch_log 方法。
print(models[0].fetch_log())
输出我们简略一部分。
1building tree 1 of 50 2building tree 2 of 50 3building tree 3 of 50 4building tree 4 of 50 5building tree 5 of 50 6building tree 6 of 50 7# 中间省略 8building tree 49 of 50 9building tree 50 of 50
要看哪个模型都可以通过这种方式。试想下,如果没有 fetch_log API,你确想看中间过程的输出有多麻烦。首先这个函数在哪个 worker 上执行,不得而知;然后,即便知道是哪个 worker,因为每个 worker 上可能有多个函数执行,这些输出就可能混杂在一起,甚至被庞大日志淹没了。fetch_log 接口让用户不需要关心在哪个 worker 上执行,也不用担心日志混合在一起。
本文为阿里云原创内容,未经允许不得转载。