如何通过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()
现有方案的核心问题
- 全量数据拉取到Driver:
collect()将所有数据拉取到Driver端,完全浪费Spark分布式计算能力,所有写入操作串行执行。 - 逐行append效率极低:每次
append()都会触发磁盘IO,无批量优化,IO开销巨大。 - 多线程无效:
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
相关产品推荐
相关产品推荐

