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

基于Python的Hadoop Streaming(2.9.0) MapReduce函数设计咨询

Solution for Hadoop Streaming MapReduce with Python (2.9.0)

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.combinations ensures we only generate unique pairs (avoids duplicates like v11 union v12 and v12 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 v11 instead 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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 09:59:14