【问题标题】:how to add Row id in pySpark dataframes [duplicate]如何在pySpark数据框中添加行ID [重复]
【发布时间】:2015-11-12 05:25:56
【问题描述】:

我有一个 csv 文件;我在pyspark中转换为DataFrame(df);经过一些改造;我想在df中添加一列;这应该是简单的行 id(从 0 或 1 到 N)。

我在 rdd 中转换了 df 并使用“zipwithindex”。我将生成的 rdd 转换回 df。这种方法有效,但它生成了 25 万个任务并且需要大量执行时间。我想知道是否有其他方法可以减少运行时间。

以下是我的代码的 sn-p;我正在处理的 csv 文件很大;包含数十亿行。

debug_csv_rdd = (sc.textFile("debug.csv")
  .filter(lambda x: x.find('header') == -1)
  .map(lambda x : x.replace("NULL","0")).map(lambda p: p.split(','))
  .map(lambda x:Row(c1=int(x[0]),c2=int(x[1]),c3=int(x[2]),c4=int(x[3]))))

debug_csv_df = sqlContext.createDataFrame(debug_csv_rdd)
debug_csv_df.registerTempTable("debug_csv_table")
sqlContext.cacheTable("debug_csv_table")

r0 = sqlContext.sql("SELECT c2 FROM debug_csv_table WHERE c1 = 'str'")
r0.registerTempTable("r0_table")

r0_1 = (r0.flatMap(lambda x:x)
    .zipWithIndex()
    .map(lambda x: Row(c1=x[0],id=int(x[1]))))

r0_df=sqlContext.createDataFrame(r0_2)
r0_df.show(10) 

【问题讨论】:

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


    【解决方案1】:

    您也可以使用 sql 包中的函数。它将生成一个唯一的 id,但它不会是连续的,因为它取决于分区的数量。我相信它在 Spark 1.5 + 中可用

    from pyspark.sql.functions import monotonicallyIncreasingId
    
    # This will return a new DF with all the columns + id
    res = df.withColumn("id", monotonicallyIncreasingId())
    

    编辑:2017 年 1 月 19 日

    正如@Sean评论的那样

    使用 monotonically_increasing_id() 代替 Spark 1.6 及更高版本

    【讨论】:

    • 在我的场景中,生成的行 ID 都是相同的值。全部 1.
    • @Matthias 它对我很有效
    • 使用 monotonically_increasing_id 代替 Spark 1.6 +
    • @Matthias 查看编辑
    • 请注意,这并不能真正回答问题,因为 OP 要求 0-(N-1) 索引。 monotonically_increasing_id 不保证连续的 id。
    猜你喜欢
    • 2019-02-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-06-26
    • 2018-06-16
    • 2019-07-17
    • 1970-01-01
    • 2019-04-09
    相关资源
    最近更新 更多