在 Hadoop 集群上运行 mapreduce
集群准备
在 Hadoop® 集群上运行 mapreduce 之前,请确保集群和客户端计算机已正确配置。请咨询您的系统管理员,或参阅 配置 Hadoop 集群 (MATLAB Parallel Server)。
输出格式和顺序
在以二进制输出(默认设置)在 Hadoop 集群上运行 mapreduce 时,生成的 KeyValueDatastore 指向 Hadoop 序列文件,而不是其他环境中由 mapreduce 生成的二进制 MAT 文件。有关更多信息,请参阅 mapreduce 参考页面上的 'OutputType' 参量描述。
在 Hadoop 集群上运行 mapreduce 时,输出中的键值对的顺序与在其他环境中运行 mapreduce 不同。如果您的应用程序依赖于输出中的数据排列,则必须根据自己的要求对数据进行排序。
计算平均延迟
此示例显示如何修改用于计算平均航班延误的 MATLAB® 示例,以便在 Hadoop 集群上运行。
首先,您必须根据您的特定 Hadoop 配置设置适当的环境变量和集群属性。请联系您的系统管理员,了解这些属性以及向集群提交作业所需的其他属性的具体值,或参阅 配置 Hadoop 集群 (MATLAB Parallel Server)。
setenv('HADOOP_HOME','/path/to/hadoop/install') cluster = parallel.cluster.Hadoop;
创建一个 MapReducer 对象,以指定 mapreduce 必须使用您的 Hadoop 集群。
mr = mapreducer(cluster);
创建并预览数据存储。
ds = datastore('airlinesmall.csv','TreatAsMissing','NA',... 'SelectedVariableNames','ArrDelay','ReadSize',1000); preview(ds)
ArrDelay
________
8
8
21
13
4
59
3
11接下来,指定您的输出文件夹,输出 outds 并调用 mapreduce 在 mr 指定的 Hadoop 集群上执行。
outputFolder = 'hdfs:///home/myuser/out1'; outds = mapreduce(ds,@myMapperFcn,@myReducerFcn,... 'OutputFolder',outputFolder); meanDelay = mapreduce(ds,@meanArrivalDelayMapper,... @meanArrivalDelayReducer,mr,... 'OutputFolder',outputFolder)
Parallel mapreduce execution on the Hadoop cluster:
********************************
* MAPREDUCE PROGRESS *
********************************
Map 0% Reduce 0%
Map 66% Reduce 0%
Map 100% Reduce 66%
Map 100% Reduce 100%
meanDelay =
KeyValueDatastore with properties:
Files: {
' .../tmp/myuser/tpc00621b1_4eef_4abc_8078_646aa916e7d9/part0.seq'
}
ReadSize: 1 key-value pairs
FileType: 'seq'
读取结果。
readall(meanDelay)
Key Value
__________________ ________
'MeanArrivalDelay' [7.1201]
虽然出于演示目的,此示例使用了本地数据集,但在使用 Hadoop 时,您的数据集很可能存储在 HDFS™ 文件系统中。同样,您可能需要将 mapreduce 输出存储在 HDFS 中。有关在 MATLAB 中访问 HDFS 的详细信息,请参阅 处理远程数据。
支持函数
meanArrivalDelayMapper 映射函数用于计算每个数据模块中到达延迟的个数和总和。然后,映射器将这些值作为与键 "PartialCountSumDelay" 相关的中间值进行存储。
function meanArrivalDelayMapper(data,info,intermKVStore) % Data is an n-by-1 table of the ArrDelay. Remove missing values first: data(isnan(data.ArrDelay),:) = []; % Record the partial counts and sums and the reducer will accumulate them. partCountSum = [length(data.ArrDelay),sum(data.ArrDelay)]; add(intermKVStore,"PartialCountSumDelay",partCountSum); end
meanArrivalDelayReducer 还原函数接收映射器存储的每个模块的计数和总和。它对这些值进行求和,以得出总计数和总和。总体平均到达延迟是通过简单除法计算得出的。mapreduce 仅调用该归约器一次,因为映射器只添加了一个唯一的键。该还原器使用 add 将一个键值对添加到输出中。function meanArrivalDelayReducer(intermKey,intermValIter,outKVStore) count = 0; sum = 0; while hasnext(intermValIter) countSum = getnext(intermValIter); count = count + countSum(1); sum = sum + countSum(2); end meanDelay = sum/count; % The key-value pair added to outKVStore will become the output of mapreduce add(outKVStore,"MeanArrivalDelay",meanDelay); end