【问题标题】:For each key/value pair in a dictionary, check if the value is of type pyspark Dataframe对于字典中的每个键/值对,检查值是否为 pyspark Dataframe 类型
【发布时间】:2021-04-28 08:20:28
【问题描述】:

我有一个数据框和一个字典列表如下 -

    x=spark.createDataFrame(["10","11","13"], "string").toDF("age")
    results = [
         {'type': 'check_datatype',
          'kwargs': {'table': x, 'columns': ['car_id','index'], 'd_type': 'str'},
          'datasource_path': '/cars_dataset_ok/',
          'Result': False},
        {'type': 'check_string_consistency',
          'kwargs': {'table': 'cars', 'columns': x, 'string_length': 6},
          'datasource_path': '/cars_dataset_ok/',
          'Result': False}
        ]

我想实现两件事-

1.
For each key/value pair in kwargs key, check if any value for all keys is of type pyspark Dataframe. If so, replace the value by the string "invalid dataframe"

2. For each key/value pair in kwargs key, check if value in 'table' key is of type pyspark Dataframe. If so, replace the value by the string "invalid dataframe"

预期输出 -

  1. 替换所有键的数据框值

最终结果1 =

[
             {'type': 'check_datatype',
              'kwargs': {'table': "invalid dataframe", 'columns': ['car_id','index'], 'd_type': 'str'},
              'datasource_path': '/cars_dataset_ok/',
              'Result': False},
            {'type': 'check_string_consistency',
              'kwargs': {'table': 'cars', 'columns': "invalid dataframe", 'string_length': 6},
              'datasource_path': '/cars_dataset_ok/',
              'Result': False}
            ]
  1. 在 kwargs 中的 key 为 'table' 的地方替换数据框

最终结果2=

    [
             {'type': 'check_datatype',
              'kwargs': {'table': "invalid dataframe", 'columns': ['car_id','index'], 'd_type': 'str'},
              'datasource_path': '/cars_dataset_ok/',
              'Result': False},
            {'type': 'check_string_consistency',
              'kwargs': {'table': 'cars', 'columns': x, 'string_length': 6},
              'datasource_path': '/cars_dataset_ok/',
              'Result': False}
            ]

【问题讨论】:

  • @mck 我创建了一个 pyspark 数据框“x”以供您理解。
  • 任何预期输出的例子?
  • 输出将仅是相同的字典列表,但数据框 x 将被字符串值“无效数据框”替换
  • @mck 我在 2 个场景中添加了预期的输出

标签: python dataframe apache-spark pyspark apache-spark-sql


【解决方案1】:

你可以使用一些列表/字典理解:

from pyspark.sql import DataFrame

result1 = [
    {
        k: v 
        if k != 'kwargs' 
        else {
            k2: "invalid dataframe" 
            if isinstance(v2, DataFrame) 
            else v2 
            for (k2, v2) in v.items()
        } 
        for (k, v) in d.items()
    } for d in results
]

print(result1)
# [{'type': 'check_datatype', 'kwargs': {'table': 'invalid dataframe', 'columns': ['car_id', 'index'], 'd_type': 'str'}, 'datasource_path': '/cars_dataset_ok/', 'Result': False}, {'type': 'check_string_consistency', 'kwargs': {'table': 'cars', 'columns': 'invalid dataframe', 'string_length': 6}, 'datasource_path': '/cars_dataset_ok/', 'Result': False}]

result2 = [
    {
        k: v 
        if k != 'kwargs' 
        else {
            k2: "invalid dataframe" 
            if k2 == 'table' and isinstance(v2, DataFrame) 
            else v2 
            for (k2, v2) in v.items()
        } 
        for (k, v) in d.items()
    } 
    for d in results
]

print(result2)
# [{'type': 'check_datatype', 'kwargs': {'table': 'invalid dataframe', 'columns': ['car_id', 'index'], 'd_type': 'str'}, 'datasource_path': '/cars_dataset_ok/', 'Result': False}, {'type': 'check_string_consistency', 'kwargs': {'table': 'cars', 'columns': DataFrame[age: string], 'string_length': 6}, 'datasource_path': '/cars_dataset_ok/', 'Result': False}]

请注意,x 在打印时将显示为DataFrame[age: string]。它应该仍然是原始数据框。

【讨论】:

  • 此类情况下的嵌套推导几乎会使代码不可读,尤其是当您在几个月后返回代码时。
  • 缩进会让理解更容易阅读。请检查我编辑的答案。
  • DataFrame[age: string] 仍在打印语句中。为什么 gettig 不替换成字符串值?
  • 你能澄清一下吗?您能看到在我的打印语句中,某些数据框在所需位置被invalid dataframe 替换了吗?
【解决方案2】:

您可以使用isinstance(variable,class) 来检查变量是否属于某种类型。 PySpark docs 表示 DataFrames 属于 pyspark.sql.DataFrame

所以我们需要做的就是循环使用isinstance进行检查

对于案例 1 首先创建新字典的深层副本

import copy

finalresult1 = copy.deepcopy(result)
for dic in finalresult1:
    for k,v in dic:
        if (isinstance(v,pyspark.sql.DataFrame)):
            dic[k] = "invalid dataframe"
        

对于案例 2,我们只需要检查表

    finalresult2 = copy.deepcopy(result)
    for dic in finalresult2:
        if (isinstance(dic['table'],pyspark.sql.DataFrame)):
            dic[k] = "invalid dataframe"

如果这引发错误,您可以使用 type(x) 为 DataFrame 找到正确的类

【讨论】:

  • 请注意,此答案不会将结果分配给任何东西。
  • 在编辑后的答案中创建了深层副本以防止更改原始字典。
  • 尝试运行您的代码并通过打印出变量来证明它给出了正确的结果
猜你喜欢
  • 1970-01-01
  • 2012-03-27
  • 2016-03-20
  • 2020-06-10
  • 2022-11-02
  • 2021-12-07
  • 1970-01-01
  • 2011-01-12
  • 1970-01-01
相关资源
最近更新 更多