【问题标题】:How do I flattern a pySpark dataframe by one array column? [duplicate]如何按一个数组列展平 pySpark 数据帧? [复制]
【发布时间】:2016-08-14 00:29:23
【问题描述】:

我有一个像这样的 spark 数据框:

+------+--------+--------------+--------------------+
|   dbn|    boro|total_students|                sBus|
+------+--------+--------------+--------------------+
|17K548|Brooklyn|           399|[B41, B43, B44-SB...|
|09X543|   Bronx|           378|[Bx13, Bx15, Bx17...|
|09X327|   Bronx|           543|[Bx1, Bx11, Bx13,...|
+------+--------+--------------+--------------------+

如何展平它以便为 sBus 中的每个元素复制每一行,并且 sBus 将是一个普通的字符串列?

所以结果会是这样的:

+------+--------+--------------+--------------------+
|   dbn|    boro|total_students|                sBus|
+------+--------+--------------+--------------------+
|17K548|Brooklyn|           399| B41                |
|17K548|Brooklyn|           399| B43                |
|17K548|Brooklyn|           399| B44-SB             |
+------+--------+--------------+--------------------+

等等……

【问题讨论】:

  • 你能提供预期的输出吗?您是否期望得到sBus 和sSw 之间的笛卡尔积?
  • 谢谢!添加了预期的结果。为简单起见,删除了 sSw 列
  • 好吧,你可以使用explode(例如stackoverflow.com/q/36484385/1560062),但如果你有多个列,那就没那么简单了。

标签: python apache-spark pyspark


【解决方案1】:

如果不将其转换为 RDD,我想不出办法。

# convert df to rdd
rdd = df.rdd

def extract(row, key):
    """Takes dictionary and key, returns tuple of (dict w/o key, dict[key])."""
    _dict = row.asDict()
    _list = _dict[key]
    del _dict[key]
    return (_dict, _list)


def add_to_dict(_dict, key, value):
    _dict[key] = value
    return _dict


# preserve rest of values in key, put list to flatten in value
rdd = rdd.map(lambda x: extract(x, 'sBus'))
# make a row for each item in value
rdd = rdd.flatMapValues(lambda x: x)
# add flattened value back into dictionary
rdd = rdd.map(lambda x: add_to_dict(x[0], 'sBus', x[1]))
# convert back to dataframe
df = sqlContext.createDataFrame(rdd)

df.show()

棘手的部分是将其他列与新展平的值保持在一起。我通过将每一行映射到(dict of other columns, list to flatten) 的元组然后调用flatMapValues 来做到这一点。这会将值列表的每个元素拆分为单独的行,但保持键附加,即

(key, ['A', 'B', 'C'])

变成

(key, 'A')
(key, 'B')
(key, 'C')

然后,我将展平的值移回其他列的字典中,并将其重新转换回 DataFrame。

【讨论】:

    猜你喜欢
    • 2020-05-15
    • 2020-01-17
    • 1970-01-01
    • 1970-01-01
    • 2021-05-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多