如何使用Apache Spark处理纯文本账单文件并转换为逗号分隔格式文件
Spark账单文本格式化实现方案
1. 实现逻辑
- 读取原始账单文本文件,每一行作为单独的字符串字段加载
- 过滤所有匹配
===Start Page: X===、===End Page: X===格式的分页标记行 - 从过滤后的有效行中取前4行作为表头,剩余行作为交易数据
- 将交易数据行的空格分隔符替换为逗号,生成符合要求的CSV格式内容
- 合并表头和处理后的交易数据,输出为最终文件
2. PySpark实现代码
from pyspark.sql import SparkSession from pyspark.sql.functions import col, regexp_replace if __name__ == "__main__": # 初始化Spark会话 spark = SparkSession.builder.appName("BillProcess").getOrCreate() # 读取原始文本文件 raw_df = spark.read.text("替换为你的原始账单文件路径") # 过滤分页标记行 filtered_df = raw_df.filter( ~col("value").rlike(r"===Start Page: \d+===") & ~col("value").rlike(r"===End Page: \d+===") ) # 提取前4行作为表头 header_list = [row["value"] for row in filtered_df.limit(4).collect()] header_rdd = spark.sparkContext.parallelize(header_list) # 处理交易数据:空格替换为逗号 data_df = filtered_df.exceptAll(filtered_df.limit(4)) processed_data_rdd = data_df.select( regexp_replace(col("value"), " ", ",").alias("value") ).rdd.map(lambda row: row["value"]) # 合并表头和数据后输出 header_rdd.union(processed_data_rdd).coalesce(1).saveAsTextFile("替换为输出文件路径") spark.stop()
3. Scala Spark实现代码
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions.{col, regexp_replace} object BillProcess { def main(args: Array[String]): Unit = { val spark = SparkSession.builder() .appName("BillProcess") .getOrCreate() // 读取原始文本 val rawDF = spark.read.text("替换为你的原始账单文件路径") // 过滤分页行 val filteredDF = rawDF.filter( !col("value").rlike("===Start Page: \\d+===") && !col("value").rlike("===End Page: \\d+===") ) // 提取表头 val headerList = filteredDF.limit(4).collect().map(_.getAs[String]("value")) val headerRDD = spark.sparkContext.parallelize(headerList) // 处理交易数据 val dataDF = filteredDF.except(filteredDF.limit(4)) val processedDataRDD = dataDF.select( regexp_replace(col("value"), " ", ",").as("value") ).as[String].rdd // 合并输出 headerRDD.union(processedDataRDD) .coalesce(1) .saveAsTextFile("替换为输出文件路径") spark.stop() } }
4. 注意事项
- 如果原始账单字段之间存在多个连续空格,可将替换逻辑改为先合并连续空格再替换:
regexp_replace(regexp_replace(col("value"), "\\s+", " "), " ", ","),避免生成多余空字段 coalesce(1)用于将输出合并为单个文件,仅适合小数据量场景,大数据量场景删除该配置即可输出为多个分片文件- 如果表头结构固定,可直接手动构造表头RDD,无需动态读取文件前4行,执行效率更高
内容的提问来源于stack exchange,提问作者neel
相关产品推荐
相关产品推荐

