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

基于首列关键字将文本文件数据导入Spark SQL单表不同列

使用Spark SQL合并不同类型记录到单行

需求分析

原始文本按首列的record_type(1/2/3)区分记录类型,同一id(第二列)的不同类型记录需要合并到同一行的对应列中:

  • record_type=1:提供基础信息(国家、名称、编码)
  • record_type=2:提供会员信息(会员编码、会员类型、是否激活)
  • record_type=3:提供额外数据(部分id无此记录)

原始数据示例

1|16255|usa|Test||TEST806282|||||
2|16255|806282|company_member|False|true
3|16255|my_data
1|18299|usa|Test||TEST260092|||||
2|18299|260092|company_member|False|False

解决方案

方法1:拆分后关联(直观易理解)

  1. 读取并解析原始数据
import org.apache.spark.sql.types.{IntegerType, BooleanType}

// 读取文本文件,按|分割列
val rawDF = spark.read.text("/path/to/your/data.txt")
  .select(split($"value", "\\|").alias("cols"))
  .select(
    $"cols"(0).cast(IntegerType).alias("record_type"),
    $"cols"(1).cast(IntegerType).alias("id"),
    $"cols"(2).alias("col3"),
    $"cols"(3).alias("col4"),
    $"cols"(5).alias("col6")
  )
  1. 拆分不同类型的数据集
// 提取类型1的基础信息
val type1DF = rawDF.filter($"record_type" === 1)
  .select($"id", $"col3".alias("country"), $"col4".alias("name"), $"col6".alias("code"))

// 提取类型2的会员信息
val type2DF = rawDF.filter($"record_type" === 2)
  .select(
    $"id",
    $"col3".alias("member_code"),
    $"col4".alias("member_type"),
    $"col6".cast(BooleanType).alias("is_active")
  )

// 提取类型3的额外数据
val type3DF = rawDF.filter($"record_type" === 3)
  .select($"id", $"col3".alias("extra_data"))
  1. 关联合并数据集
// 用inner join保证同时有类型1和2的记录,left join保留无类型3的记录
val resultDF = type1DF
  .join(type2DF, Seq("id"), "inner")
  .join(type3DF, Seq("id"), "left")

// 查看结果
resultDF.show()

方法2:使用Pivot(更简洁高效)

利用Spark的pivot功能直接按record_type转置列:

val pivotResult = rawDF.groupBy("id")
  .pivot("record_type")
  .agg(
    first($"col3").alias("col3"),
    first($"col4").alias("col4"),
    first($"col6").alias("col6")
  )
  .select(
    $"id",
    $"1_col3".alias("country"),
    $"1_col4".alias("name"),
    $"1_col6".alias("code"),
    $"2_col3".alias("member_code"),
    $"2_col4".alias("member_type"),
    $"2_col6".cast(BooleanType).alias("is_active"),
    $"3_col3".alias("extra_data")
  )

pivotResult.show()

输出结果

idcountrynamecodemember_codemember_typeis_activeextra_data
16255usaTestTEST806282806282company_membertruemy_data
18299usaTestTEST260092260092company_memberfalsenull

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 19:43:03