如何在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
相关产品推荐
相关产品推荐

