如何使用.NET for Apache Spark加载包含多区块的固定位置格式文件
实现思路与方案
首先明确:完全可以通过UDF配合窗口函数/分区遍历实现该需求,核心逻辑是基于文件行顺序关联最近的01头信息,具体实现思路如下:
第一步:保留原始行顺序并标记行类型
读取原始文本文件时,先为每一行生成全局唯一的递增行号,确保文件的原始顺序不被分布式处理打乱,同时拆分每行前2位作为行类型标识:
- 01:发明者头信息行
- 02:子头行(直接过滤即可,无业务价值)
- 03:发明条目行
第二步:解析01头信息(可通过UDF实现)
编写UDF专门解析固定宽度的01行,提取出发明者姓名、邮箱、电话三个字段,非01行的这三个字段统一设为null。
固定宽度解析逻辑示例(按你提供的样例对齐):
// UDF入参是原始行字符串,返回值是三元组(姓名,邮箱,电话) Func<Column, Column> parse01Row = Udf<string, (string, string, string)>((row) => { if (row.StartsWith("01")) { string name = row.Substring(2, 30).Trim(); string email = row.Substring(32, 70).Trim(); string phone = row.Substring(102).Trim(); return (name, email, phone); } return (null, null, null); });
第三步:关联最近的01头信息
有两种常用实现方式,按需选择即可:
方式1:窗口函数实现(DataFrame API)
定义按行号升序排序的全局窗口,使用last函数配合ignorenulls=true参数,自动将当前行之前最近的非空01头信息填充到所有行:
var window = Window.OrderBy("line_id"); var filledDf = rawDf .WithColumn("owner_info", parse01Row(Col("raw_line"))) .WithColumn("name", Last(Col("owner_info.Item1"), true).Over(window)) .WithColumn("email", Last(Col("owner_info.Item2"), true).Over(window)) .WithColumn("phone", Last(Col("owner_info.Item3"), true).Over(window)) .Filter(Col("line_type") == "03") // 只保留发明条目行 .WithColumn("invention", Trim(Col("raw_line").Substr(3, 100))) // 提取03行的发明名称 .Select("name", "email", "phone", "invention");
方式2:RDD分区遍历实现(性能更高)
如果要直接生成结构化RDD,可以用MapPartitions处理每个分区的行,因为单个分区内的行顺序是连续的,遍历过程中缓存最近遇到的01头信息,遇到03行直接输出关联后的结果即可,不需要全局排序,性能更好:
var structuredRdd = spark.TextFile("your_file_path") .MapPartitions(rows => { string currentName = null; string currentEmail = null; string currentPhone = null; var result = new List<(string, string, string, string)>(); foreach (var row in rows) { if (row.StartsWith("01")) { // 解析更新当前发明者信息 currentName = row.Substring(2, 30).Trim(); currentEmail = row.Substring(32, 70).Trim(); currentPhone = row.Substring(102).Trim(); } else if (row.StartsWith("03")) { var invention = row.Substring(2).Trim(); result.Add((currentName, currentEmail, currentPhone, invention)); } // 02行直接跳过 } return result; });
注意事项
- 如果你处理的是单个大文件拆分成多个分区的场景,要保证文件拆分后分区的顺序和原始文件一致,避免行顺序错乱导致关联错误
- 固定宽度截取的偏移量可以根据你的实际文件格式调整,上面的偏移量是基于你提供的样例对齐的
内容的提问来源于stack exchange,提问作者Bruno Moreira
相关产品推荐
相关产品推荐

