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

使用Python实现运行时可选两个MapReduce作业的方法咨询

实现方案

你可以通过新增一个统一的Python启动入口脚本实现需求,不需要修改已有的mapper、reducer代码,侵入性极低。

1. 提前梳理作业调用逻辑

先把两个作业的运行命令整理好,两种常见的运行场景对应命令如下:

  • Hadoop集群运行(使用Hadoop Streaming):
    电影平均评分作业命令示例:
    hadoop jar ${HADOOP_STREAMING_JAR_PATH} \
    -files mapper_movie.py,reducer_movie.py \
    -mapper mapper_movie.py \
    -reducer reducer_movie.py \
    -input ${INPUT_RATINGS_PATH} \
    -output ${MOVIE_AVG_OUTPUT_PATH}
    
    用户平均评分作业仅需要替换对应的mapper、reducer文件名和输出路径即可。
  • 本地调试运行(用管道模拟MapReduce流程):
    电影平均评分作业命令示例:
    cat ratings.txt | python mapper_movie.py | sort | python reducer_movie.py > movie_avg_result.txt
    

2. 编写入口脚本

新建run_jobs.py作为统一启动入口,代码示例如下:

import os
import sys

# 公共配置可统一在这里定义,避免重复修改
HADOOP_STREAMING_JAR = "/opt/hadoop/share/hadoop/tools/lib/hadoop-streaming-3.3.4.jar"
INPUT_PATH = "/data/ratings"
LOCAL_INPUT_FILE = "ratings.txt"

def run_movie_avg_job():
    run_mode = input("请选择运行模式:1. Hadoop集群 2. 本地调试 >> ")
    if run_mode == "1":
        output_path = input("请输入HDFS输出路径 >> ")
        # 提前校验输出路径是否存在,避免报错
        if os.popen(f"hdfs dfs -test -e {output_path}; echo $?").read().strip() == "0":
            del_choice = input(f"输出路径{output_path}已存在,是否删除后继续?y/n >> ")
            if del_choice.lower() == "y":
                os.system(f"hdfs dfs -rm -r {output_path}")
            else:
                print("请更换输出路径后重试")
                return
        cmd = f'''
        hadoop jar {HADOOP_STREAMING_JAR} \
        -files mapper_movie.py,reducer_movie.py \
        -mapper mapper_movie.py \
        -reducer reducer_movie.py \
        -input {INPUT_PATH} \
        -output {output_path}
        '''
        os.system(cmd)
        print(f"作业执行完成,结果存储在HDFS路径:{output_path}")
    elif run_mode == "2":
        output_file = input("请输入本地结果存储文件名 >> ")
        cmd = f"cat {LOCAL_INPUT_FILE} | python mapper_movie.py | sort | python reducer_movie.py > {output_file}"
        os.system(cmd)
        print(f"作业执行完成,结果存储在本地文件:{output_file}")
    else:
        print("无效的模式选择,退出程序")
        sys.exit(1)

def run_user_avg_job():
    run_mode = input("请选择运行模式:1. Hadoop集群 2. 本地调试 >> ")
    if run_mode == "1":
        output_path = input("请输入HDFS输出路径 >> ")
        if os.popen(f"hdfs dfs -test -e {output_path}; echo $?").read().strip() == "0":
            del_choice = input(f"输出路径{output_path}已存在,是否删除后继续?y/n >> ")
            if del_choice.lower() == "y":
                os.system(f"hdfs dfs -rm -r {output_path}")
            else:
                print("请更换输出路径后重试")
                return
        cmd = f'''
        hadoop jar {HADOOP_STREAMING_JAR} \
        -files mapper_user.py,reducer_user.py \
        -mapper mapper_user.py \
        -reducer reducer_user.py \
        -input {INPUT_PATH} \
        -output {output_path}
        '''
        os.system(cmd)
        print(f"作业执行完成,结果存储在HDFS路径:{output_path}")
    elif run_mode == "2":
        output_file = input("请输入本地结果存储文件名 >> ")
        cmd = f"cat {LOCAL_INPUT_FILE} | python mapper_user.py | sort | python reducer_user.py > {output_file}"
        os.system(cmd)
        print(f"作业执行完成,结果存储在本地文件:{output_file}")
    else:
        print("无效的模式选择,退出程序")
        sys.exit(1)

if __name__ == "__main__":
    print("=== MapReduce作业选择入口 ===")
    print("1. 计算每部电影的平均评分")
    print("2. 计算每个用户的平均评分")
    choice = input("请输入要运行的作业编号 >> ")
    if choice == "1":
        run_movie_avg_job()
    elif choice == "2":
        run_user_avg_job()
    else:
        print("无效的作业编号,退出程序")
        sys.exit(1)

3. 运行方式

直接在终端执行命令即可:

python run_jobs.py

按照终端提示选择对应作业、运行模式和输出路径即可启动对应的MapReduce作业。

可选优化点

  • 可以把作业的mapper、reducer文件名也抽到公共配置中,进一步减少重复代码
  • 新增作业执行状态校验,执行完成后自动判断作业是否成功运行
  • 支持批量作业执行,可同时选择多个作业依次运行

内容的提问来源于stack exchange,提问作者Adam

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 16:15:02