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

如何通过PySpark高效将Hadoop(Hive)数据写入DBF文件?

高效将Spark数据写入DBF文件的优化方案

问题背景

我们需要从Hadoop(Hive)查询数据并保存为DBF文件,采用PySpark(Python 3.4)作为处理引擎,借助dbf包实现DBF写入。测试发现该过程耗时极长(可达20分钟),远慢于CSV、ORC等格式。现有两种实现方式:

  • 基础单线程写法(耗时约20分钟)
  • 多线程写法(效果相近)

已将Spark Driver内存设为7GB,但写入DBF时仅占用约1GB内存,寻求高效写入的调优项或替代方案。

现有基础实现代码

import dbf
from datetime import datetime

collections = spark.sql("SELECT JENISKEGIA, JUMLAHUM_A, ... , URUTAN, WEIGHT FROM silastik.sakernas_2022_8").collect()

filename2="/home/sak202208_"+str(datetime.now())+"_tes.dbf"

header2 = "JENISKEGIA N(8,0); JUMLAHUM_A N(8,0); ... , URUTAN N(7,0); WEIGHT N(8,0)"

new_table2 = dbf.Table(filename2, header2)

new_table2.open(dbf.READ_WRITE)

for row in collections:
    new_table2.append(row)

new_table2.close()

现有多线程实现代码

import dbf
from datetime import datetime
import concurrent.futures
import os

collections = spark.sql("SELECT JENISKEGIA, JUMLAHUM_A, ... , URUTAN, WEIGHT FROM silastik.sakernas_2022_8").collect()

filename2="/home/sak202208_"+str(datetime.now())+"_tes.dbf"

header2 = "JENISKEGIA N(8,0); JUMLAHUM_A N(8,0); ... , URUTAN N(7,0); WEIGHT N(8,0)"

new_table2 = dbf.Table(filename2, header2)

new_table2.open(dbf.READ_WRITE)

def append_row(table, record):
    table.append(record)

with concurrent.futures.ThreadPoolExecutor(max_workers=min(32, (os.cpu_count() or 1) + 4)) as executor:
    for row in collections:
        executor.submit(append_row(new_table2, row))

new_table2.close()

现有方案的核心问题

  1. 全量数据拉取到Driver:collect()将所有数据拉取到Driver端,完全浪费Spark分布式计算能力,所有写入操作串行执行。
  2. 逐行append效率极低:每次append()都会触发磁盘IO,无批量优化,IO开销巨大。
  3. 多线程无效:dbf.Table并非线程安全,多线程写入会引发锁竞争或数据损坏,实际无法并行执行。

优化方案

一、优化dbf包使用方式

1. 批量写入替代逐行append

使用append_many()一次性写入多条记录,大幅减少磁盘IO次数:

import dbf
from datetime import datetime

collections = spark.sql("SELECT JENISKEGIA, JUMLAHUM_A, ... , URUTAN, WEIGHT FROM silastik.sakernas_2022_8").collect()

filename2 = f"/home/sak202208_{datetime.now()}_tes.dbf"
header2 = "JENISKEGIA N(8,0); JUMLAHUM_A N(8,0); ... , URUTAN N(7,0); WEIGHT N(8,0)"

new_table2 = dbf.Table(filename2, header2)
new_table2.open(dbf.READ_WRITE)

# 批量写入所有记录
new_table2.append_many(collections)

new_table2.close()

2. 分批次拉取数据(避免Driver内存溢出)

若数据量过大,分批次拉取并批量写入,降低Driver内存压力:

import dbf
from datetime import datetime

query = spark.sql("SELECT JENISKEGIA, JUMLAHUM_A, ... , URUTAN, WEIGHT FROM silastik.sakernas_2022_8")
batch_size = 10000  # 根据Driver内存调整批次大小

filename2 = f"/home/sak202208_{datetime.now()}_tes.dbf"
header2 = "JENISKEGIA N(8,0); JUMLAHUM_A N(8,0); ... , URUTAN N(7,0); WEIGHT N(8,0)"

new_table2 = dbf.Table(filename2, header2)
new_table2.open(dbf.READ_WRITE)

# 按分区分批次处理
for batch in query.rdd.mapPartitions(lambda iter: [list(iter)]).collect():
    new_table2.append_many(batch)

new_table2.close()

二、利用Spark分布式写入(推荐)

放弃Driver端集中写入,改为Executor端分布式生成DBF文件,再合并为单文件:

1. 分布式生成分区DBF文件

每个Executor处理一个数据分区,写入本地临时DBF文件:

from datetime import datetime
import dbf
import os

def write_dbf_partition(iterator):
    # 获取分区ID生成唯一临时文件名
    partition_id = os.environ.get("SPARK_PARTITION_ID", "0")
    filename = f"/tmp/sak202208_part_{partition_id}_{datetime.now()}.dbf"
    header = "JENISKEGIA N(8,0); JUMLAHUM_A N(8,0); ... , URUTAN N(7,0); WEIGHT N(8,0)"
    
    table = dbf.Table(filename, header)
    table.open(dbf.READ_WRITE)
    table.append_many(list(iterator))
    table.close()
    return [filename]

# 执行分布式写入
result_files = spark.sql("SELECT JENISKEGIA, JUMLAHUM_A, ... , URUTAN, WEIGHT FROM silastik.sakernas_2022_8")\
    .rdd.mapPartitions(write_dbf_partition)\
    .collect()

2. 合并分区DBF文件

将多个临时DBF文件合并为最终单文件:

def merge_dbf_files(output_file, input_files):
    header = "JENISKEGIA N(8,0); JUMLAHUM_A N(8,0); ... , URUTAN N(7,0); WEIGHT N(8,0)"
    main_table = dbf.Table(output_file, header)
    main_table.open(dbf.READ_WRITE)
    
    for file in input_files:
        temp_table = dbf.Table(file)
        temp_table.open(dbf.READ_ONLY)
        main_table.append_many(temp_table)
        temp_table.close()
        os.remove(file)
    
    main_table.close()

# 生成最终文件
final_file = f"/home/sak202208_{datetime.now()}_final.dbf"
merge_dbf_files(final_file, result_files)

替代方案

1. 先写CSV再转DBF

利用Spark快速分布式写入CSV,再批量转DBF:

# 1. Spark写入CSV
csv_path = "/tmp/sak202208_temp"
spark.sql("SELECT JENISKEGIA, JUMLAHUM_A, ... , URUTAN, WEIGHT FROM silastik.sakernas_2022_8")\
    .write.mode("overwrite").option("header", "false").csv(csv_path)

# 2. CSV转DBF
import dbf
from datetime import datetime
import csv
import glob

filename2 = f"/home/sak202208_{datetime.now()}_tes.dbf"
header2 = "JENISKEGIA N(8,0); JUMLAHUM_A N(8,0); ... , URUTAN N(7,0); WEIGHT N(8,0)"

new_table2 = dbf.Table(filename2, header2)
new_table2.open(dbf.READ_WRITE)

# 读取所有CSV分区文件
for csv_file in glob.glob(f"{csv_path}/part-*"):
    with open(csv_file, "r") as f:
        reader = csv.reader(f)
        records = [tuple(row) for row in reader]
        new_table2.append_many(records)

new_table2.close()

2. 使用支持分布式DBF写入的第三方库

若条件允许,可选用支持Spark分布式写入DBF的开源/商业扩展库,简化实现逻辑。


内容的提问来源于stack exchange,提问作者m hanif f

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 21:37:03