【问题标题】:Using pyspark to create a segment array from a flat record使用 pyspark 从平面记录创建段数组
【发布时间】:2020-09-02 09:01:11
【问题描述】:

我有一个稀疏填充的表,其中包含唯一用户 ID 的各个段的值。我只需要创建一个包含 unique_id 和相关段标头的数组

请注意,这只是一个指示性数据集。我有数百个这样的片段。

------------------------------------------------
| user_id   | seg1 | seg2 | seg3 | seg4 | seg5 |
------------------------------------------------
| 100       |   M  |  null|   25 |  null|  30  |
| 200       |  null|  null|   43 |  null|  250 |
| 300       |   F  |  3000|  null|  74  |  null|
------------------------------------------------

我希望输出是

-------------------------------
| user_id| segment_array      |
-------------------------------
| 100    | [seg1, seg3, seg5] |
| 200    | [seg3, seg5]       |
| 300    | [seg1, seg2, seg4] |
-------------------------------

pyspark-sql 的 pyspark 中是否有可用的函数来完成此任务?

感谢您的帮助!

【问题讨论】:

    标签: arraylist pyspark apache-spark-sql record


    【解决方案1】:

    我找不到直接的方法,但你可以这样做。

    cols= df.columns[1:]
    
    r = df.withColumn('array', array(*[when(col(c).isNotNull(), lit(c)).otherwise('notmatch') for c in cols])) \
      .withColumn('array', array_remove('array', 'notmatch'))
    r.show()
    +-------+----+----+----+----+----+------------------+
    |user_id|seg1|seg2|seg3|seg4|seg5|             array|
    +-------+----+----+----+----+----+------------------+
    |    100|   M|null|  25|null|  30|[seg1, seg3, seg5]|
    |    200|null|null|  43|null| 250|      [seg3, seg5]|
    |    300|   F|3000|null|  74|null|[seg1, seg2, seg4]|
    +-------+----+----+----+----+----+------------------+
    

    【讨论】:

      【解决方案2】:

      不确定这是最好的方法,但我会这样攻击它:

      collect_set 函数将始终为您汇总的值列表提供唯一值。

      为每个段做一个联合:

      df_seg_1 = df.select(
        'user_id', 
        fn.when(
          col('seg1').isNotNull(), 
          lit('seg1)
        ).alias('segment')
      )
      # repeat for all segments
      
      df = df_seg_1.union(df_seg_2).union(...)
      
      df.groupBy('user_id').agg(collect_list('segment'))
      

      【讨论】:

      • 感谢您的回复汉斯。这无济于事,因为我有数百个片段:(
      • 它可能会影响速度,但您可以遍历 [col for col in df.columns if col.startswith('seg')]
      猜你喜欢
      • 1970-01-01
      • 2013-02-04
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-12-08
      • 1970-01-01
      相关资源
      最近更新 更多