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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 07:06:01