【问题标题】:How do I group by hour in SparkR?如何在 SparkR 中按小时分组?
【发布时间】:2017-08-19 07:32:47
【问题描述】:

我正在尝试使用 SparkR 和 Spark 2.1.0 按小时汇总一些日期。 我的数据如下:

                       created_at
1  Sun Jul 31 22:25:01 +0000 2016
2  Sun Jul 31 22:25:01 +0000 2016
3  Fri Jun 03 10:16:57 +0000 2016
4  Mon May 30 19:23:55 +0000 2016
5  Sat Jun 11 21:00:07 +0000 2016
6  Tue Jul 12 16:31:46 +0000 2016
7  Sun May 29 19:12:26 +0000 2016
8  Sat Aug 06 11:04:29 +0000 2016
9  Sat Aug 06 11:04:29 +0000 2016
10 Sat Aug 06 11:04:29 +0000 2016

我希望输出是:

Hour      Count
22         2
10         1
19         1
11         3
....

我试过了:

sumdf <- summarize(groupBy(df, df$created_at), count = n(df$created_at))
head(select(sumdf, "created_at", "count"),10)

但分组到最接近的秒数:

                       created_at count
1  Sun Jun 12 10:24:54 +0000 2016     1
2  Tue Aug 09 14:12:35 +0000 2016     2
3  Fri Jul 29 19:22:03 +0000 2016     2
4  Mon Jul 25 21:05:05 +0000 2016     2

我试过了:

sumdf <- summarize(groupBy(df, hr=hour(df$created_at)), count = n(hour(df$created_at)))
head(select(sumdf, "hour(created_at)", "count"),20)

但这给出了:

  hour(created_at) count
1               NA     0

我试过了:

sumdf <- summarize(groupBy(df, df$created_at), count = n(hour(df$created_at)))
head(select(sumdf, "created_at", "count"),10)

但这给出了:

                       created_at count
1  Sun Jun 12 10:24:54 +0000 2016     0
2  Tue Aug 09 14:12:35 +0000 2016     0
3  Fri Jul 29 19:22:03 +0000 2016     0
4  Mon Jul 25 21:05:05 +0000 2016     0
...

如何使用小时功能来实现这一点,或者有更好的方法吗?

【问题讨论】:

  • 您应该尝试将您的created_at 列拆分为hour 值,然后在hour 列上使用groupBy。

标签: r apache-spark sparkr


【解决方案1】:

我会用to_timestamp(Spark 2.2)或unix_timestamp %&gt;% cast("timestamp")(早期版本)解析日期并访问hour

df <- createDataFrame(data.frame(created_at="Sat Aug 19 12:33:26 +0000 2017"))
head(count(group_by(df, 
  alias(hour(to_timestamp(column("created_at"), "EEE MMM d HH:mm:ss Z yyyy")), "hour")
)))
##  hour count
## 1   14     1

【讨论】:

  • 谢谢。我在使用 1.6.3 时做了一些更改: df % cast("timestamp")), "hour") ))) 知道如何克服:错误在 as.POSIXlt.default(x, tz = tz(x)) 中:不知道如何将“x”转换为“POSIXlt”类?
【解决方案2】:

假设您的本地表是df,这里真正的问题是从created_at 列中提取小时,然后使用您的分组代码。为此,您可以使用dapply

library(SparkR)
sc1 <- sparkR.session()
df2 <- createDataFrame(df)

#with dapply you need to specify the schema i.e. the data.frame that will come out
#of the applied function - i.e. substringDF in our case
schema <- structType(structField('created_at', 'string'), structField('time', 'string'))

#a function that will be applied to each partition of the spark data frame.
#remember that each partition is a data.frame itself.
substringDF <- function(DF) {

 DF$time <- substr(DF$created_at, 15, 16)

 DF

}

#and then we use the above in dapply
df3 <- dapply(df2, substringDF, schema)
head(df3)
#                        created_at time
#1 1  Sun Jul 31 22:25:01 +0000 2016   22
#2 2  Sun Jul 31 22:25:01 +0000 2016   22
#3 3  Fri Jun 03 10:16:57 +0000 2016   10
#4 4  Mon May 30 19:23:55 +0000 2016   19
#5 5  Sat Jun 11 21:00:07 +0000 2016   21
#6 6  Tue Jul 12 16:31:46 +0000 2016   16

然后只需应用您的正常分组代码:

sumdf <- summarize(groupBy(df3, df3$time), count = n(df3$time))
head(select(sumdf, "time", "count"))
#  time count
#1   11     3
#2   22     2
#3   16     1
#4   19     2
#5   10     1
#6   21     1

【讨论】:

  • 谢谢。我应该提到我正在使用 Spark 1.6.3,我认为 dapply 不存在。
  • 我决定使用 Spark2 并尝试一下您的解决方案。 df 已经是一个火花数据框。因此,我使用您的代码减去 forst 3 行并将 df2 更改为 df 并将 df3 更改为 df2。 head(df3) 给出:ERROR Executor: Exception in task 0.0 in stage 10.0 (TID 12) java.lang.ArrayIndexOutOfBoundsException: 2 at org.apache.spark.sql.api.r.SQLUti 有什么想法吗?
  • 如果您将我的df3 更改为df2,您应该使用head(df2) 对吗?另外,尝试使用我的代码,看看它是否有效(即命名你的 spark data.frame df2)。我在独立版本和 YARN 版本上都试过了,它们都工作正常。您确定 Spark 2 已正确安装和配置吗?
  • 是的head(df2)。我觉得Spark2还可以,用了一段时间。这很好用:head(substringDF(df))。我将尝试重新启动 Spark2。
  • 因此,可能是架构的情况。您需要告诉dapply sparkdata.frame 的结构是什么,即它需要知道列名和列的数据类型(例如,在上述情况下,列名timestring)。确保您已正确设置。
【解决方案3】:

这里是代码SCALA,我想你可以参考一下。

    var index = ss.sparkContext.parallelize( Seq(
  (1,"Sun Jul 31 22:25:01 +0000 2016"),
  (2,"Sun Jul 31 22:25:01 +0000 2016"),
  (3,"Fri Jun 03 10:16:57 +0000 2016"),
  (4,"Mon May 30 19:23:55 +0000 2016"),
  (5,"Sat Jun 11 21:00:07 +0000 2016"),
  (6,"Tue Jul 12 16:31:46 +0000 2016"),
  (7,"Sun May 29 19:12:26 +0000 2016"),
  (8,"Sat Aug 06 11:04:29 +0000 2016"),
  (9,"Sat Aug 06 11:04:29 +0000 2016"),
  (10,"Sat Aug 06 11:04:29 +0000 2016"))
).toDF("ID", "time")

val getHour = udf( (s : String) => {
  s.substring( 11, 13)
})
index.withColumn("hour", getHour($"time")).groupBy( "hour").agg( count("*").as("count")).show

【讨论】:

  • 这是题外话,即使对于 Scala 也是不好的做法。我相信 OP 对 sparkR 很准确
  • 是的,所以我说“这里是代码 SCALA,我想你可以参考一下。”
  • 这仍然是糟糕的代码@Robin。一个人应该接受批评。我记得我没有侮辱你。我没有像其他人那样对你投反对票。
  • 我不懂SparkR,我懂一点Scala。所以我把我的想法放在这里,希望能给提问者一些启发。
  • 这不是问题。但是,如果您想查看您的 Scala 代码,这里是:将日期视为字符串是不好的做法。如果格式发生变化,您应该使用日期库和解析器。
猜你喜欢
  • 2019-03-31
  • 1970-01-01
  • 2011-06-27
  • 2019-11-24
  • 2011-01-07
  • 2021-10-26
  • 1970-01-01
  • 1970-01-01
  • 2018-03-08
相关资源
最近更新 更多