【问题标题】:Spark SQL RDD loads in pyspark but not in spark-submit: "JDBCRDD: closed connection"Spark SQL RDD 在 pyspark 中加载,但不在 spark-submit 中:“JDBCRDD:关闭连接”
【发布时间】:2017-02-16 08:33:48
【问题描述】:

我有以下简单的代码,用于将我的 Postgres 数据库中的表加载到 RDD 中。

# this setup is just for spark-submit, will be ignored in pyspark
from pyspark import SparkConf, SparkContext
from pyspark.sql import SQLContext
conf = SparkConf().setAppName("GA")#.setMaster("localhost")
sc = SparkContext(conf=conf)
sqlContext = SQLContext(sc)

# func for loading table
def get_db_rdd(table):
    url = "jdbc:postgresql://localhost:5432/harvest?user=postgres"
    print(url)
    lower = 0
    upper = 1000
    ret = sqlContext \
      .read \
      .format("jdbc") \
      .option("url", url) \
      .option("dbtable", table) \
      .option("partitionColumn", "id") \
      .option("numPartitions", 1024) \
      .option("lowerBound", lower) \
      .option("upperBound", upper) \
      .option("password", "password") \
      .load()
    ret = ret.rdd
    return ret

# load table, and print results
print(get_db_rdd("mytable").collect())

我运行./bin/pyspark,然后将其粘贴到解释器中,它会按预期打印出我的表中的数据。

现在,如果我将该代码保存到名为 test.py 的文件中,然后执行 ./bin/spark-submit test.py,它会开始运行,但随后我会看到这些消息永远在我的控制台中发送:

17/02/16 02:24:21 INFO Executor: Running task 45.0 in stage 0.0 (TID 45)
17/02/16 02:24:21 INFO JDBCRDD: closed connection
17/02/16 02:24:21 INFO Executor: Finished task 45.0 in stage 0.0 (TID 45). 1673 bytes result sent to driver

编辑:这是在一台机器上。我还没有开始任何主人或奴隶; spark-submit 是我在系统启动后运行的唯一命令。我尝试使用具有相同结果的主/从设置。 我的spark-env.sh 文件如下所示:

export SPARK_WORKER_INSTANCES=2
export SPARK_WORKER_CORES=2
export SPARK_WORKER_MEMORY=800m
export SPARK_EXECUTOR_MEMORY=800m
export SPARK_EXECUTOR_CORES=2
export SPARK_CLASSPATH=/home/ubuntu/spark/pg_driver.jar # Postgres driver I need for SQLContext
export PYTHONHASHSEED=1337 # have to make workers use same seed in Python3

如果我 spark-submit 一个仅从列表或其他内容创建 RDD 的 Python 文件,它会起作用。我只有在尝试使用 JDBC RDD 时才会遇到问题。我错过了什么?

【问题讨论】:

    标签: apache-spark jdbc pyspark


    【解决方案1】:

    当使用spark-submit 时,您应该将jar 提供给执行者。

    spark 2.1 JDBC documents中所述:

    要开始使用,您需要包含 JDBC 驱动程序 spark类路径上的特定数据库。例如,要连接到 来自 Spark Shell 的 postgres,您将运行以下命令:

    bin/spark-shell --driver-class-path postgresql-9.4.1207.jar --jars postgresql-9.4.1207.jar
    

    注意:spark-submit 命令也应如此

    疑难解答

    JDBC 驱动程序类必须对原始类加载器可见 在客户端会话和所有执行者上。这是因为 Java 的 DriverManager 类进行安全检查,导致它忽略 当一个驱动程序运行时,所有驱动程序对原始类加载器都不可见 打开连接。 一种方便的方法是修改 所有工作节点上的compute_classpath.sh 以包含您的驱动程序JAR。

    【讨论】:

    • 我的$SPARK_CLASSPATH 设置应该已经这样做了,但我还是尝试了你的建议。取消设置该环境变量并运行spark-submit --driver-class-path pg_driver.jar --jars pg_driver.jar test.py,它也有同样的问题。我认为如果它缺少驱动程序,它会抛出一些其他错误,例如“找不到合适的驱动程序”。顺便说一句,这是在一台机器上(更新我的问题)。
    【解决方案2】:

    这是一个可怕的黑客攻击。我不认为这是答案,但它确实有效。

    好的,只有 pyspark 有效吗?好,那我们就用它。写了这个 Bash 脚本:

    cat $1 | $SPARK_HOME/bin/pyspark # pipe the Python file into pyspark
    

    我在提交作业的 Python 脚本中运行该脚本。此外,我还包括了用于在进程之间传递参数的代码,以防它对某人有所帮助:

    new_env = os.environ.copy()
    new_env["pyspark_argument_1"] = "some param I need in my Spark script" # etc...
    p = subprocess.Popen(["pyspark_wrapper.sh {}".format(py_fname)], shell=True, env=new_env)
    

    在我的 Spark 脚本中:

    something_passed_from_submitter = os.environ["pyspark_argument_1"]
    # do stuff in Spark...
    

    我觉得 Spark 得到了更好的支持,并且(如果这是一个错误)Scala 比 Python 3 的错误更少,因此这可能是目前更好的解决方案。但是我的脚本使用了我们在 Python 3 中编写的一些文件,所以...

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2015-06-17
      • 1970-01-01
      • 1970-01-01
      • 2017-01-24
      • 2020-10-07
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多