【发布时间】:2017-01-30 18:16:50
【问题描述】:
所以,我在 Spark 中遇到了臭名昭著的 Task Not Serializable 错误。这是相关的代码块:
val labeledPoints: RDD[LabeledPoint] = events.map(event => {
var eventsPerEntity = try {
HBaseHelper.scan(...filter entity here...)(sc).map(newEvent => {
Try(new Object(...))
}).filter(_.isSuccess).map(_.get)
} catch {
case e: Exception => {
logger.error(s"Failed to convert event ${event}." +
s"Exception: ${e}.")
throw e
}
}
})
基本上我想要实现的是我正在访问sc,这是我在map 中的Spark Context 对象。在运行时,我收到Task Not Serializable 错误。
这是我能想到的一个潜在解决方案:
在没有sc 的情况下查询 HBase,我可以这样做,但反过来我会得到一个列表。 (如果我尝试并行化;我必须再次使用sc)。有一个列表会导致我无法使用reduceByKey,在我的另一个问题中建议here。所以我也无法成功实现这一目标,因为我不知道如果没有reduceByKey,我将如何实现this。另外我真的很想使用 RDD :)
所以我正在寻找另一种解决方案 + 询问我是否做错了什么。提前致谢!
更新
所以基本上,我的问题变成了这样:
我有一个名为events 的RDD。这是整个 HBase 表。注意:每个event 都由performerId 执行,这又是event 中的一个字段,即event.performerId。
对于events 中的每个event,我需要计算event.numericColumn 与events(events 的子集)的平均值的比率,这些events 由相同的@ 执行987654344@.
我在映射events 时尝试这样做。在map 内,我试图根据他们的performerId 过滤事件。
基本上,我正在尝试将每个event 转换为LabeledPoint,上面的比率将成为我的向量中的一个特征。即对于每一个事件,我都试图得到
// I am trying to calculate the average, but cannot use filter, because I am in map block.
LabeledPoint(
event.someColumn,
Vectors.dense(
averageAbove,
...
)
)
我将不胜感激。谢谢!
【问题讨论】:
标签: java scala apache-spark filter rdd