GPU上的随机森林:比Apache Spark快2000倍

作者|Aaron Richter 编译|VK 来源|Towards Data Science

随机森林是一种机器学习算法,以其鲁棒性、准确性和可扩展性而受到许多数据科学家的信赖。

该算法通过bootstrap聚合训练出多棵决策树,然后通过集成对输出进行预测。由于其集成特征的特点,随机森林是一种可以在分布式计算环境中实现的算法。树可以在集群中跨进程和机器并行训练,结果比使用单个进程的训练时间快得多。

在本文中,我们探索了使用Apache Spark在CPU机器集群上实现分布式随机森林训练,并将其与使用NVIDIA RAPIDS和Dask的GPU机器集群上的训练性能进行了比较。

虽然GPU计算传统上是为深度学习应用而保留的,但RAPIDS是一个在GPU上执行数据处理和非深度学习ML工作的库,与在cpu上执行相比,它可以大大提高性能。

我们使用3亿个实例训练了一个随机森林模型:Spark在20个节点CPU集群上耗时37分钟,而RAPIDS在20个节点GPU集群上耗时1秒。GPU的速度提高了2000倍以上!

实验概述

我们使用公共可用的纽约出租车数据集,并训练一个随机森林回归器,该回归器可以使用与乘客接送相关的属性来预测出租车的票价金额。以2017年、2018年和2019年的出租车出行量为训练集,共计300700143个实例。

数据集链接:https://www1.nyc.gov/site/tlc/about/tlc-trip-record-data.page

Spark和RAPIDS代码可以在Jupyter Notebook中找到。

硬件

Spark集群使用Amazon EMR进行管理,而Dask/RAPIDS集群则使用Saturn Cloud进行管理。

两个集群都有20个工作节点,具有以下AWS实例类型:

Spark:r5.2xlarge

  • 8个CPU,64 GB RAM

  • 按需价格:0.504美元/小时

RAPIDS:g4dn.xlarge

  • 4个CPU,16 GB RAM

  • 1个GPU,16 GB GPU RAM(NVIDIA T4)

  • 按需价格:0.526美元/小时

Saturn Cloud也可以用NVIDIA特斯拉V100 GPU来启动Dask集群,但我们在这个练习中选择了g4dn.xlarge,保持与Spark集群相似的小时成本概况。

Spark

Apache Spark是一个在Scala中构建的开源大数据处理引擎,它有一个Python接口,可以调用Scala/JVM代码。

它是Hadoop处理生态系统中的一个重要组成部分,围绕MapReduce范例构建,并且具有用于数据帧和机器学习的接口。

设置Spark集群不在本文的讨论范围之内,但是一旦准备好集群,就可以在Jupyter Notebook中运行以下命令来初始化Spark:

1import findspark 2findspark.init() 3 4from pyspark.sql import SparkSession 5 6spark = (SparkSession 7 .builder 8 .config('spark.executor.memory', '36g') 9 .getOrCreate())

findspark包检测系统上的Spark安装位置;如果可以知道Spark包的安装位置,则可能不需要这样做。

要获得有性能的Spark代码,需要设置几个配置设置,这取决于集群设置和工作流。在这种情况下,我们设置spark.executor.memory以确保我们不会遇到任何内存溢出或Java堆错误。

RAPIDS

NVIDIA RAPIDS是一个开源的Python框架,它在gpu而不是cpu上执行数据科学代码。类似于在训练深度学习模型时所看到的,这将为数据科学工作带来巨大的性能提升。

RAPIDS有数据帧、ML、图形分析等接口。RAPIDS使用Dask来处理与具有多个gpu的机器的并行化,以及每个具有一个或多个gpu的机器集群。

设置GPU机器可能有点棘手,但是Saturn Cloud已经为启动GPU集群预构建了映像,所以你只需几分钟就可以启动并运行了!要初始化指向群集的Dask客户端,可以运行以下命令:

1from dask.distributed import Client 2from dask_saturn import SaturnCluster 3 4cluster = SaturnCluster() 5client = Client(cluster)

要自己设置Dask集群,请参阅此docs页面:https://docs.dask.org/en/latest/setup.html

数据加载

数据文件托管在一个公共的S3 bucket上,因此我们可以直接从那里读取csv。S3 bucket的所有文件都在同一个目录中,所以我们使用s3fs来选择我们想要的文件:

1import s3fs 2fs = s3fs.S3FileSystem(anon=True) 3files = [f"s3://{x}" for x in fs.ls('s3://nyc-tlc/trip data/') 4 if 'yellow' in x and ('2019' in x or '2018' in x or '2017' in x)] 5 6cols = ['VendorID', 'tpep_pickup_datetime', 'tpep_dropoff_datetime', 'passenger_count', 'trip_distance', 7 'RatecodeID', 'store_and_fwd_flag', 'PULocationID', 'DOLocationID', 'payment_type', 'fare_amount', 8 'extra', 'mta_tax', 'tip_amount', 'tolls_amount', 'improvement_surcharge', 'total_amount']

使用Spark,我们需要单独读取每个CSV文件,然后将它们组合在一起:

1import functools 2from pyspark.sql.types import * 3import pyspark.sql.functions as F 4from pyspark.sql import DataFrame 5 6# 手动指定模式,因为read.csv中的inferSchema非常慢 7schema = StructType([ 8 StructField('VendorID', DoubleType()), 9 StructField('tpep_pickup_datetime', TimestampType()), 10 ... 11 # 参考notebook获得完整对象模式 12]) 13 14def read_csv(path): 15 df = spark.read.csv(path, 16 header=True, 17 schema=schema, 18 timestampFormat='yyyy-MM-dd HH:mm:ss', 19 ) 20 df = df.select(cols) 21 return df 22 23dfs = [] 24for tf in files: 25 df = read_csv(tf) 26 dfs.append(df) 27 28taxi = functools.reduce(DataFrame.unionAll, dfs) 29taxi.count()

使用Dask+RAPIDS,我们可以一次性读取所有CSV文件:

1import dask_cudf 2 3taxi = dask_cudf.read_csv(files, 4 assume_missing=True, 5 parse_dates=[1,2], 6 usecols=cols, 7 storage_options={'anon': True}) 8len(taxi)

特征工程

我们将根据时间生成一些特征,然后保存数据帧。在这两个框架中,这将执行所有CSV加载和预处理,并将结果存储在RAM中(在RAPIDS的情况下是GPU RAM)。我们将用于训练的特征包括:

1features = ['pickup_weekday', 'pickup_hour', 'pickup_minute', 2 'pickup_week_hour', 'passenger_count', 'VendorID', 3 'RatecodeID', 'store_and_fwd_flag', 'PULocationID', 4 'DOLocationID']

对于Spark,我们需要将特征收集到向量类中:

1from pyspark.ml.feature import VectorAssembler 2from pyspark.ml.pipeline import Pipeline 3 4taxi = taxi.withColumn('pickup_weekday', F.dayofweek(taxi.tpep_pickup_datetime).cast(DoubleType())) 5 6taxi = taxi.withColumn('pickup_hour', F.hour(taxi.tpep_pickup_datetime).cast(DoubleType())) 7 8taxi = taxi.withColumn('pickup_minute', F.minute(taxi.tpep_pickup_datetime).cast(DoubleType())) 9 10taxi = taxi.withColumn('pickup_week_hour', ((taxi.pickup_weekday * 24) + taxi.pickup_hour).cast(DoubleType())) 11 12taxi = taxi.withColumn('store_and_fwd_flag', F.when(taxi.store_and_fwd_flag == 'Y', 1).otherwise(0)) 13 14taxi = taxi.withColumn('label', taxi.total_amount) 15taxi = taxi.fillna(-1) 16 17assembler = VectorAssembler( 18 inputCols=features, 19 outputCol='features', 20) 21 22pipeline = Pipeline(stages=[assembler]) 23assembler_fitted = pipeline.fit(taxi) 24X = assembler_fitted.transform(taxi) 25X.cache() 26X.count()

对于RAPIDS,我们将所有浮点值转换为float32,以便进行GPU计算:

1from dask import persist 2from dask.distributed import wait 3 4taxi['pickup_weekday'] = taxi.tpep_pickup_datetime.dt.weekday 5taxi['pickup_hour'] = taxi.tpep_pickup_datetime.dt.hour 6taxi['pickup_minute'] = taxi.tpep_pickup_datetime.dt.minute 7taxi['pickup_week_hour'] = (taxi.pickup_weekday * 24) + taxi.pickup_hour 8taxi['store_and_fwd_flag'] = (taxi.store_and_fwd_flag == 'Y').astype(float) 9taxi = taxi.fillna(-1) 10 11X = taxi[features].astype('float32') 12y = taxi['total_amount'] 13X, y = persist(X, y) 14_ = wait([X, y]) 15len(X)

训练随机森林

我们只需要几行代码就可以训练随机森林。

Spark:

1from pyspark.ml.regression import RandomForestRegressor 2rf = RandomForestRegressor(numTrees=100, maxDepth=10, seed=42) 3fitted = rf.fit(X)

RAPIDS:

1from cuml.dask.ensemble import RandomForestRegressor 2rf = RandomForestRegressor(n_estimators=100, max_depth=10, seed=42) 3_ = rf.fit(X, y)

结果

我们对Spark(CPU)和RAPIDS(GPU)集群上的300700143个纽约出租车数据实例训练了一个随机森林模型。两个集群都有20个工作节点,每小时价格大致相同。以下是工作流每个部分的结果:

Task

Spark

RAPIDS

Load/rowcount

20.6 seconds

25.5 seconds

Feature engineering

54.3 seconds

23.1 seconds

Random forest

36.9 minutes

1.02 seconds

37分钟的_Spark_ 与1秒的_RAPIDS_!

GPU胜利!想一想,一次拟合你不需要等待37分钟了,这将加快之后迭代和改进模型的速度。而在CPU上,一旦添加了超参数调优或测试不同的模型,迭代都很容易累积到数小时或数天。

你需要看到才能相信吗?你可以在这里找到Notebook,然后自己运行测试:https://github.com/saturncloud/saturn-cloud-examples/tree/main/machine_learning/random_forest

你需要更快的随机森林吗

对!你可以在几秒钟内用Saturn Cloud进入Dask/RAPIDS。Saturn处理所有工具基础设施、安全性和部署方面的难题,让你立即启动并运行RAPIDS。点击这里在你的AWS帐户免费试用Saturn:https://manager.aws.saturnenterprise.io/register

原文链接:https://towardsdatascience.com/random-forest-on-gpus-2000x-faster-than-apache-spark-9561f13b00ae

欢迎关注磐创AI博客站: http://panchuang.net/

sklearn机器学习中文官方文档: http://sklearn123.com/

欢迎关注磐创博客资源汇总站: http://docs.panchuang.net/

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )