主要内容

本页采用了机器翻译。点击此处可查看英文原文。

在并行池上运行 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 数据存储对象。要计算对客户产生影响的事件的持续时间,您只需使用 EventTimeRestorationTimeCustomersAffected 这三个变量即可。

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

另请参阅

函数

主题