基于空行分割大文本文件的PySpark实现问询
PySpark 读取大文本文件并按空行分割数据块
文件情况与格式
我正在处理一个近2GB的大型文本文件,内容格式如下:
#*MOSFET table look-up models for circuit simulation #t1984 #cIntegration, the VLSI Journal #index1 #*The verification of the protection mechanisms of high-level language machines #@Virgil D. Gligor #t1984 #cInternational Journal of Parallel Programming #index2 #*Another view of functional and multivalued dependencies in the relational database model #@M. Gyssens, J. Paredaens #t1984 #cInternational Journal of Parallel Programming #index3 #*Entity-relationship diagrams which are in BCNF #@Sushil Jajodia, Peter A. Ng, Frederick N. Springsteel #t1984 #cInternational Journal of Parallel Programming #index4
需求
希望在PySpark中读取该文件,并按空行分割为对应的数据块,例如:#*Entity-relationship diagrams which are in BCNF #@Sushil Jajodia, Peter A. Ng, Frederick N. Springsteel #t1984 #cInternational Journal of Parallel Programming #index4
当前代码问题
当前编写的代码:
rdd = sc.textFile('acm.txt').flatMap( lambda x : x.split("\n\n") )
这个写法无法实现需求,因为textFile是按行读取文件,每个x代表单独的一行,单一行中不存在\n\n分隔符,所以split("\n\n")不会有任何分割效果,最终得到的还是原文件的每一行数据。
解决方案
方法1:使用wholeTextFiles读取整个文件
如果文件规模在单分区处理能力范围内,可以用wholeTextFiles读取整个文件内容,再按空行分割:
# 读取文件,获取文件名和内容 whole_rdd = sc.wholeTextFiles('acm.txt') # 提取文件内容,按空行分割成数据块并展平 block_rdd = whole_rdd.flatMap(lambda x: x[1].split('\n\n')) # 将每个数据块的换行替换为空格,匹配需求格式 formatted_rdd = block_rdd.map(lambda block: block.replace('\n', ' ').strip())
方法2:按行读取并分区内拼接数据块(适合超大文件)
通过自定义分区内的拼接逻辑,避免单分区加载整个文件的内存压力:
def merge_lines(iterator): current_block = [] for line in iterator: line = line.strip() if not line: # 遇到空行,输出当前拼接好的块 if current_block: yield ' '.join(current_block) current_block = [] else: current_block.append(line) # 处理分区内最后一个未完成的数据块 if current_block: yield ' '.join(current_block) # 按行读取文件,通过mapPartitions处理分区内的行拼接 rdd = sc.textFile('acm.txt').mapPartitions(merge_lines)
内容的提问来源于stack exchange,提问作者ShivanshVerma
相关产品推荐
相关产品推荐

