【问题标题】:How to share data among JavaRDD partitions in Spark?如何在 Spark 中的 JavaRDD 分区之间共享数据?
【发布时间】: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


    【解决方案1】:

    为了在机器之间共享数据,数据必须在源端序列化,通过网络传输,并在目标端反序列化。所以你不能传输不可序列化的对象。

    如果MeasureReportingData 不可序列化,则必须将其转换为可序列化对象,共享该对象,然后在函数内将其转换回MeasureReportingData

    【讨论】:

    • 好的,看起来序列化是一种方式。我知道如果 MeasureReportingData 是可序列化的,它会起作用。但是我面临的问题在 Nashorn 的 Serialize a Predicate 的链接中提到,对此有什么想法吗?
    • 您根本无法共享不可序列化的对象。没有办法解决这个问题。
    • @Dikei 好吧,你总是可以不分享,而是原地初始化。
    • 确实如此。我建议他在我的回答中这样做。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2015-04-15
    • 2015-07-14
    • 2018-10-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多