【问题标题】:Returning multiple columns from a single pyspark dataframe从单个 pyspark 数据帧返回多列
【发布时间】:2020-06-13 13:39:08
【问题描述】:

我正在尝试解析 pyspark 数据框的单列并获取具有多列的数据框。我的数据框如下:

   a  b               dic
0  1  2  {'d': 1, 'e': 2}
1  3  4  {'d': 7, 'e': 0}
2  5  6  {'d': 5, 'e': 4}

我想解析 dic 列并获取数据框,如下所示。如果可能,我期待使用 pandas UDF。我的预期输出如下:

   a  b  c  d
0  1  2  1  2
1  3  4  7  0
2  5  6  5  4

这是我的解决方案:

schema = StructType([
    StructField("c", IntegerType()),
    StructField("d", IntegerType())])

@pandas_udf(schema,PandasUDFType.GROUPED_MAP)
def do_someting(dic_col):
    return (pd.DataFrame(dic_col))

df.apply(add_json).show(10)

但这给出了错误'DataFrame' object has no attribute 'apply'

【问题讨论】:

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


    【解决方案1】:

    试试:

    #to convert pyspark df into pandas:
    df=df.toPandas()
    
    df["d"]=df["dic"].str.get("d")
    df["e"]=df["dic"].str.get("e")
    df=df.drop(columns=["dic"])
    

    返回:

       a  b  d  e
    0  1  2  1  2
    1  3  4  7  0
    2  5  6  5  4
    

    【讨论】:

    • 我想使用 spark 而不是通过将 pyspark 数据帧转换为 pandas 来将其作为 pandas UDF 来完成
    • pandas UDF 是什么意思?
    • @GrzegorzSkibinski link : spark.apache.org/docs/latest/api/python/… ,它的用户定义函数在群组中使用。它比转换为 pandas 数据帧更快,因为您可以将其应用于 spark 数据帧,但它比内置功能中的 spark 慢得多
    【解决方案2】:

    您可以先将单引号转换为双引号,然后使用from_json 将其转换为结构或映射列。

    如果你知道 dict 的架构,你可以这样做:

    data = [
        (1,   2,  "{'c': 1, 'd': 2}"),
        (3,   4,  "{'c': 7, 'd': 0}"),
        (5,   6,  "{'c': 5, 'd': 4}")
    ]
    
    df = spark.createDataFrame(data, ["a", "b", "dic"])
    
    schema = StructType([
        StructField("c", StringType(), True),
        StructField("d", StringType(), True)
    ])
    
    df = df.withColumn("dic", from_json(regexp_replace(col("dic"), "'", "\""), schema))
    
    df.select("a", "b", "dic.*").show(truncate=False)
    
    #+---+---+---+---+
    #|a  |b  |c  |d  |
    #+---+---+---+---+
    #|1  |2  |1  |2  |
    #|3  |4  |7  |0  |
    #|5  |6  |5  |4  |
    #+---+---+---+---+
    

    如果您不知道所有的键,您可以将其转换为映射而不是结构,然后将其分解并旋转以获取键作为列:

    df = df.withColumn("dic", from_json(regexp_replace(col("dic"), "'", "\""), MapType(StringType(), StringType())))\
           .select("a", "b", explode("dic"))\
           .groupBy("a", "b")\
           .pivot("key")\
           .agg(first("value"))
    

    【讨论】:

      猜你喜欢
      • 2023-03-24
      • 2021-06-20
      • 1970-01-01
      • 2019-12-22
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-02-15
      • 1970-01-01
      相关资源
      最近更新 更多