【发布时间】:2015-11-27 14:15:53
【问题描述】:
我正在使用 Spark/Scala 编写一个应用程序,我需要在其中计算一列的指数移动平均值。
EMA_t = (price_t * 0.4) + (EMA_t-1 * 0.6)
我面临的问题是我需要同一列的先前计算的值(EMA_t-1)。通过 mySQL,这可以通过使用 MODEL 或通过创建一个 EMA 列来实现,然后您可以更新每行的行,但我已经尝试过了,并且既不能使用 Spark SQL 也不能使用 Hive 上下文......有什么办法我可以访问这个 EMA_t-1?
我的数据如下所示:
timestamp price
15:31 132.3
15:32 132.48
15:33 132.76
15:34 132.66
15:35 132.71
15:36 132.52
15:37 132.63
15:38 132.575
15:39 132.57
所以我需要添加一个新列,其中我的第一个值只是第一行的价格,然后我需要使用以前的值:EMA_t = (price_t * 0.4) + (EMA_t-1 * 0.6)计算该列中的以下行。 我的 EMA 列必须是:
EMA
132.3
132.372
132.5272
132.58032
132.632192
132.5873152
132.6043891
132.5926335
132.5835801
我目前正在尝试使用 Spark SQL 和 Hive 来实现它,但如果可以以其他方式实现它,这将同样受欢迎!我也想知道如何使用 Spark Streaming 做到这一点。我的数据在数据框中,我使用的是 Spark 1.4.1。
非常感谢您提供的任何帮助!
【问题讨论】:
-
我认为您的用例不太适合大数据环境,因为它的功能是并行处理数据,而您的用例不允许这样做...您有多少记录有吗?
-
@mark91 每个数据集大约有 100 000 行,需要分析大约 100 个数据集。我需要计算这个的原因是我需要这个作为输入来计算一个特征。我需要使用这些功能来训练带有随机森林的模型。我的模型必须根据几个特征的值来预测价格是上涨还是下跌。我还需要在未来实现这一点。
-
那么在我看来,最好的选择是分配数据,例如每个数据集都在一个分区中(每个分区都包含一个数据集),然后您独立处理每个分区..
标签: scala apache-spark hive apache-spark-sql spark-dataframe