在并行池上运行 MapReduce
启动并行池
如果您安装了 Parallel Computing Toolbox™,执行 mapreduce 可以在您的默认配置文件指定的集群上打开一个并行池,作为执行环境。
您可以设置并行设置,使池不会自动打开。在这种情况下,如果您希望 mapreduce 使用它来并行化其工作,则必须明确启动一个池。要了解有关并行设置的详细信息,请参阅 指定并行设置。
例如,此概念代码启动了一个有 12 个工作单元的池。然后它使用 mapreducer 将执行环境设置为池,从而创建 MapReducer 对象 mr。最后,它使用 mr 在转换后的数据存储 tds 上运行 mapreduce。
p = parpool('Processes',12);
mr = mapreducer(p);
outds = mapreduce(tds,@mapFun,@reduceFun,mr)注意
mapreduce可以在任何支持并行池的集群上运行。本主题中的示例使用了一个本地集群,该集群适用于所有 Parallel Computing Toolbox 安装。在集群上并行运行
mapreduce时,输出中的键值对的顺序与在 MATLAB® 中运行mapreduce不同。如果您的应用程序依赖于输出中的数据排列,则必须根据自己的要求对数据进行排序。
比较并行 MapReduce 的性能
本示例演示了如何在并行工作单元池上运行 MapReduce 计算,并比较 mapreduce 在客户端和并行工作单元池上的性能。例如,在此示例中,您将使用 mapreduce 函数计算数据集的分组均值。
首先,在 MATLAB ® 客户端会话上运行 mapreduce 函数,然后在本地集群上并行运行该函数。您可以使用 mapreducer 函数来显式控制执行环境。
首先,使用“本地进程”配置文件启动一个包含四个工作单元的并行池。
p = parpool("Processes",4);创建两个 MapReducer 对象,用于为 mapreduce 指定不同的执行环境。
首先,在客户端创建一个用于串行执行的 SerialMapReducer 对象。
onClient = mapreducer(0);
然后,在已打开的并行池上创建一个用于并行执行的 ParallelMapReducer 对象。
onPool = mapreducer(p);
准备样本数据集
示例数据集文件 energyGridEvents.parquet 包含代表大型能源电网中每个变电站和馈线事件的数据。使用数据存储来预览文件中的数据。该数据共有 10 列。
filename = "energyGridEvents.parquet";
preview(parquetDatastore(filename))ans=8×10 table
01-Jan-2000 00:33:59 "North" "S173" "F1057" "EquipmentFault" 4 0 01-Jan-2000 00:59:59 "Mitigated" "Equipment"
01-Jan-2000 00:39:19 "West" "S270" "F0022" "FrequencyDeviation" 1 876 01-Jan-2000 00:47:19 "Resolved" "Equipment"
01-Jan-2000 00:46:13 "East" "S284" "F1123" "Outage" 5 115885 04-Jan-2000 16:31:13 "Mitigated" "Equipment"
01-Jan-2000 02:13:30 "East" "S089" "F0044" "StormImpact" 4 53753 02-Jan-2000 22:33:30 "Open" "Weather"
01-Jan-2000 03:05:16 "Central" "S254" "F1738" "VoltageDip" 2 1931 01-Jan-2000 03:16:16 "Resolved" "HumanError"
01-Jan-2000 03:10:49 "Central" "S371" "F0593" "Outage" 3 25535 02-Jan-2000 03:57:49 "Mitigated" "Weather"
01-Jan-2000 03:40:03 "North" "S102" "F1322" "EquipmentFault" 4 0 01-Jan-2000 05:16:03 "Resolved" "Unknown"
01-Jan-2000 07:23:46 "East" "S272" "F1723" "StormImpact" 4 87741 02-Jan-2000 18:38:46 "Resolved" "Weather"
energyGridEvents.parquet 文件是一个相对较小的数据集,因此与 MATLAB 客户端会话相比,使用并行池时不太可能观察到性能的提升。为了增加数据集的大小,您可以创建该文件的多个副本。
该代码会在当前文件夹中创建一个子文件夹,并生成 20 个 energyGridEvents.parquet 文件的副本,这些文件大约占用 100 MB 的磁盘空间。若要确认您确实要增加数据集的大小,请在运行示例之前从下拉列表中选择 "true"。若要使用单份 energyGridEvents.parquet 文件运行该示例,请从下拉列表中选择 "false"。
increaseDatasetIfTrue =select; if increaseDatasetIfTrue if ~isfolder("dataFolder") mkdir("dataFolder") end for c = 1:50 filepath = fullfile("dataFolder", ... strcat("events_",num2str(c),".parquet")); copyfile(filename,filepath); end dataRoot = "dataFolder"; else dataRoot = filename; end
在更大规模的数据集上创建数据存储
使用较大的事件数据集创建一个 ParquetDatastore 数据存储对象。要计算对客户产生影响的事件的持续时间,您只需使用 EventTime、RestorationTime 和 CustomersAffected 这三个变量即可。
selectedVariables = ["EventTime","RestorationTime","CustomersAffected"]; ds = parquetDatastore(dataRoot, ... SelectedVariableNames=selectedVariables);
使用 ParquetDatastore 对象创建一个行过滤器。然后,使用行过滤器过滤出 CustomerAffected 值大于 0 的行。
rf = rowfilter(ds); ds.RowFilter = rf.CustomersAffected>0;
预览数据存储。
preview(ds)
ans=8×3 table
01-Jan-2000 00:39:19 01-Jan-2000 00:47:19 876
01-Jan-2000 00:46:13 04-Jan-2000 16:31:13 115885
01-Jan-2000 02:13:30 02-Jan-2000 22:33:30 53753
01-Jan-2000 03:05:16 01-Jan-2000 03:16:16 1931
01-Jan-2000 03:10:49 02-Jan-2000 03:57:49 25535
01-Jan-2000 07:23:46 02-Jan-2000 18:38:46 87741
01-Jan-2000 12:57:06 02-Jan-2000 14:56:06 16922
01-Jan-2000 13:23:56 02-Jan-2000 15:52:56 6694
在客户端上运行 MapReduce 操作
运行 MapReduce 计算,按事件开始的星期几进行分组,计算事件平均持续时间(单位:小时)。要在 MATLAB 客户端会话上运行 mapreduce 函数,请将 SerialMapReducer 对象 onClient 指定为执行环境。该示例的 map 和 reduce 函数定义在示例末尾。
serialMeanEventDurationByDay = mapreduce(ds, ... @meanEventDurationMapper,... @meanEventDurationReducer,onClient);
******************************** * MAPREDUCE PROGRESS * ******************************** Map 0% Reduce 0% Map 10% Reduce 0% Map 20% Reduce 0% Map 30% Reduce 0% Map 40% Reduce 0% Map 50% Reduce 0% Map 60% Reduce 0% Map 70% Reduce 0% Map 80% Reduce 0% Map 90% Reduce 0% Map 100% Reduce 0% Map 100% Reduce 14% Map 100% Reduce 29% Map 100% Reduce 43% Map 100% Reduce 57% Map 100% Reduce 71% Map 100% Reduce 86% Map 100% Reduce 100%
mapreduce 函数返回一个 KeyValueDatastore 对象 serialMeanEventDurationByDay,该对象指向它在当前文件夹中创建的结果文件。
从输出数据存储库 serialMeanEventDurationByDay 中读取最终结果。
readall(serialMeanEventDurationByDay)
ans=7×2 table
'Sunday' 23.3503
'Monday' 15.7961
'Tuesday' 10.5672
'Wednesday' 10.6704
'Thursday' 10.4908
'Friday' 15.7035
'Saturday' 23.1165
在池中运行 mapreduce 计算
接下来,通过在 mapreduce 函数中指定 ParallelMapReducer 对象 onPool,在已打开的并行池上运行 MapReduce 计算。该示例的 map 和 reduce 函数定义在示例末尾。请注意,显示的文本表明 mapreduce 函数正在并行运行。
parallelMeanEventDurationByDay = mapreduce(ds, ... @meanEventDurationMapper, ... @meanEventDurationReducer,onPool);
Parallel mapreduce execution on the parallel pool: ******************************** * MAPREDUCE PROGRESS * ******************************** Map 0% Reduce 0% Map 1% Reduce 0% Map 2% Reduce 0% Map 3% Reduce 0% Map 4% Reduce 0% Map 5% Reduce 0% Map 6% Reduce 0% Map 7% Reduce 0% Map 8% Reduce 0% Map 9% Reduce 0% Map 10% Reduce 0% Map 11% Reduce 0% Map 12% Reduce 0% Map 13% Reduce 0% Map 14% Reduce 0% Map 15% Reduce 0% Map 16% Reduce 0% Map 17% Reduce 0% Map 18% Reduce 0% Map 19% Reduce 0% Map 20% Reduce 0% Map 21% Reduce 0% Map 22% Reduce 0% Map 23% Reduce 0% Map 24% Reduce 0% Map 25% Reduce 0% Map 26% Reduce 0% Map 27% Reduce 0% Map 28% Reduce 0% Map 29% Reduce 0% Map 30% Reduce 0% Map 31% Reduce 0% Map 32% Reduce 0% Map 33% Reduce 0% Map 34% Reduce 0% Map 35% Reduce 0% Map 36% Reduce 0% Map 37% Reduce 0% Map 38% Reduce 0% Map 39% Reduce 0% Map 40% Reduce 0% Map 41% Reduce 0% Map 42% Reduce 0% Map 43% Reduce 0% Map 44% Reduce 0% Map 45% Reduce 0% Map 46% Reduce 0% Map 47% Reduce 0% Map 48% Reduce 0% Map 49% Reduce 0% Map 50% Reduce 0% Map 51% Reduce 0% Map 52% Reduce 0% Map 53% Reduce 0% Map 54% Reduce 0% Map 55% Reduce 0% Map 56% Reduce 0% Map 57% Reduce 0% Map 58% Reduce 0% Map 59% Reduce 0% Map 60% Reduce 0% Map 61% Reduce 0% Map 62% Reduce 0% Map 63% Reduce 0% Map 64% Reduce 0% Map 65% Reduce 0% Map 66% Reduce 0% Map 67% Reduce 0% Map 68% Reduce 0% Map 69% Reduce 0% Map 70% Reduce 0% Map 71% Reduce 0% Map 72% Reduce 0% Map 73% Reduce 0% Map 74% Reduce 0% Map 75% Reduce 0% Map 76% Reduce 0% Map 77% Reduce 0% Map 78% Reduce 0% Map 79% Reduce 0% Map 80% Reduce 0% Map 81% Reduce 0% Map 82% Reduce 0% Map 83% Reduce 0% Map 84% Reduce 0% Map 85% Reduce 0% Map 86% Reduce 0% Map 87% Reduce 0% Map 88% Reduce 0% Map 89% Reduce 0% Map 90% Reduce 0% Map 91% Reduce 0% Map 92% Reduce 0% Map 93% Reduce 0% Map 94% Reduce 0% Map 95% Reduce 0% Map 96% Reduce 0% Map 97% Reduce 0% Map 98% Reduce 0% Map 99% Reduce 0% Map 100% Reduce 0% Map 100% Reduce 100%
mapreduce 函数返回一个 KeyValueDatastore 对象 parallelMeanEventDurationByDay,该对象指向它在当前文件夹中创建的四个结果文件。每个结果文件对应着每个工作单元计算出的结果。
从输出数据存储库 parallelMeanEventDurationByDay 中读取最终结果。
readall(parallelMeanEventDurationByDay)
ans=7×2 table
'Friday' 15.7035
'Monday' 15.7961
'Wednesday' 10.6704
'Sunday' 23.3503
'Thursday' 10.4908
'Saturday' 23.1165
'Tuesday' 10.5672
比较执行时间
使用 timeit 函数,测量在 MATLAB 客户端和并行池上执行 MapReduce 计算所需的时间。timeit 函数需要几分钟才能完成。
tClient = timeit(@() ... mapreduce(ds,@meanEventDurationMapper, ... @meanEventDurationReducer,onClient,Display="off")); tPool = timeit(@() ... mapreduce(ds,@meanEventDurationMapper, ... @meanEventDurationReducer,onPool,Display="off"));
比较在 MATLAB 客户端和并行池上执行 MapReduce 计算的情况。
disp("Speedup of mapreduce computation on" + ... " a parallel pool compared to the MATLAB client: " + ... round(tClient/tPool) + "x")
Speedup of mapreduce computation on a parallel pool compared to the MATLAB client: 3x
figure executionEnvironment = ["Client" "Pool"]; bar(executionEnvironment,[tClient tPool]) xlabel("Execution Environment") ylabel("MapReduce Computation Time (s)")

清理文件
删除当前文件夹中的事件数据集和 mapreduce 结果 MAT 文件。
rmdir("dataFolder","s") delete("result*")
支持函数
map 函数 meanEventDurationMapper 根据事件的开始日期对事件进行分组,移除所有包含缺失值的行,并计算每天的事件数量以及事件持续时间的总和(单位为小时)。该函数将结果作为中间键值对添加到 KeyValueStore 对象中。其中,键是星期几,例如 "Sunday"",值是一个包含两个元素的向量,其中包含该天事件持续时间的计数和总和。
function meanEventDurationMapper(data,~,intermKVStore) eventTime = data.EventTime; % Calculate event duration in hours eventDuration = data.RestorationTime - data.EventTime; eventDuration = hours(eventDuration); % Remove rows with missing values notNaN = ~isnan(eventDuration); eventTime = eventTime(notNaN); eventDuration = eventDuration(notNaN); % Compute the count and sum of event duration per event start day [s,dayOfWeek,n] = groupsummary(eventDuration,eventTime,"dayname","sum"); dayOfWeek = string(dayOfWeek); sumAndCount = num2cell([n,s],2); addmulti(intermKVStore,dayOfWeek,sumAndCount); end
reduce 函数 meanEventDurationReducer 接收一个列表,其中包含由输入键 intermKey 指定的当天各延迟的中间计数和总和,并将这些值汇总为总计数和总和。随后,该函数计算总体均值,并向输出对象 KeyValueStore 中添加最后一组键值对。该键值对表示该星期几的事件平均持续时间的值。
function meanEventDurationReducer(intermKey,intermValIter,outKVStore) totalCount = 0; totalSum = 0; % Accumulate intermediate results for one day while hasnext(intermValIter) intermValue = getnext(intermValIter); totalCount = totalCount + intermValue(1); totalSum = totalSum + intermValue(2); end % Add results to the output datastore meanDuration = totalSum/totalCount; add(outKVStore,intermKey,meanDuration); end
