【发布时间】:2018-06-18 13:36:00
【问题描述】:
我有一个要求,我必须从 Neo4j 中提取数据并从该数据中创建 Spark RDD。我在我的项目中使用 Python。 this 连接器用于相同目的,但它是用 Scala 编写的。所以我现在可以考虑以下解决方法 -
从neo4j中以小块/批量查询数据,使用
parallize()方法将每个块转换为Spark RDD。最后使用union()方法合并/组合所有RDD以获得单个RDD。然后我可以对它们进行转换和操作。另一种方法是从 Neo4j 读取数据并从中创建一个 Kafka 生产者。然后使用 Kafka 作为 Spark 的数据源。例如
Neo4j -> 卡夫卡 -> Spark
我想知道哪个对大块数据更有效?如果有更好的方法来解决这个问题,请帮助我。
注意:我确实尝试扩展 pyspark API 以便在 python 中创建自定义 RDD。与 Spark 的 Scala/Java API 相比,pyspark 的 API 非常不同。对于 Scala API,可以通过扩展 RDD 类并覆盖 compute() 和 getPartitions() 方法来创建自定义 RDD。但是在pyspark API中,我在rdd.py的RDD类下找不到compute()
【问题讨论】:
标签: python apache-spark neo4j pyspark apache-kafka