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

如何将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数据处理流程中

注意事项

  1. 确保Hadoop命令行工具(方案一)或hdfs库(方案二)在运行环境中可用
  2. 如果文件编码不是UTF-8,需调整代码中的encoding参数
  3. 大文件场景下,两种方案均不会一次性加载全量数据,内存压力小
  4. 运行前可检查本地输出文件是否存在,避免覆盖重要数据

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.12 22:10:29