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

如何在Spark中读取无分隔符的固定字节长度文件?

Spark读取固定字节长度分隔的无分隔符文件解决方案

你需要读取按固定72字节分隔的无分隔符文件,但直接使用spark.read()会将所有内容合并为一行,以下是两种可行的解决方案,能生成对应3条独立记录的Dataset<Row>:

方法一:读取整文件后拆分固定长度记录

先读取整个文件为单条数据,再按指定字节长度拆分每条记录,适合中小文件场景:

import org.apache.spark.sql.SparkSession

val spark = SparkSession.builder().appName("FixedLengthReader").getOrCreate()

// 读取文件为单一行的Dataset
val rawFile = spark.read.text("path/to/your/file.txt").as[String]

// 定义每条记录的字节长度
val recordByteLength = 72

// 按字节拆分记录,注意匹配文件实际编码(示例用UTF-8)
val splitRecords = rawFile.flatMap { content =>
  val contentBytes = content.getBytes("UTF-8")
  contentBytes.grouped(recordByteLength).map(new String(_, "UTF-8"))
}.filter(_.trim.nonEmpty) // 过滤可能存在的空记录

// 转换为Dataset<Row>,如需解析字段可继续处理
splitRecords.toDF("raw_record").show(false)

方法二:使用Hadoop FixedLengthInputFormat(推荐大文件)

直接在输入阶段按固定字节拆分记录,无需先加载整个文件,适合超大文件场景:

import org.apache.spark.sql.SparkSession
import org.apache.hadoop.mapreduce.lib.input.FixedLengthInputFormat
import org.apache.hadoop.io.{LongWritable, Text}

val spark = SparkSession.builder().appName("HadoopFixedLengthReader").getOrCreate()
val sc = spark.sparkContext

// 配置每条记录的字节长度
sc.hadoopConfiguration.set(FixedLengthInputFormat.FIXED_RECORD_LENGTH, "72")

// 通过Hadoop输入格式读取文件
val hadoopRDD = sc.newAPIHadoopFile(
  "path/to/your/file.txt",
  classOf[FixedLengthInputFormat],
  classOf[LongWritable],
  classOf[Text]
)

// 提取记录内容并转换为Dataset
val fixedLengthDS = spark.createDataset(hadoopRDD.map(_._2.toString))
fixedLengthDS.toDF("raw_record").show(false)

可选:解析记录字段

如果需要将每条72字节的记录拆分为具体字段(比如序列号、姓名等),可以基于上述结果继续处理:

import org.apache.spark.sql.functions._

val parsedDF = fixedLengthDS.toDF("raw_record")
  .select(
    substring(col("raw_record"), 1, 1).alias("serial_number"),
    substring(col("raw_record"), 2, 10).trim.alias("name"),
    substring(col("raw_record"), 12, 15).trim.alias("position"),
    substring(col("raw_record"), 27, 20).trim.alias("address"),
    substring(col("raw_record"), 47, 25).trim.alias("category")
  )

parsedDF.show(false)

内容的提问来源于stack exchange,提问作者Jay

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 15:01:15