【问题标题】:SparkR dubt and Broken pipe exceptionSparkR dubt 和 Broken pipe 异常
【发布时间】:2016-01-11 10:06:03
【问题描述】:

您好,我正在使用纱线集群以分布式模式开发 SparkR。

我有两个问题:

1) 如果我制作了一个包含 R 行代码和 SparkR 行代码的脚本,它会只分发 SparkR 代码还是简单的 R 代码?

这是脚本。我读了一个 csv 并只记录了 100k 的第一条记录。 我清理它(使用 R 函数)删除 NA 值并创建了一个 SparkR 数据框。
这就是它的作用:foreach Lineset 获取 LineSet 出现的每个 TimeInterval 并对某些属性(数字属性)求和,然后将它们全部放入矩阵中。

这是带有 R 和 SparkR 代码的脚本。在独立模式下需要 7 小时,在分布式模式下需要 60 小时(被 java.net.SocketException: Broken Pipe 杀死)

LineSmsInt<-fread("/home/sentiment/Scrivania/LineSmsInt.csv")
Short<-LineSmsInt[1:100000,]
Short[is.na(Short)] <- 0
Short$TimeInterval<-Short$TimeInterval/1000
ShortDF<-createDataFrame(sqlContext,Short)
UniqueLineSet<-unique(Short$LINESET)
UniqueTime<-unique(Short$TimeInterval)
UniqueTime<-as.numeric(UniqueTime)
Row<-length(UniqueLineSet)*length(UniqueTime)
IntTemp<-matrix(nrow =Row,ncol=7)
k<-1

colnames(IntTemp)<-c("LINESET","TimeInterval","SmsIN","SmsOut","CallIn","CallOut","Internet")
Sys.time()
for(i in 1:length(UniqueLineSet)){
  SubSetID<-filter(ShortDF,ShortDF$LINESET==UniqueLineSet[i])
  for(j in 1:length(UniqueTime)){
    SubTime<-filter(SubSetID,SubSetID$TimeInterval==UniqueTime[j])       
    IntTemp[k,1]<-UniqueLineSet[i]
    IntTemp[k,2]<-as.numeric(UniqueTime[j])
    k3<-collect(select(SubTime,sum(SubTime$SmsIn)))
    IntTemp[k,3]<-k3[1,1]
    k4<-collect(select(SubTime,sum(SubTime$SmsOut)))
    IntTemp[k,4]<-k4[1,1]
    k5<-collect(select(SubTime,sum(SubTime$CallIn)))
    IntTemp[k,5]<-k5[1,1]
    k6<-collect(select(SubTime,sum(SubTime$CallOut)))
    IntTemp[k,6]<-k6[1,1]
    k7<-collect(select(SubTime,sum(SubTime$Internet)))
    IntTemp[k,7]<-k7[1,1]
    k<-k+1
  }
  print(UniqueLineSet[i])
  print(i)
}

这是脚本 R,唯一改变的是子集函数,当然是普通的 R data.frame 而不是 SparkR 数据帧。 在独立模式下需要 1.30 分钟。
为什么它在 R 中如此之快,而在 SparkR 中却如此缓慢?

for(i in 1:length(UniqueLineSet)){
  SubSetID<-subset.data.frame(LineSmsInt,LINESET==UniqueLineSet[i])
  for(j in 1:length(UniqueTime)){
    SubTime<-subset.data.frame(SubSetID,TimeInterval==UniqueTime[j])
    IntTemp[k,1]<-UniqueLineSet[i]
    IntTemp[k,2]<-as.numeric(UniqueTime[j])
    IntTemp[k,3]<-sum(SubTime$SmsIn,na.rm = TRUE)
    IntTemp[k,4]<-sum(SubTime$SmsOut,na.rm = TRUE)
    IntTemp[k,5]<-sum(SubTime$CallIn,na.rm = TRUE)
    IntTemp[k,6]<-sum(SubTime$CallOut,na.rm = TRUE)
    IntTemp[k,7]<-sum(SubTime$Internet,na.rm=TRUE)
    k<-k+1
  }
  print(UniqueLineSet[i])
  print(i)
}

2) 分布式模式下的第一个脚本被以下人员杀死:

java.net.SocketException: 断管

有时也会出现这种情况:

java.net.SocketTimeoutException: 接受超时

可能是因为配置不好?建议?

谢谢。

【问题讨论】:

    标签: r apache-spark sparkr


    【解决方案1】:

    不要误以为这是一段写得不好的代码。使用核心 R 已经是低效的,并且将 SparkR 添加到等式中会使情况变得更糟。

    如果我制作了一个包含 R 行代码和 SparkR 行代码的脚本,它会只分发 SparkR 代码还是简单的 R 代码?

    除非您使用分布式数据结构和对这些结构进行操作的函数,否则它只是在主服务器上的单个线程中执行的纯 R 代码。

    为什么它在 R 中如此之快,而在 SparkR 中却如此缓慢?

    对于初学者,您为LINESETUniqueTime 和列的每个组合执行一个作业。每次 Spark 扫描所有记录并获取数据给驱动程序。

    此外,使用 Spark 处理可以在单机内存中轻松处理的数据根本没有意义。在这种情况下运行作业的成本通常远高于实际处理的成本。

    建议?

    如果您真的想使用 SparkR,只需 groupByagg

    group_by(Short, Short$LINESET, Short$TimeInterval) %>% agg(
      sum(Short$SmsIn), sum(Short$SmsOut), sum(Short$CallIn),
      sum(Short$CallOut), sum(Short$Internet))
    

    如果您担心缺少(LINESETTimeInterval)对,请使用joinunionAll 填写这些。

    在实践中,它会简单地使用 data.table 并在本地聚合:

    Short[, lapply(.SD, sum, na.rm=TRUE), by=.(LINESET, TimeInterval)]
    

    【讨论】:

    • 非常感谢。我尝试了这个差异函数(我不知道),它当然可以工作。现在速度挺快的。我是 sparkR 的新手(也是 spark)。你有一些关于 sparkR 的链接(通过 sparkR 官方文档)吗?
    • 不是真的,但我不使用 SparkR,所以我从来没有寻找过这样的东西。不过,您可以使用任何通用的 Spark / SparkSQL 参考。这里没有区别。
    • 我还有一个 dubt,createDataFrame 是分布在集群上的吗?如果我从 hdfs 或简单的 csv 读取文件有什么区别?(对于我的意思是分布)
    • 是的,有区别。一般来说,除了原型设计和测试之外,您应该避免将本地数据结构并行化。
    猜你喜欢
    • 2014-01-29
    • 2011-02-16
    • 2012-04-13
    • 2018-10-29
    • 1970-01-01
    • 2020-01-24
    • 2014-05-22
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多