基于Hadoop的Python MapReduce矩阵乘法Reduce函数实现咨询
Hey there! Let's work through the Reduce function for your matrix multiplication using Hadoop MapReduce in Python. I’ll break this down step by step, including code examples and a test plan, since you already have your Map function ready.
First, let’s recap the core idea: matrix multiplication ( C = A \times B ) means each element ( C[i][j] = \sum_{k} A[i][k] \times B[k][j] ). Your Map function should already be emitting key-value pairs that group all ( A[i][k] ) and ( B[k][j] ) elements by the shared index ( k ). For example:
- For an input line
A,0,1,2.0(A's row 0, column 1, value 2.0), your Map should output1\tA,0,2.0 - For
B,1,3,5.0(B's row 1, column 3, value 5.0), your Map should output1\tB,3,5.0
The Reduce function’s job is to take all these grouped elements for a single ( k ), calculate the partial products ( A[i][k] \times B[k][j] ), and emit these partial results. We’ll then need a second Reduce step to sum all partial products for each ( (i,j) ) pair to get the final ( C[i][j] ).
This Reduce processes all elements for a single ( k ), stores A's row values and B's column values, then computes every possible product between them. Here’s the Python code:
import sys def main(): current_k = None a_elements = {} # Key: A's row index, Value: element value b_elements = {} # Key: B's column index, Value: element value for line in sys.stdin: line = line.strip() if not line: continue # Split into key (k) and value (type, index, value) key_part, value_part = line.split('\t', 1) k = key_part type_flag, idx, val = value_part.split(',', 2) val = float(val) idx = int(idx) # When we hit a new k, process the previous group first if current_k is not None and k != current_k: # Compute all partial products for the current k for i in a_elements: for j in b_elements: product = a_elements[i] * b_elements[j] print(f"C,{i},{j},{product}") # Reset for the new k group a_elements.clear() b_elements.clear() current_k = k else: current_k = k # Add the element to the correct dictionary if type_flag == 'A': a_elements[idx] = val elif type_flag == 'B': b_elements[idx] = val # Don't forget to process the last k group! if current_k is not None: for i in a_elements: for j in b_elements: product = a_elements[i] * b_elements[j] print(f"C,{i},{j},{product}") if __name__ == "__main__": main()
What this does:
- Tracks the current ( k ) we’re processing
- Stores all A and B elements for that ( k ) in dictionaries
- When moving to a new ( k ), it calculates every ( A[i][k] \times B[k][j] ) and emits a partial result in your desired format (
C,i,j,product)
The first Reduce outputs partial products for each ( k ). To get the final sum for each ( C[i][j] ), we need a quick second MapReduce job:
Mapper 2 (Simple Key Extraction)
This just rekeys the partial results by the ( (i,j) ) pair:
import sys for line in sys.stdin: line = line.strip() if not line: continue # Split the line into components _, i, j, val = line.split(',', 3) # Emit (i,j) as the key, and the partial product as the value print(f"{i},{j}\t{val}")
Reducer 2 (Sum Partial Products)
This sums all partial products for each ( (i,j) ) pair to get the final ( C[i][j] ):
import sys def main(): current_pair = None total = 0.0 for line in sys.stdin: line = line.strip() if not line: continue # Split into key (i,j) and value (partial product) key_part, val = line.split('\t', 1) i, j = key_part.split(',', 1) val = float(val) current_key = (i, j) # When we hit a new (i,j) pair, output the total for the previous one if current_pair is not None and current_key != current_pair: print(f"C,{current_pair[0]},{current_pair[1]},{total}") total = 0.0 current_pair = current_key else: current_pair = current_key total += val # Process the last (i,j) pair if current_pair is not None: print(f"C,{current_pair[0]},{current_pair[1]},{total}") if __name__ == "__main__": main()
Let’s use a simple 2x2 matrix pair to verify:
Input Data (A + B):
A,0,0,1.0 A,0,1,2.0 A,1,0,3.0 A,1,1,4.0 B,0,0,5.0 B,0,1,6.0 B,1,0,7.0 B,1,1,8.0
Expected Final Output:
C,0,0,19.0 C,0,1,22.0 C,1,0,43.0 C,1,1,50.0
How to Run with Hadoop Streaming:
- First job (partial products):
hadoop jar $HADOOP_HOME/share/hadoop/tools/lib/hadoop-streaming-*.jar \ -files mapper.py,reducer1.py \ -mapper mapper.py \ -reducer reducer1.py \ -input /path/to/input/matrix_data \ -output /tmp/matrix_intermediate
- Second job (sum partial products):
hadoop jar $HADOOP_HOME/share/hadoop/tools/lib/hadoop-streaming-*.jar \ -files mapper2.py,reducer2.py \ -mapper mapper2.py \ -reducer reducer2.py \ -input /tmp/matrix_intermediate \ -output /path/to/final/matrix_result
- Ensure your original Map function uses tabs to separate keys and values (Hadoop Streaming’s default delimiter)
- Double-check that your Map emits the correct format:
k\tA,i,valork\tB,j,val - If you know the dimensions of your matrices upfront, you can optimize, but this approach works for arbitrary sizes
内容的提问来源于stack exchange,提问作者HHKSHD_HH

