【问题标题】:Assign value to specific cell in PySpark dataFrame为 PySpark 数据帧中的特定单元格赋值
【发布时间】:2018-10-27 20:04:06
【问题描述】:

我想使用PySpark 更改我的Spark DataFrame 的特定单元格中的值。

简单示例 - 我创建了一个模拟 Spark DataFrame

df = spark.createDataFrame(
    [
     (1, 1.87, 'new_york'), 
     (4, 2.76, 'la'), 
     (6, 3.3, 'boston'), 
     (8, 4.1, 'detroit'), 
     (2, 5.70, 'miami'), 
     (3, 6.320, 'atlanta'), 
     (1, 6.1, 'houston')
    ],
    ('variable_1', "variable_2", "variable_3")
)

Runnning display(df) 我得到这张表:

variable_1   variable_2   variable_3
    1           1.87    new_york
    4           2.76    la
    6           3.3     boston
    8           4.1     detroit
    2           5.7     miami
    3           6.32    atlanta
    1           6.1     houston

例如,我想为第 4 行第 3 列的单元格分配一个新值,即将detroit 更改为new_orleans。我知道df.iloc[4, 3] = 'new_orleans'df.loc[4, 'detroit'] = 'new_orleans' 的分配在Spark 中无效。

使用when 对我的问题的有效答案是:

from pyspark.sql.functions import when
targetDf = df.withColumn("variable_3", \
              when(((df["variable_1"] == 8) & (df["variable_2"] == 4.1)) , 'new_orleans').otherwise(df["variable_3"]))

我的问题是:这是否可以在PySpark 中以更实用的方式完成,而无需输入我只想更改 1 个单个单元格的行的所有值和列名(可能在不使用的情况下实现相同when 函数)?

提前感谢您的帮助和@useruser9806664 的反馈。

【问题讨论】:

    标签: python apache-spark dataframe pyspark


    【解决方案1】:

    考虑使用 Pandas DataFrame

    Spark DataFrame 确实是不可变的,因此它们不是为修改而设计的。 Spark Dataframes 是为处理大量数据而优化的分布式数据集合,如果您想要进行任何更改,您必须创建一个新的并进行您想要的修改。

    尽管如此,有时您可能需要修改特定行的特定单元格。对于这些情况,您可以使用 when 函数(就像您在示例中所做的那样)修改列,其中单元格的值与您要修改的特定单元格位于同一行。或者您可以考虑将 Spark DataFrame 转换为 Pandas DataFrame可变),在将新值分配给相关单元格后,将其转换回来进入 Spark DataFrame。你可以这样做:

    # Copy the schema of your Spark dataframe 
    schema = df.schema
    
    # Create Pandas Dataframe using your Spark DataFrame
    pandas_df = df.toPandas()
    
    # Assign the new value to the specific cell (you could use .at or .loc)
    pandas_df.at[3, 'variable_3'] = 'new_orleans'
    
    # Update your dataframe with the new value using the Pandas DataFrame
    df = spark.createDataFrame(pandas_df,schema=schema)
    
    # Delete the auxiliary pandas dataframe to free memory for other uses
    del pandas_df
    

    请记住,Pandas DataFrame 不是分布式的,处理大量数据时 Pandas DataFrame 中的处理速度会变慢。

    【讨论】:

      【解决方案2】:

      您可以使用底层 RDD 创建行号:

      from pyspark.sql import Row
      
      # Function to update dataframe row with a rownumber
      def create_rownum(ziprow):
          row, index=ziprow
          row=row.asDict()
          row['rownum']= index
          return(Row(**row))
      
      # First create a rownumber then add to dataframe
      df.rdd.zipWithIndex().map(create_rownum).toDF().show()
      

      现在您可以过滤 DataFrame 以获得您想要的行号。

      【讨论】:

        【解决方案3】:

        Spark DataFrames 不可变不提供随机访问,严格来说,无序。结果:

        • 您不能分配任何东西(因为属性不可变)。
        • 您无法访问特定行(因为没有随机访问)。
        • 行“索引”定义不明确(因为无序)。

        您可以做的是用新列创建一个新的数据框,替换现有的,使用一些条件表达式,您找到的答案已经涵盖了这些。

        另外,monotonically_increasing_id 不添加索引(行号)。它添加单调递增的数字,不一定是连续数字或从任何特定值开始(在空分区的情况下)。

        【讨论】:

        • 感谢@user9806664 的回答。阅读标题为 SO 帖子的“Scala/Spark:不可变数据帧和内存”时,我意识到 Spark 的不变性。不过,是否可以创建一个新的数据框来执行您以另一种形式提到的替换(不输入我们要修改 1 个单元格的行中每一列的所有值)?如果是这样,你能举个例子吗?感谢您的回答,我将编辑我的问题以使其更清楚。
        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 2016-04-22
        • 2011-01-24
        • 1970-01-01
        • 1970-01-01
        • 2019-09-03
        • 2013-12-11
        • 2020-11-12
        相关资源
        最近更新 更多