【问题标题】:Connect pyspark with Neo4j将 pyspark 与 Neo4j 连接起来
【发布时间】:2023-03-29 04:28:01
【问题描述】:

我想使用 pyspark 向 neo4j 提供数据。我试过下面的代码。

from pyspark.sql import SparkSession

spark = SparkSession.builder.master("local[1]") \
    .appName("SparkByExamples.com") \
    .getOrCreate()

df = spark.read.csv("countries.csv")

df.write \
    .format("org.neo4j.spark.DataSource") \
    .mode("ErrorIfExists") \
    .option("url", "bolt://localhost:7687") \
    .option("labels", ":Countries") \
    .save()

但它给了我这样的错误:

2020-11-11 17:49:19 WARN  Utils:66 - Set SPARK_LOCAL_IP if you need to bind to another address
2020-11-11 17:49:20 WARN  NativeCodeLoader:62 - Unable to load native-hadoop library for your platform... using builtin-java classes where applicable
Setting default log level to "WARN".
To adjust logging level use sc.setLogLevel(newLevel). For SparkR, use setLogLevel(newLevel).
Traceback (most recent call last):
  File "xxx", line 13, in <module>
    .option("labels", ":Countries") \
  File "/opt/spark/spark-2.4.0-bin-hadoop2.7/python/pyspark/sql/readwriter.py", line 734, in save
    self._jwrite.save()
  File "/opt/spark/spark-2.4.0-bin-hadoop2.7/python/lib/py4j-0.10.7-src.zip/py4j/java_gateway.py", line 1257, in __call__
  File "/opt/spark/spark-2.4.0-bin-hadoop2.7/python/pyspark/sql/utils.py", line 63, in deco
    return f(*a, **kw)
  File "/opt/spark/spark-2.4.0-bin-hadoop2.7/python/lib/py4j-0.10.7-src.zip/py4j/protocol.py", line 328, in get_return_value
py4j.protocol.Py4JJavaError: An error occurred while calling o40.save.
: java.lang.ClassNotFoundException: Failed to find data source: org.neo4j.spark.DataSource. Please find packages at http://spark.apache.org/third-party-projects.html
    at org.apache.spark.sql.execution.datasources.DataSource$.lookupDataSource(DataSource.scala:657)

谁能帮我解决这个问题? 提前谢谢你。

【问题讨论】:

    标签: csv apache-spark pyspark neo4j write


    【解决方案1】:

    问题出在您的格式上。您需要将包包含到 Spark 应用程序中

    试试

    $SPARK_HOME/bin/spark-shell --packages neo4j-contrib:neo4j-spark-connector:2.4.5-M2
    

    【讨论】:

    • 我已经包含了这个包。但它仍然给我同样的错误。
    • 可以从网站 (spark-packages.org/package/neo4j-contrib/neo4j-spark-connector) 下载 jar 文件,然后使用 --jars /path/to/jar 运行 pyspark 或 spark-submit 命令
    • 我从给定的链接下载了 jar 文件并使用 spark-submit 命令传递了 jar 文件,但它仍然给我同样的错误。
    【解决方案2】:

    您必须添加 neo4j-connector jar 并在启动 spark shell 时传递它。

    pyspark --jars neo4j-connector-apache-spark_2.12-4.0.1_for_spark_3.jar
    

    你可以从https://neo4j.com/product/connectors/apache-spark-connector/下载这个jar

    【讨论】:

      猜你喜欢
      • 2022-01-19
      • 1970-01-01
      • 2019-08-15
      • 2022-08-18
      • 1970-01-01
      • 2022-07-04
      • 2017-09-24
      • 2014-08-16
      • 2020-01-20
      相关资源
      最近更新 更多