如何使用PySpark从HDFS读取DOCX/PDF文件?
用PySpark从HDFS读取DOCX/PDF文件的解决方案
确实,pandas的格式支持范围有限,要在PySpark里处理HDFS上的DOCX和PDF文件,得结合Spark的分布式能力和Python的文档解析库来实现,我给你梳理两种实用的方案:
方案一:用Spark原生binaryFiles API读取(推荐)
这个方法不需要额外的HDFS客户端库,直接用Spark的API读取二进制文件,更适配分布式场景:
步骤1:给所有Spark节点安装依赖
首先确保每个Executor节点都装了解析库:
pip install python-docx PyPDF2
步骤2:编写PySpark代码
from pyspark import SparkContext, SparkConf import docx import PyPDF2 import io # 初始化Spark上下文 conf = SparkConf().setAppName("ReadDocxPdfFromHDFS") sc = SparkContext(conf=conf) # 用binaryFiles读取HDFS上的目标文件,返回(file_path, 二进制内容)的RDD # 这里替换成你的HDFS路径,注意用hdfs://协议(端口通常是9000) binary_rdd = sc.binaryFiles('hdfs://192.00.00.30:9000/user/user/*') def parse_file(file_tuple): file_path, binary_data = file_tuple # 将二进制内容转为字节流,方便解析库处理 file_stream = io.BytesIO(binary_data) # 根据后缀判断文件类型并解析内容 if file_path.endswith('.docx'): doc = docx.Document(file_stream) content = '\n'.join([para.text for para in doc.paragraphs]) elif file_path.endswith('.pdf'): pdf_reader = PyPDF2.PdfReader(file_stream) content = '\n'.join([page.extract_text() for page in pdf_reader.pages]) else: content = "不支持的文件格式" return (file_path, content) # 解析所有文件并转为DataFrame result_rdd = binary_rdd.map(parse_file) result_df = result_rdd.toDF(["文件路径", "文件内容"]) # 查看结果(truncate=False显示完整内容) result_df.show(truncate=False)
方案二:结合HDFS客户端+RDD处理
如果习惯用hdfs库操作HDFS,也可以用这种方式,适合需要更精细控制HDFS操作的场景:
步骤1:安装依赖
除了解析库,还要装HDFS客户端:
pip install python-docx PyPDF2 hdfs
步骤2:编写代码
from pyspark import SparkContext, SparkConf from hdfs import InsecureClient import docx import PyPDF2 import io conf = SparkConf().setAppName("ReadDocxPdfFromHDFS") sc = SparkContext(conf=conf) # 初始化HDFS客户端 client_hdfs = InsecureClient('http://192.00.00.30:50070') # 获取目标路径下的所有DOCX/PDF文件路径 def get_target_files(hdfs_path): # 遍历路径下的文件,过滤掉文件夹 all_files = [f.path for f in client_hdfs.list(hdfs_path, status=True) if f['type'] != 'DIRECTORY'] # 筛选出DOCX和PDF文件 return [path for path in all_files if path.endswith(('.docx', '.pdf'))] target_files = get_target_files('/user/user') # 转为RDD并处理,用mapPartitions减少客户端初始化次数(优化性能) def process_partition(file_paths): # 每个分区初始化一次客户端,避免重复创建连接 client = InsecureClient('http://192.00.00.30:50070') for path in file_paths: with client.read(path) as f: file_stream = io.BytesIO(f.read()) # 解析逻辑和方案一一致 if path.endswith('.docx'): doc = docx.Document(file_stream) content = '\n'.join([para.text for para in doc.paragraphs]) elif path.endswith('.pdf'): pdf_reader = PyPDF2.PdfReader(file_stream) content = '\n'.join([page.extract_text() for page in pdf_reader.pages]) else: content = "不支持的文件格式" yield (path, content) file_rdd = sc.parallelize(target_files) result_rdd = file_rdd.mapPartitions(process_partition) result_df = result_rdd.toDF(["文件路径", "文件内容"]) result_df.show(truncate=False)
注意事项
- 依赖安装:一定要确保所有Spark节点(包括Driver和Executor)都安装了所需的Python库,不然运行时会报
ModuleNotFoundError。 - 大文件优化:如果处理大量大文件,优先用
mapPartitions代替map,减少HDFS客户端或解析库的初始化次数,提升性能。 - PDF解析精度:PyPDF2对某些复杂PDF的解析可能不够完美,如果需要更高精度,可以换成
pdfplumber库(同样需要在所有节点安装)。
内容的提问来源于stack exchange,提问作者Bhagesh Arora
相关产品推荐
相关产品推荐

