【问题标题】:How to select last row and also how to access PySpark dataframe by index?如何选择最后一行以及如何按索引访问 PySpark 数据帧?
【发布时间】:2017-01-25 11:14:27
【问题描述】:

来自 PySpark SQL 数据框,例如

name age city
abc   20  A
def   30  B

如何获取最后一行。(如 df.limit(1) 我可以将第一行数据帧放入新数据帧)。

以及如何通过 index.like 行号访问数据帧行。 12 或 200 。

在熊猫中我可以做到

df.tail(1) # for last row
df.ix[rowno or index] # by index
df.loc[] or by df.iloc[]

我只是好奇如何以这种方式或替代方式访问 pyspark 数据帧。

谢谢

【问题讨论】:

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


    【解决方案1】:
    from pyspark.sql import functions as F
    
    expr = [F.last(col).alias(col) for col in df.columns]
    
    df.agg(*expr)
    

    提示:看起来您仍然具有使用 pandas 或 R 的思维方式。Spark 是我们处理数据方式的不同范例。您不再访问单个单元格内的数据,现在您可以处理整个数据块。如果你像刚才那样继续收集东西并采取行动,你就会失去 spark 提供的并行性的整个概念。看看 Spark 中转换与动作的概念。

    【讨论】:

    • 那个技巧让我觉得 spark 不是一个很好的工具来处理时间序列数据。对吗?
    • @yeliabsalohcin Spark 是处理时间序列的绝佳工具。但是,在这个问题中,用户想知道一种使用类似于 pandas 的索引来访问数据的方法,我的意思是,你应该注意 spark 与 pandas 或 R 的工作方式不同。
    【解决方案2】:

    使用以下内容获取包含单调递增、唯一、连续整数的索引列,这不是如何monotonically_increasing_id()工作。索引将按照与 DataFrame 的colName 相同的顺序升序。

    import pyspark.sql.functions as F
    from pyspark.sql.window import Window as W
    
    window = W.orderBy('colName').rowsBetween(W.unboundedPreceding, W.currentRow)
    
    df = df\
     .withColumn('int', F.lit(1))\
     .withColumn('index', F.sum('int').over(window))\
     .drop('int')\
    

    使用以下代码查看DataFrame的尾部或最后一个rownums

    rownums = 10
    df.where(F.col('index')>df.count()-rownums).show()
    

    使用以下代码查看 DataFrame 中从 start_rowend_row 的行。

    start_row = 20
    end_row = start_row + 10
    df.where((F.col('index')>start_row) & (F.col('index')<end_row)).show()
    

    zipWithIndex() 是一种 RDD 方法,它确实返回单调递增、唯一且连续的整数,但实现起来似乎要慢得多,因为您可以返回使用 id 列修改的原始 DataFrame。

    【讨论】:

      【解决方案3】:

      如何获取最后一行。

      如果您有一列可用于对数据框进行排序,例如“索引”,那么获取最后一条记录的一种简单方法是使用 SQL: 1)按降序排列你的桌子和 2) 从此订单中取第一个值

      df.createOrReplaceTempView("table_df")
      query_latest_rec = """SELECT * FROM table_df ORDER BY index DESC limit 1"""
      latest_rec = self.sqlContext.sql(query_latest_rec)
      latest_rec.show()
      

      以及如何通过 index.like 行号访问数据帧行。 12 或 200 。

      类似的方式你可以在任何行中获取记录

      row_number = 12
      df.createOrReplaceTempView("table_df")
      query_latest_rec = """SELECT * FROM (select * from table_df ORDER BY index ASC limit {0}) ord_lim ORDER BY index DESC limit 1"""
      latest_rec = self.sqlContext.sql(query_latest_rec.format(row_number))
      latest_rec.show()
      

      如果您没有“索引”列,您可以使用它来创建它

      from pyspark.sql.functions import monotonically_increasing_id
      
      df = df.withColumn("index", monotonically_increasing_id())
      

      【讨论】:

      • 非常感谢您解释清楚的答案。很高兴现在有了一种新方法。
      • monotonically_increasing_id() documentation "当前实现将分区ID放在高31位,每个分区内的记录号放在低33位。" i> 因此,对于可能跨不同分区存储的大型 DataFrame,这不会像您想象的那样起作用。您将无法引用 DataFrame 的最后一行,除非它全部在一个分区中。
      • @Clay 最后一部分比较补充。但是如果一个大的DataFrame真的很庞大,即不符合monotonically_increasing_id()的假设“数据帧有不到10亿个分区,每个分区有不到80亿条记录”,那么可以使用 sql ROW_NUMBER() OVER (PARTITION BY xxx ORDER BY yyy) 作为替代。
      • @DanyloZherebetskyy 我不是指每个分区有 >= 80 亿条记录的情况。 monotonically_increasing_id() 不保证 连续 个索引。因此,您不能使用它来创建“索引”列。如果这样做,您可能会使用filter() 查看第 20,000 行,而您使用monotonically_increasing_id() 创建的“索引”列中可能没有数字 20,000。我已经看到这个函数的输出从 526 跳到 28,622,因为它从第一个分区的末尾移动到了第二个分区的开头。
      • @Clay ,我可能看到了混乱的地方。 monotonically_increasing_id() 照它说的做。在问题的上下文中提供的 SQL 解决方案的美妙之处在于,它们不需要创建的 index 没有 jumps 或从 1 开始,而是单调增加(无论如何最后一个索引将具有最大值,第 12 大索引将具有第 12 大值,limit 负责处理)。如果需要在整个 DataFrame 上创建索引,其中索引从 1 开始并增加 1,则使用 SQL row_number() over (order by ...)
      【解决方案4】:

      如何获取最后一行。

      假定所有列都是可操作的,又长又丑:

      from pyspark.sql.functions import (
          col, max as max_, struct, monotonically_increasing_id
      )
      
      last_row = (df
          .withColumn("_id", monotonically_increasing_id())
          .select(max(struct("_id", *df.columns))
          .alias("tmp")).select(col("tmp.*"))
          .drop("_id"))
      

      如果不是所有列都可以排序,你可以试试:

      with_id = df.withColumn("_id", monotonically_increasing_id())
      i = with_id.select(max_("_id")).first()[0]
      
      with_id.where(col("_id") == i).drop("_id")
      

      注意。 pyspark.sql.functions/ `o.a.s.sql.functions 中有last 函数,但考虑到description of the corresponding expressions,这里不是一个好的选择。

      如何通过 index.like 访问数据框行

      你不能。 Spark DataFrame 并可通过索引访问。 You can add indices using zipWithIndex 并稍后过滤。请记住这个 O(N) 操作。

      【讨论】:

      • 嗨,目前我正在通过自动增量 ID 列添加方式或小 df 处理最后一行,我使用的是 toPandas().tail(1)。无论如何感谢您的回答。我所要求的数据帧的索引访问是因为,有时我必须替换列值(通过一些 col 值相等条件,并且为此我得到了 udf 的帮助)。但是如果我只想替换一个实例(特定的索引号行),那么我没有办法做到这一点。现在我可以按照建议使用“zipWithIndex”。谢谢。
      猜你喜欢
      • 2021-11-28
      • 2018-01-01
      • 2019-02-10
      • 2021-04-20
      • 2013-12-13
      • 2022-07-06
      • 2021-09-29
      • 2022-11-18
      相关资源
      最近更新 更多