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

如何使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.23 17:24:00