MapReduce单Reducer正常,多Reducer失效问题排查与修复求助
问题描述
在MapReduce框架中实现**矩阵乘法(M*N)后执行减法(X-MN)**的任务,遇到以下问题:
- 当设置
-D mapred.reduce.tasks=1(仅1个Reducer)时,代码运行正常; - 增加Reducer数量后任务直接失败。
矩阵说明
- 均为同尺寸方阵(示例为3*3);
- 三个输入矩阵通过每行首值区分:M首值为
1,N首值为2,X首值为3; - 示例M.txt内容:
1, 0, 1, 2, 3 1, 1, 4, 5, 6 1, 2, 7, 8, 9
现有代码及运行脚本
Mapper代码
m_matrix_num = 1 n_matrix_num = 2 x_matrix_num = 3 row = int(sys.argv[1]) col = int(sys.argv[2]) def fill_dictionaries_from_txt(txt_list, matrix_num): val_list = [] val_list.append(txt_list.split(",")) val = [int(elem) for elem in val_list[0][2:]] i = int(val_list[0][1]) for j in range(col): key_name = (matrix_num, i, j) print(key_name, val[j]) for line in sys.stdin: if line !="\n": if line[0]=="1": listM_dict = fill_dictionaries_from_txt(line, m_matrix_num) if line[0]=="2": listN_dict = fill_dictionaries_from_txt(line, n_matrix_num) if line[0]=="3": listX_dict = fill_dictionaries_from_txt(line, x_matrix_num)
Reducer代码
import sys m_matrix_num = 1 n_matrix_num = 2 x_matrix_num = 2 # 笔误:应为3 mn_matrix_num = 4 listM_dict = {} listN_dict = {} listX_dict = {} listMN_dict = {} for line in sys.stdin: key_str, val_str = line.split(":") key_str = [elem.strip() for elem in key_str if elem.isalnum()] key_str = [int(elem) for elem in key_str] key = tuple(key_str) val = val_str[3:-3] val = int(val) if key[0] == 1: listM_dict[key] = val if key[0] == 2: listN_dict[key] = val if key[0] == 3: listX_dict[key] = val row = max(listM_dict, key=listM_dict.get)[1] + 1 col = max(listM_dict, key=listM_dict.get)[2] + 1 for i in range(row): for j in range(col): sum = 0 for k in range(row): key_m = (m_matrix_num, i , k) key_n = (n_matrix_num, k, j) key_mn = (mn_matrix_num, i, j) if key_mn in listMN_dict: listMN_dict[key_mn] += listM_dict.get(key_m) * listN_dict.get(key_n) else: listMN_dict[key_mn] = listM_dict.get(key_m) * listN_dict.get(key_n) for i in range(row): for j in range(col): key_mn = (mn_matrix_num, i, j) key_x = (x_matrix_num, i, j) print('%s %s %s' % (i, j, listX_dict.get(key_x) - listMN_dict.get(key_mn)))
运行脚本
#!/bin/bash row=3 column=3 hadoop jar ./hadoop-streaming-3.1.4.jar \ -D mapred.reduce.tasks=3 \ -file ./mapper.py \ -mapper ./mapper.py $row $column \ -file ./reducer.py \ -reducer ./reducer.py \ -input /input \ -output /output
猜测原因:Mapper输出被分散到多个Reducer中,导致单个Reducer无法获取矩阵M、N、X的全部数据,无法完成完整的矩阵乘法和减法。
原因分析
你的猜测完全正确,核心问题出在Mapper输出的Key设计和Reducer的数据分配逻辑:
- 当前Mapper输出的Key是
(matrix_num, i, j),MapReduce会根据这个Key的哈希值分配到不同Reducer; - 当Reducer数量>1时,每个Reducer只会拿到部分矩阵的片段数据(比如某个Reducer只拿到M的部分元素、N的部分元素);
- 原Reducer逻辑依赖完整的M、N、X矩阵才能计算MN和X-MN,单个Reducer缺少数据就会导致计算失败(比如
listM_dict.get(key_m)返回None,乘法时抛出异常)。
另外,原代码还有两个明显bug:
- Reducer中
x_matrix_num = 2是笔误,应为3,否则无法正确读取X矩阵数据; - Mapper输出格式为
(1, 0, 0) 1,但Reducer用line.split(":")分割,会导致分割失败。
修改方案
要支持多Reducer,必须重新设计Key,让计算同一个结果元素所需的所有数据都分配到同一个Reducer。以下是调整后的代码:
修改后的Mapper代码
import sys def main(): row = int(sys.argv[1]) col = int(sys.argv[2]) for line in sys.stdin: line = line.strip() if not line: continue parts = [p.strip() for p in line.split(",")] matrix_flag = int(parts[0]) current_row = int(parts[1]) values = list(map(int, parts[2:])) if matrix_flag == 1: # M矩阵:元素M[current_row][k],需和所有N[k][j]相乘得到MN[current_row][j] for k in range(col): val = values[k] for j in range(col): print(f"{current_row}\t{j}\tM\t{k}\t{val}") elif matrix_flag == 2: # N矩阵:元素N[current_row][j],需和所有M[i][current_row]相乘得到MN[i][j] j_idx = 0 for val in values: for i in range(row): print(f"{i}\t{j_idx}\tN\t{current_row}\t{val}") j_idx += 1 elif matrix_flag == 3: # X矩阵:元素X[current_row][j],直接关联到结果位置(i,j) j_idx = 0 for val in values: print(f"{current_row}\t{j_idx}\tX\t{val}") j_idx += 1 if __name__ == "__main__": main()
修改后的Reducer代码
import sys from collections import defaultdict def main(): current_key = None x_val = 0 m_elements = defaultdict(int) n_elements = defaultdict(int) for line in sys.stdin: line = line.strip() if not line: continue parts = line.split("\t") i = int(parts[0]) j = int(parts[1]) flag = parts[2] key = (i, j) # Key切换时,计算上一个Key的结果 if current_key is not None and key != current_key: # 计算MN[i][j] mn_val = 0 for k in m_elements: mn_val += m_elements[k] * n_elements.get(k, 0) # 计算X-MN print(f"{current_key[0]}\t{current_key[1]}\t{x_val - mn_val}") # 重置变量 x_val = 0 m_elements.clear() n_elements.clear() current_key = key if flag == "M": k = int(parts[3]) val = int(parts[4]) m_elements[k] = val elif flag == "N": k = int(parts[3]) val = int(parts[4]) n_elements[k] = val elif flag == "X": x_val = int(parts[3]) # 处理最后一个Key的结果 if current_key is not None: mn_val = 0 for k in m_elements: mn_val += m_elements[k] * n_elements.get(k, 0) print(f"{current_key[0]}\t{current_key[1]}\t{x_val - mn_val}") if __name__ == "__main__": main()
修改逻辑说明
- Key设计优化:
- 矩阵乘法中,将计算
MN[i][j]所需的所有M[i][k]和N[k][j]都标记为同一个Key(i,j); - X矩阵的元素
X[i][j]也标记为Key(i,j),确保和对应的MN结果进入同一个Reducer;
- 矩阵乘法中,将计算
- Reducer逻辑调整:
- 同一Key
(i,j)的所有数据会被分配到同一个Reducer; - 收集该Key对应的所有M、N元素,计算
MN[i][j],再结合X元素计算X[i][j]-MN[i][j];
- 同一Key
- 多Reducer支持:每个Reducer只负责计算若干个
(i,j)对应的结果,不需要完整的矩阵数据,支持任意数量的Reducer。
内容的提问来源于stack exchange,提问作者Jaimee-lee Lincoln
相关产品推荐
相关产品推荐

