【问题标题】:neo4j with Flink and Scalaneo4j 与 Flink 和 Scala
【发布时间】:2017-08-31 16:24:44
【问题描述】:

我正在使用 Scala 2.11.7 和 Flink 1.3.2 处理数据。现在我想将生成的 org.apache.flink.api.scala.DataSet 存储在 neo4j 图形数据库中。

为了兼容有 Github 项目:

  • 使用 neo4j 进行 Flink:https://github.com/s1ck/flink-neo4j
  • 带有 neo4j 的 Scala:_https://github.com/FaKod/neo4j-scala
  • Flink 的图形库“Gelly”与 neo4j:_https://github.com/albertodelazzari/gelly-neo4j

最有希望的方法是什么?还是直接使用neo4j的REST API更好?

(顺便说一句:为什么stackoverflow会限制postet的链接数量...?)

我试过flink-neo4j,但是混合Java和Scala类似乎有一些问题:

package dummy.neo4j

import org.apache.flink.api.common.io.OutputFormat
import org.apache.flink.api.java.io.neo4j.Neo4jOutputFormat
import org.apache.flink.api.java.tuple.{Tuple, Tuple2}
import org.apache.flink.api.scala._

object Neo4jDummyWriter {

  def main(args: Array[String]) {
    val env = ExecutionEnvironment.getExecutionEnvironment

    val outputFormat: OutputFormat[_ <: Tuple] = Neo4jOutputFormat.buildNeo4jOutputFormat.setRestURI("http://localhost:7474/db/data/")
  .setConnectTimeout(1000).setReadTimeout(1000).setCypherQuery("UNWIND {inserts} AS i CREATE (a:User {name:i.name, born:i.born})")
  .addParameterKey(0, "name").addParameterKey(1, "born").setTaskBatchSize(1000).finish

    val tuple1: Tuple = new Tuple2("abc", 1)
    val tuple2: Tuple = new Tuple2("def", 2)

    val test = env.fromElements[Tuple](tuple1, tuple2)
    println("test: " + test.getClass)
    test.output(outputFormat)
  }

}

线程“main”中的异常 java.lang.ClassCastException: [Ljava.lang.Object;无法转换为 [Lorg.apache.flink.api.common.typeinfo.TypeInformation; 在 dummy.neo4j.Neo4jDummyWriter$.main(Neo4jDummyWriter.scala:20) 在 dummy.neo4j.Neo4jDummyWriter.main(Neo4jDummyWriter.scala)

类型不匹配,预期:OutputFormat[Tuple],实际:OutputFormat[_ <: tuple>

【问题讨论】:

    标签: scala neo4j apache-flink


    【解决方案1】:

    解决办法是不要把 Tuple2 对象改成 Tuple:

    package dummy.neo4j
    
    import org.apache.flink.api.common.io._
    import org.apache.flink.api.java.io.neo4j.Neo4jOutputFormat
    import org.apache.flink.api.java.tuple.Tuple2
    import org.apache.flink.api.scala._
    
    object Neo4jDummyWriter {
    
      def main(args: Array[String]) {
        val env = ExecutionEnvironment.getExecutionEnvironment
    
        val tuple1 = ("user9", 1978)
        val tuple2 = ("user10", 1996)
        val datasetWithScalaTuples = env.fromElements(tuple1, tuple2)
        val dataset: DataSet[Tuple2[String, Int]] = datasetWithScalaTuples.map(tuple => new Tuple2(tuple._1, tuple._2))
    
        val outputFormat = Neo4jOutputFormat.buildNeo4jOutputFormat.setRestURI("http://localhost:7474/db/data/").setUsername("neo4j").setPassword("...")
      .setConnectTimeout(1000).setReadTimeout(1000).setCypherQuery("UNWIND {inserts} AS i CREATE (a:User {name:i.name, born:i.born})")
      .addParameterKey(0, "name").addParameterKey(1, "born").setTaskBatchSize(1000).finish.asInstanceOf[OutputFormat[Tuple2[String, Int]]]
    
        dataset.output(outputFormat)
        env.execute
      }
    
    }
    

    【讨论】:

      猜你喜欢
      • 2016-05-24
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-03-26
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多