【问题标题】:PySpark: Convert Map Column Keys Using DictionaryPySpark:使用字典转换地图列键
【发布时间】:2022-10-13 01:06:44
【问题描述】:

我有一个 PySpark DataFrame,其地图列如下:

root
 |-- id: long (nullable = true)
 |-- map_col: map (nullable = true)
 |    |-- key: string
 |    |-- value: double (valueContainsNull = true)

map_col 具有需要根据字典进行转换的键。例如,字典可能是:

mapping = {'a': '1', 'b': '2', 'c': '5', 'd': '8' }

因此,DataFrame 需要更改为:

[Row(id=123, map_col={'a': 0.0, 'b': -42.19}),
  Row(id=456, map_col={'a': 13.25, 'c': -19.6, 'd': 15.6})]

到以下:

[Row(id=123, map_col={'1': 0.0, '2': -42.19}),
  Row(id=456, map_col={'1': 13.25, '5': -19.6, '8': 15.6})]

如果我可以写出字典,我看到transform_keys 是一个选项,但它太大并且在工作流程的早期动态生成。我认为explode/pivot 也可以工作,但似乎表现不佳?

有任何想法吗?

编辑:添加了一点以显示map_colmap 的大小不统一。

【问题讨论】:

  • 你到底从哪里得到0.0-42.19等?当“映射”具有重复键时会发生什么?或者你将a重命名为b,而b已经存在于map_col中?
  • 它太大并且在工作流程的早期动态生成- 为什么这些很重要?您可以广播 dict 以使其跨执行者访问
  • 您的地图列是否始终包含相同数量的键?或者它至少受到一个已知数字的限制?
  • @OneCricketeer 我正在从早期的流程/作业中捕获整个 DataFrame。映射不会有重复的键(在map_colmapping 字典中。重命名的值也保证不会重叠。关于太大,我的意思是我对transform_key 的理解是它必须是作为expr 的一部分写出来。不过,我当然可以广播这本词典。
  • @Arseny 不-键是更大的一组独特可能性的某个子集-不一定是相同的长度。它们受到已知数量的限制 - 最多可能有大约 400 个左右的唯一键。

标签: apache-spark pyspark


【解决方案1】:

一种使用 RDD 转换的方法。

def updateKey(theDict, mapDict):
    """
    update theDict's key using mapDict
    """

    updDict = []
    for item in theDict.items():
        updDict.append((mapDict[item[0]] if item[0] in mapDict.keys() else item[0], item[1]))
    
    return dict(updDict)

data_sdf.rdd. 
    map(lambda r: (r[0], r[1], updateKey(r[1], mapping))). 
    toDF(['id', 'map_col', 'new_map_col']). 
    show(truncate=False)

# +---+-----------------------------------+-----------------------------------+
# |id |map_col                            |new_map_col                        |
# +---+-----------------------------------+-----------------------------------+
# |123|{a -> 0.0, b -> -42.19, e -> 12.12}|{1 -> 0.0, 2 -> -42.19, e -> 12.12}|
# |456|{a -> 13.25, c -> -19.6, d -> 15.6}|{8 -> 15.6, 1 -> 13.25, 5 -> -19.6}|
# +---+-----------------------------------+-----------------------------------+

【讨论】:

    【解决方案2】:

    transform_keys 可以使用lambda,如example 所示,不仅限于expr。但是,lambda 或 Python 可调用函数将需要使用在 pyspark.sql.functions 中定义的函数、Column 方法或 Scala UDF,因此不使用引用 mapping 字典对象的 Python UDF目前可以使用这种机制。但是,我们可以使用when 函数来应用映射,方法是将mapping 中的键值对展开为链式when 条件。请参阅下面的示例来说明这个想法:

    from typing import Dict, Callable
    from functools import reduce
    
    from pyspark.sql.functions import Column, when, transform_keys
    from pyspark.sql import SparkSession
    
    def apply_mapping(mapping: Dict[str, str]) -> Callable[[Column, Column], Column]:
    
        def convert_mapping_into_when_conditions(key: Column, _: Column) -> Column:
            initial_key, initial_value = mapping.popitem()
            initial_condition = when(key == initial_key, initial_value)
            return reduce(lambda x, y: x.when(key == y[0], y[1]), mapping.items(), initial_condition)
    
        return convert_mapping_into_when_conditions
    
    
    if __name__ == "__main__":
        spark = SparkSession
            .builder
            .appName("Temp")
            .getOrCreate()
        df = spark.createDataFrame([(1, {"foo": -2.0, "bar": 2.0})], ("id", "data"))
        mapping = {'foo': 'a', 'bar': 'b'}
        df.select(transform_keys(
            "data", apply_mapping(mapping)).alias("data_transformed")
                  ).show(truncate=False)
    

    上面的输出是:

    +---------------------+
    |data_transformed     |
    +---------------------+
    |{b -> 2.0, a -> -2.0}|
    +---------------------+
    

    这表明定义的映射 (foo -> a, bar -> b) 已成功应用于列。 apply_mapping 函数应该足够通用,可以在您自己的管道中复制和使用。

    【讨论】:

    • 这很聪明。非常好-感谢您的帮助!
    • 确定的事!实际上,这是一个有趣的问题。 :)
    猜你喜欢
    • 1970-01-01
    • 2017-06-26
    • 2021-05-21
    • 1970-01-01
    • 2023-03-12
    • 2021-01-15
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多