You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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的数据分配逻辑:

  1. 当前Mapper输出的Key是(matrix_num, i, j),MapReduce会根据这个Key的哈希值分配到不同Reducer;
  2. 当Reducer数量>1时,每个Reducer只会拿到部分矩阵的片段数据(比如某个Reducer只拿到M的部分元素、N的部分元素);
  3. 原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()

修改逻辑说明

  1. Key设计优化:
    • 矩阵乘法中,将计算MN[i][j]所需的所有M[i][k]和N[k][j]都标记为同一个Key(i,j);
    • X矩阵的元素X[i][j]也标记为Key(i,j),确保和对应的MN结果进入同一个Reducer;
  2. Reducer逻辑调整:
    • 同一Key(i,j)的所有数据会被分配到同一个Reducer;
    • 收集该Key对应的所有M、N元素,计算MN[i][j],再结合X元素计算X[i][j]-MN[i][j];
  3. 多Reducer支持:每个Reducer只负责计算若干个(i,j)对应的结果,不需要完整的矩阵数据,支持任意数量的Reducer。

内容的提问来源于stack exchange,提问作者Jaimee-lee Lincoln

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.17 22:25:25