基于Python的Hadoop Streaming(2.9.0) MapReduce函数设计咨询
Hey there! Let's walk through building this MapReduce job to generate pairwise unions of value sets for each key. Here's a step-by-step solution tailored to your needs.
1. Mapper Script (map.py)
The mapper's core job is to parse input lines, extract key-value pairs, and emit them in a format Hadoop can use to shuffle all values for the same key together.
Case 1: Single key-value pair per line
If your input lines look like k1,v11 (one pair per line), use this code:
#!/usr/bin/env python3 import sys def main(): for line in sys.stdin: line = line.strip() if not line: continue # Split key and value set (handle commas in values with split(',', 1)) key, value_set = line.split(',', 1) # Emit key + tab-separated value set (Hadoop's default delimiter) print(f"{key}\t{value_set}") if __name__ == "__main__": main()
Case 2: Multiple key-value pairs per line
If your input lines are like k1,v11 k2,v21 k1,v12 (space-separated pairs), modify the mapper to split each line into individual pairs first:
#!/usr/bin/env python3 import sys def main(): for line in sys.stdin: line = line.strip() if not line: continue # Split line into individual key-value pairs pairs = line.split() for pair in pairs: key, value_set = pair.split(',', 1) print(f"{key}\t{value_set}") if __name__ == "__main__": main()
2. Reducer Script (reduce.py)
The reducer receives all value sets for a single key, generates every unique pair of sets, computes their union, and formats the output as you specified.
#!/usr/bin/env python3 import sys from itertools import combinations def main(): current_key = None value_sets = [] for line in sys.stdin: line = line.strip() if not line: continue key, value_str = line.split('\t', 1) # Convert value string to a Python set (adjust delimiter if your values use spaces) value_set = set(value_str.split(',')) if current_key == key: value_sets.append(value_set) else: # Process the previous key if it exists if current_key is not None: process_key(current_key, value_sets) # Reset for new key current_key = key value_sets = [value_set] # Process the last key after loop ends if current_key is not None: process_key(current_key, value_sets) def process_key(key, value_sets): # Skip if there are fewer than 2 value sets (no pairs to form) if len(value_sets) < 2: return # Generate all unique pairwise combinations union_results = [] for vs1, vs2 in combinations(value_sets, 2): # Compute union, sort for consistent output, then format union = ','.join(sorted(vs1.union(vs2))) # Format as "v11 union v12 = [union content]" (adjust to match your preferred style) union_results.append(f"{','.join(sorted(vs1))} union {','.join(sorted(vs2))} = {union}") # Output key followed by comma-separated union results print(f"{key}, {', '.join(union_results)}") if __name__ == "__main__": main()
Key Reducer Notes:
itertools.combinationsensures we only generate unique pairs (avoids duplicates likev11 union v12andv12 union v11).- Sorting union results makes output predictable, since Hadoop doesn't guarantee the order of values received by the reducer.
- If you need to keep original value labels (like
v11instead of the actual set content), modify the mapper to emit both the label and set (e.g.,print(f"{key}\t{v_label}\t{v_set}")) and update the reducer to track labels alongside sets.
3. Run the Job
Step 1: Make scripts executable
chmod +x map.py reduce.py
Step 2: Test locally (without Hadoop)
Mimic Hadoop's shuffle/sort stage with the sort command:
cat input.txt | ./map.py | sort | ./reduce.py
Step 3: Run on Hadoop
Use the Hadoop Streaming jar (adjust paths to match your installation):
hadoop jar $HADOOP_HOME/share/hadoop/tools/lib/hadoop-streaming-2.9.0.jar \ -files map.py,reduce.py \ -mapper ./map.py \ -reducer ./reduce.py \ -input /hdfs/path/to/input \ -output /hdfs/path/to/output
Example Output
For input lines:
k1,1,2 k2,5,6 k1,3,4 k2,7,8 k1,2,5
The output will look like:
k1, 1,2 union 3,4 = 1,2,3,4, 1,2 union 2,5 = 1,2,5, 3,4 union 2,5 = 2,3,4,5 k2, 5,6 union 7,8 = 5,6,7,8
内容的提问来源于stack exchange,提问作者cdt

