【发布时间】:2016-10-18 00:18:22
【问题描述】:
在火花流中,我想在处理每个批次之前查询数据库,将结果存储在可以序列化并通过网络发送到执行程序的哈希图中。
class ExecutingClass implements Serializable {
init(DB db) {
try(JavaStreamingContext jsc = new JavaStreamingContext(...)) {
JavaPairInputDStream<String,String> kafkaStream = getKafkaStream(jsc);
kafkaStream.foreachRDD(rdd -> {
// this part is supposed to execute in the driver
Map<String, String> indexMap = db.getIndexMap();// connects to a db, queries the results as a map
JavaRDD<String> results = processRDD(rdd, indexMap);
...
}
}
JavaRDD<String> processRDD(JavaPairRDD<String, String> rdd, Map<String,String> indexMap) {
...
}
}
在上面的代码中,indexMap 应该在驱动程序中初始化,生成的映射用于处理 rdd。当我在 foreachRDD 闭包之外声明 indexMap 时我没有问题,但是当我在里面执行它时出现序列化错误。这是什么原因?
我想做这样的事情的原因是为了确保我从数据库中获得每个批次的最新值。我怀疑这是由于 foreachRDD 的关闭试图序列化关闭之外的所有内容。
【问题讨论】:
-
为什么不能为此目的使用累加器(读写)/广播(只读)?在这种情况下,因为它是读写累加器,所以有意义不是吗?
-
闭包内的代码将被序列化并发送给执行器。所以我假设
db.getIndexMap()不能为此目的进行序列化。 -
@LiMuBei 这就是问题所在。对于每一批数据,我们先查询数据库得到indexMap,然后只传递indexMap进行处理。
标签: java serialization apache-spark spark-streaming