使用Python实现运行时可选两个MapReduce作业的方法咨询
实现方案
你可以通过新增一个统一的Python启动入口脚本实现需求,不需要修改已有的mapper、reducer代码,侵入性极低。
1. 提前梳理作业调用逻辑
先把两个作业的运行命令整理好,两种常见的运行场景对应命令如下:
- Hadoop集群运行(使用Hadoop Streaming):
电影平均评分作业命令示例:
用户平均评分作业仅需要替换对应的mapper、reducer文件名和输出路径即可。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} - 本地调试运行(用管道模拟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
相关产品推荐
相关产品推荐

