【发布时间】:2016-07-13 11:27:53
【问题描述】:
我有一些对象要在 apache spark 的分区之间共享。下面是代码 sn-p 和我面临的问题。
private static void processDataWithResult() throws IOException {
JavaRDD<Long> idRDD = createIdRDDUsingDb();
final MeasureReportingData measureReporingData = getMeasureReportingData(jobConfiguration);
resultRDD = idRDD.mapPartitions(new FlatMapFunction<Iterator<Long>, Boolean>() {
@Override
public Iterable<Boolean> call(Iterator<Long> idIterator) throws Exception {
MeasureReportingData mrd = measureReporingData;
final List<Boolean> dummyList = new ArrayList<>();
long minId = idIterator.next();
engine.processInBatch(minId, minId + BATCH_SIZE - 1);
return (Iterable<Boolean>) dummyList;
}
});
resultRDD.count();
}
我想将measureReportingData 对象分发到所有分区?
我收到序列化错误,因为MeasureReportingData 包含不是Serializable 的实例成员。这个问题中指定了问题的模拟:How to serialize a Predicate<T> from Nashorn engine in java 8
还有其他方法可以在分区之间共享 measureReportingData 吗?
【问题讨论】:
标签: java serialization apache-spark