【问题标题】:Inserting records in a spark dataframe在火花数据框中插入记录
【发布时间】:2018-03-26 22:47:00
【问题描述】:

我在 pyspark 中有一个数据框。这是它的样子,

+---------+---------+
|timestamp| price   |
+---------+---------+
|670098928|  50     |
|670098930|  53     |
|670098934|  55     |
+---------+---------+

我想用之前的状态来填补时间戳的空白,这样我就可以得到一个完美的集合来计算时间加权平均值。输出应该是这样的 -

+---------+---------+
|timestamp| price   |
+---------+---------+
|670098928|  50     |
|670098929|  50     | 
|670098930|  53     |
|670098931|  53     |
|670098932|  53     |
|670098933|  53     |
|670098934|  55     |
+---------+---------+

最终,我想将这个新数据帧保存在磁盘上并可视化我的分析。

如何在 pyspark 中执行此操作? (为简单起见,我只保留了 2 列。在填补空白之前,我的实际数据框有 89 列,约 6.7 亿条记录。)

【问题讨论】:

  • 你可以用 scipy 进行插值。我不太确定 PySpark 能做你想做的事
  • @cricket_007 spark 无法做到这一点。 Veenit,我不知道你为什么要这么做?
  • @eliasah 我正在尝试创建一个数据框,其中包含每个时间戳(最低级别粒度)的记录,这样如果我想进行时间加权平均,这非常方便。

标签: apache-spark pyspark


【解决方案1】:

您可以生成时间戳范围、展平它们并选择行

import pyspark.sql.functions as func

from pyspark.sql.types import IntegerType, ArrayType


a=sc.parallelize([[670098928, 50],[670098930, 53], [670098934, 55]])\
.toDF(['timestamp','price'])

f=func.udf(lambda x:range(x,x+5),ArrayType(IntegerType()))

a.withColumn('timestamp',f(a.timestamp))\
.withColumn('timestamp',func.explode(func.col('timestamp')))\
.groupBy('timestamp')\
.agg(func.max(func.col('price')))\
.show()

+---------+----------+
|timestamp|max(price)|
+---------+----------+
|670098928|        50|
|670098929|        50|
|670098930|        53|
|670098931|        53|
|670098932|        53|
|670098933|        53|
|670098934|        55|
|670098935|        55|
|670098936|        55|
|670098937|        55|
|670098938|        55|
+---------+----------+

【讨论】:

  • 当我执行f=func.udf(lambda x:range(x,x+5),ArrayType(IntegerType()))时得到AttributeError: 'JavaMember' object has no attribute 'parseDataType'
  • 没有。它没有。您使用的是哪个版本的 Spark?我在 2.0.0
  • 我在 1.6.0,但如果你不能定义一个简单的 udf,那么你的环境就有问题。
  • 可以去掉udf,用RDD上的map替换,把a.withColumn('timestamp',f(a.timestamp))\替换成a.map(lambda row:( range(row[0],row[0]+5),row[1])).toDF(['timestamp','price'])\
  • 以上代码有效。但是,并不能从本质上解决我的问题。在 UDF 中,硬编码 x+5。那么,如果两个数字之间的差距大于5 怎么办?对此的一个答案是将5 替换为 Integer.MAX_VALUE 但随后需要对最后一个数字进行约束。最终,时间戳是“时间”,所以我想在“秒”或“毫秒”上分解它
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-09-19
  • 2016-04-10
  • 2021-09-22
  • 1970-01-01
  • 1970-01-01
  • 2021-11-22
相关资源
最近更新 更多