【问题标题】:Using Neo4j with Apache Spark将 Neo4j 与 Apache Spark 结合使用
【发布时间】:2015-03-06 10:31:12
【问题描述】:

我正在尝试将 Neo4j 与 Apache Spark Streaming 一起使用,但我发现可串行化是一个问题。

基本上,我希望 Apache Spark 实时解析和捆绑我的数据。之后,数据已被捆绑,它应该存储在数据库 Neo4j 中。但是,我收到此错误:

org.apache.spark.SparkException: Task not serializable
    at org.apache.spark.util.ClosureCleaner$.ensureSerializable(ClosureCleaner.scala:166)
    at org.apache.spark.util.ClosureCleaner$.clean(ClosureCleaner.scala:158)
    at org.apache.spark.SparkContext.clean(SparkContext.scala:1264)
    at org.apache.spark.api.java.JavaRDDLike$class.foreach(JavaRDDLike.scala:297)
    at org.apache.spark.api.java.JavaPairRDD.foreach(JavaPairRDD.scala:45)
    at twoGrams.Main$4.call(Main.java:102)
    at twoGrams.Main$4.call(Main.java:1)
    at org.apache.spark.streaming.api.java.JavaDStreamLike$$anonfun$foreachRDD$2.apply(JavaDStreamLike.scala:282)
    at org.apache.spark.streaming.api.java.JavaDStreamLike$$anonfun$foreachRDD$2.apply(JavaDStreamLike.scala:282)
    at org.apache.spark.streaming.dstream.ForEachDStream$$anonfun$1.apply$mcV$sp(ForEachDStream.scala:41)
    at org.apache.spark.streaming.dstream.ForEachDStream$$anonfun$1.apply(ForEachDStream.scala:40)
    at org.apache.spark.streaming.dstream.ForEachDStream$$anonfun$1.apply(ForEachDStream.scala:40)
    at scala.util.Try$.apply(Try.scala:161)
    at org.apache.spark.streaming.scheduler.Job.run(Job.scala:32)
    at org.apache.spark.streaming.scheduler.JobScheduler$JobHandler.run(JobScheduler.scala:172)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1145)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:615)
    at java.lang.Thread.run(Thread.java:745)
Caused by: java.io.NotSerializableException: org.neo4j.kernel.EmbeddedGraphDatabase
    at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1183)
    at java.io.ObjectOutputStream.defaultWriteFields(ObjectOutputStream.java:1547)
    at java.io.ObjectOutputStream.writeSerialData(ObjectOutputStream.java:1508)
    at java.io.ObjectOutputStream.writeOrdinaryObject(ObjectOutputStream.java:1431)
    at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1177)
    at java.io.ObjectOutputStream.defaultWriteFields(ObjectOutputStream.java:1547)
    at java.io.ObjectOutputStream.writeSerialData(ObjectOutputStream.java:1508)
    at java.io.ObjectOutputStream.writeOrdinaryObject(ObjectOutputStream.java:1431)
    at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1177)
    at java.io.ObjectOutputStream.defaultWriteFields(ObjectOutputStream.java:1547)
    at java.io.ObjectOutputStream.writeSerialData(ObjectOutputStream.java:1508)
    at java.io.ObjectOutputStream.writeOrdinaryObject(ObjectOutputStream.java:1431)
    at java.io.ObjectOutputStream.writeObject0(ObjectOutputStream.java:1177)
    at java.io.ObjectOutputStream.writeObject(ObjectOutputStream.java:347)
    at org.apache.spark.serializer.JavaSerializationStream.writeObject(JavaSerializer.scala:42)
    at org.apache.spark.serializer.JavaSerializerInstance.serialize(JavaSerializer.scala:73)
    at org.apache.spark.util.ClosureCleaner$.ensureSerializable(ClosureCleaner.scala:164)
    ... 17 more

这是我的代码:

output a stream of type: JavaPairDStream<String, ArrayList<String>>

output.foreachRDD(
                new Function2<JavaPairRDD<String,ArrayList<String>>,Time,Void>(){

                    @Override
                    public Void call(
                            JavaPairRDD<String, ArrayList<String>> arg0,
                            Time arg1) throws Exception {
                        // TODO Auto-generated method stub

                        arg0.foreach(
                                new VoidFunction<Tuple2<String,ArrayList<String>>>(){

                                    @Override
                                    public void call(
                                            Tuple2<String, ArrayList<String>> arg0)
                                            throws Exception {
                                        // TODO Auto-generated method stub
                                        try( Transaction tx = graphDB.beginTx()){
                                            if(Neo4jOperations.getHMacFromValue(graphDB, arg0._1)!=null)
                                                System.out.println("Alread in Database:" + arg0._1);
                                            else{
                                                Neo4jOperations.createHMac(graphDB, arg0._1);
                                            }
                                            tx.success();
                                        }
                                    }

                        });
                        return null;
                    }



                });

Neo4jOperations 类:

public class Neo4jOperations{

public static Node getHMacFromValue(GraphDatabaseService graphDB,String value){
        try(ResourceIterator<Node> HMacs=graphDB.findNodesByLabelAndProperty(DynamicLabel.label("HMac"), "value", value).iterator()){
            return HMacs.next();
        }
    }

    public static void createHMac(GraphDatabaseService graphDB,String value){
        Node HMac=graphDB.createNode(DynamicLabel.label("HMac"));
        HMac.setProperty("value", value);
        HMac.setProperty("time", new SimpleDateFormat("yyyyMMdd_HHmmss").format(Calendar.getInstance().getTime()));
    }
}

我知道我必须对 Neo4jOperations 类进行序列化,但我能够弄清楚如何。或者还有其他方法可以实现吗?

【问题讨论】:

    标签: java serialization apache-spark neo4j


    【解决方案1】:

    在涉及到外部系统的连接或处理不可序列化对象时,可以直接在工作线程上创建这些对象并避免序列化的需要。

    Given: val stream: DStream = ???
    stream.forEachRDD{rdd =>
       rdd.forEachPartition{iter =>
           val nonSerializableConn = new NonSerializableDriver(ip, port)
           iter.foreach(elem => nonSerializableConn.doStuff(elem)
       }
    }
    

    这种模式通过在每个分区(将包含许多元素)只执行一次来分摊对象创建

    在诸如 Spark Streaming 之类的长期进程中,我们可以通过保持每个 VM 的资源缓存来进一步减少开销:

    stream.forEachRDD{rdd =>
       rdd.forEachPartition{iter =>
           val nonSerializableConn = NonSerializableDriver.getConnection(ip, port)
           iter.foreach(elem => nonSerializableConn.doStuff(elem)
       }
    }
    

    在后一种情况下,我们需要在VM终止时进行连接管理和关闭资源。

    【讨论】:

    • 有趣的模式。我认为这种方法的问题在于我所知道的 Neo4j 没有 Java 不可序列化驱动程序。 Neo4j 的嵌入式图形数据库可能会被捆绑到一个胖 JAR 中并发送到 Spark。我认为使用 TinkerPop 来做类似的事情更有意义。最后,Neo4j 服务器公开了一个 REST API,因此驱动程序是包装 HTTP 请求的客户端库。您必须使请求持久并使用批处理 API。不理想。
    【解决方案2】:

    没有可能的方法来序列化 Neo4jOperations 类中包含的传递依赖项。不幸的是,Spark 不能那样工作。

    问题是 Neo4j 遍历 API 无法序列化或捆绑并分派到 Spark。即使您尝试将 Spark 捆绑到 Neo4j 中,您也会遇到与 Jetty servlet 版本的依赖冲突。

    这就是我创建Neo4j Mazerunner 的原因。在创建扩展 Spark RDD 包的基类的 Neo4j Spark 连接器之前,没有一种简单的方法可以将来自 Neo4j 的数据导入 Spark 的运行时。

    请参阅 Couchbase's Spark Connector 以了解这样做所涉及的内容。

    Mazerunner 尚不支持流媒体功能,但我计划在未来实现这一点

    【讨论】:

    • 所以,基本上你的建议是我无法连接 Spark Streaming 和 Neo4j ?
    猜你喜欢
    • 2015-08-22
    • 2019-03-03
    • 1970-01-01
    • 2015-12-09
    • 1970-01-01
    • 1970-01-01
    • 2018-01-24
    • 1970-01-01
    • 2014-01-01
    相关资源
    最近更新 更多