如何将Hadoop中带重复表头的多份CSV合并为本地单文件?
合并Hadoop目录下多份带重复表头的CSV到本地单个文件(仅保留一份表头)
针对Spark生成多CSV文件、coalesce(1)内存不足、hadoop getmerge重复表头的问题,这里提供两个高效的Python解决方案,均无需加载全量数据到内存:
方案一:结合Hadoop命令行工具处理(无需额外Python库)
通过调用hadoop fs命令获取文件列表,分批次下载并合并,仅保留第一个文件的表头:
import subprocess # 替换为你的实际路径 HADOOP_CSV_DIR = "/your/hadoop/csv/directory" LOCAL_OUTPUT = "/local/path/to/merged.csv" # 获取HDFS目录下所有有效CSV文件(排除_SUCCESS和目录) ls_cmd = f"hadoop fs -ls {HADOOP_CSV_DIR} | grep -v '_SUCCESS' | grep -v '^d' | awk '{{print $NF}}'" try: file_paths = subprocess.check_output(ls_cmd, shell=True).decode().strip().split('\n') file_paths = [path for path in file_paths if path.endswith('.csv')] except subprocess.CalledProcessError: print("Failed to list files in HDFS directory") exit() if not file_paths: print("No CSV files found in target HDFS directory") exit() # 写入第一个文件(包含完整表头) with open(LOCAL_OUTPUT, 'w', encoding='utf-8') as out_file: cat_cmd = f"hadoop fs -cat {file_paths[0]}" first_file_content = subprocess.check_output(cat_cmd, shell=True).decode('utf-8') out_file.write(first_file_content) # 追加剩余文件(跳过第一行表头) for path in file_paths[1:]: # 用tail -n +2跳过第一行 cat_cmd = f"hadoop fs -cat {path} | tail -n +2" try: content = subprocess.check_output(cat_cmd, shell=True).decode('utf-8') with open(LOCAL_OUTPUT, 'a', encoding='utf-8') as out_file: out_file.write(content) except subprocess.CalledProcessError: print(f"Warning: Failed to process file {path}, skipping") print(f"Merged CSV saved to {LOCAL_OUTPUT}")
优势
- 无需额外安装Python库,依赖Hadoop原生命令即可运行
- 逐文件流式处理,内存占用极低,适合超大文件场景
- 自动排除Spark生成的
_SUCCESS冗余文件
方案二:使用hdfs Python库(更优雅的API调用)
如果你的环境允许安装第三方库,可通过hdfs库直接操作HDFS,避免依赖shell命令:
首先安装依赖:
pip install hdfs
然后执行合并代码:
from hdfs import InsecureClient # 替换为你的NameNode地址和用户名 client = InsecureClient('http://your-namenode:50070', user='your-username') HADOOP_CSV_DIR = "/your/hadoop/csv/directory" LOCAL_OUTPUT = "/local/path/to/merged.csv" # 筛选有效CSV文件 file_paths = [] for status in client.list(HADOOP_CSV_DIR, status=True): if status['type'] != 'FILE': continue filename = status['pathSuffix'] if filename.endswith('.csv') and filename != '_SUCCESS': file_paths.append(f"{HADOOP_CSV_DIR}/{filename}") if not file_paths: print("No CSV files found") exit() # 写入第一个文件(含表头) with client.read(file_paths[0]) as hdfs_file, open(LOCAL_OUTPUT, 'w', encoding='utf-8') as local_file: local_file.write(hdfs_file.read().decode('utf-8')) # 追加剩余文件(跳过表头) for path in file_paths[1:]: with client.read(path) as hdfs_file, open(LOCAL_OUTPUT, 'a', encoding='utf-8') as local_file: # 读取并跳过第一行表头 hdfs_file.readline() # 写入剩余内容 local_file.write(hdfs_file.read().decode('utf-8')) print(f"Merged file created at {LOCAL_OUTPUT}")
优势
- 纯Python API操作,避免shell命令的兼容性问题
- 同样采用流式读取,内存效率高
- 更易于集成到现有Python数据处理流程中
注意事项
- 确保Hadoop命令行工具(方案一)或hdfs库(方案二)在运行环境中可用
- 如果文件编码不是UTF-8,需调整代码中的
encoding参数 - 大文件场景下,两种方案均不会一次性加载全量数据,内存压力小
- 运行前可检查本地输出文件是否存在,避免覆盖重要数据
内容的提问来源于stack exchange,提问作者Lazy_Nerd
相关产品推荐
相关产品推荐

