【问题标题】:Problems with reading Postgres database using JDBC connection in Spark在 Spark 中使用 JDBC 连接读取 Postgres 数据库的问题
【发布时间】:2021-04-07 17:19:45
【问题描述】:

我目前在使用 (Py)Spark 中的 JDBC 连接从 Postgres 数据库读取数据时遇到了一些问题。我在 Postgres 中有一个表,我想在 Spark 中读取它,对其进行处理,然后将结果作为 .parquet 文件保存在 AWS S3 存储桶中。

我创建了一个执行一些基本逻辑的示例脚本(以免问题过于复杂):

from pyspark.sql import SparkSession
from pyspark.sql.functions import length
import argparse
import uuid
import datetime

def parse_arguments():
    parser = argparse.ArgumentParser()
    parser.add_argument(
        "--loc",
        type=str,
        default="./",
        help="Output location"
    )

    args = parser.parse_known_args()[0]
    return args

if __name__=="__main__":
    args = parse_arguments()

    spark = SparkSession.builder. \
        appName("test-script"). \
        config("spark.jars.packages", "com.johnsnowlabs.nlp:spark-nlp_2.11:2.6.5,org.postgresql:postgresql:42.1.1"). \
        getOrCreate()

    CENTRAL_ID = "id"
    PG_HOST = spark.conf.get("spark.yarn.appMasterEnv.PG_HOST")
    PG_PORT = spark.conf.get("spark.yarn.appMasterEnv.PG_PORT")
    PG_USER = spark.conf.get("spark.yarn.appMasterEnv.PG_USER")
    PG_DB = spark.conf.get("spark.yarn.appMasterEnv.PG_DB")
    PG_PASS = spark.conf.get("spark.yarn.appMasterEnv.PG_PASS")
    PG_MAX_CONCURRENT = spark.conf.get("spark.yarn.appMasterEnv.PG_MAX_CONCURRENT")

    table = "test_schema.test_table"
    partitions = spark.sparkContext.defaultParallelism
    fetch_size = 2000

    data = spark.read.format("jdbc"). \
        option("url", "jdbc:postgresql://{}:{}/{}".format(PG_HOST, PG_PORT, PG_DB)). \
        option("dbtable", "(SELECT *, MOD({}, {}) AS p FROM {}) AS t".format(CENTRAL_ID, partitions, table)). \
        option("user", PG_USER). \
        option("password", PG_PASS). \
        option("driver", "org.postgresql.Driver"). \
        option("partitionColumn", "p"). \
        option("lowerBound", 0). \
        option("upperBound", partitions). \
        option("numPartitions", PG_MAX_CONCURRENT). \
        option("fetchSize", fetch_size). \
        load()
    data = data.repartition(partitions)

    # Cache data
    data.cache()

    # Calculate on data
    out1 = data.withColumn("abstract_length", length("abstract"))
    out2 = data.withColumn("title_length", length("title"))

    # Create a timestamp
    time_stamp = datetime.datetime.utcnow().isoformat()
    save_id = "{}-{}".format(uuid.uuid4(), time_stamp)

    out1.select([CENTRAL_ID, "abstract_length"]).write.mode("overwrite").parquet("{}/{}/abstract".format(args.loc, save_id))
    out2.select([CENTRAL_ID, "title_length"]).write.mode("overwrite").parquet("{}/{}/title".format(args.loc, save_id))

    spark.stop()

    print("run successfull!")

程序尝试将 Potsgres 数据库的并发查询数限制为 8(请参阅 PG_MAX_CONCURRENT)。这是为了不使数据库过载。加载后,我重新分区到更多 (360) 个分区,以便将数据分发给所有工作人员。

EMR集群配置如下:

  • 主:m5.xlarge 4 vCore,16 GiB 内存,仅 EBS 存储 EBS 存储:64 GiB
  • 6 x 从站:m5.4xlarge 16 vCore,64 GiB 内存,仅 EBS 存储 EBS 存储:64 GiB
  • Spark 2.4.7

spark-submit 参数如下:

spark-submit \
--master yarn \
--conf 'spark.yarn.appMasterEnv.PG_HOST=<<host>>' \
--conf 'spark.yarn.appMasterEnv.PG_PORT=<<port>>' \
--conf 'spark.yarn.appMasterEnv.PG_DB=<<db>>' \
--conf 'spark.yarn.appMasterEnv.PG_USER=<<user>>' \
--conf 'spark.yarn.appMasterEnv.PG_PASS=<<password>>' \
--conf 'spark.yarn.appMasterEnv.PG_MAX_CONCURRENT=8' \
--conf 'spark.executor.cores=3' \
--conf 'spark.executor.instances=30' \
--conf 'spark.executor.memory=12g' \
--conf 'spark.driver.memory=12g' \
--conf 'spark.default.parallelism=360' \
--conf 'spark.kryoserializer.buffer.max=1000M' \
--conf 'spark.serializer=org.apache.spark.serializer.KryoSerializer' \
--conf 'spark.dynamicAllocation.enabled=false' \
--packages 'com.johnsnowlabs.nlp:spark-nlp_2.11:2.6.5,org.postgresql:postgresql:42.1.1' \
program.py \
--loc s3a://<<bucket>>/

我遇到的第一种错误:

Job aborted due to stage failure: Task 0 in stage 0.0 failed 4 times, most recent failure: Lost task 0.4 in stage 0.0 (TID 26, ip-172-31-35-159.eu-central-1.compute.internal, executor 9): ExecutorLostFailure (executor 9 exited caused by one of the running tasks) Reason: Executor heartbeat timed out after 140265 ms
Driver stacktrace:

我不确定这意味着什么。可能是从表中获取数据需要太长时间吗?还是有其他原因?

我得到的第二种错误是:

    ExecutorLostFailure (executor 7 exited caused by one of the running tasks) Reason: Container from a bad node: container_1609406414316_0002_01_000013 on host: ip-172-31-44-127.eu-central-1.compute.internal. Exit status: 137. Diagnostics: [2020-12-31 10:08:38.093]Container killed on request. Exit code is 137

这似乎表明存在 OOM 问题。但我不明白为什么会出现 OOM 错误,因为在我看来,我为执行程序和驱动程序分配了足够的内存。此外,当我查看集群的统计信息时,我知道它有足够的内存:

会不会是这样的情况,当对 Postgres 服务器使用 8 个并发查询时,它会将 1/8 的数据发送给每个执行器,所以执行器应该准备好接收总大小的 1/8?或者fetchSize 是否限制发送给执行程序的数据大小以避免内存问题?或者也许还有其他原因?我尝试处理的整个表约为 110 GB。

有人可以帮忙吗?提前致谢!

【问题讨论】:

    标签: postgresql apache-spark jdbc pyspark out-of-memory


    【解决方案1】:

    您似乎正在超时。您是否尝试过在 spark 提交参数中增加超时配置?

    --conf 'spark.network.timeout=10000000' \
    

    Spark cluster full of heartbeat timeouts, executors exiting on their own

    【讨论】:

    • 谢谢,我会重试的。但是,这并不能解决 OOM 错误,对吧?你知道是什么原因造成的吗?
    • 它似乎工作得更好,但仍然出现ExecutorLostFailure: Container killed on request. Exit code is 137 错误......现在,超时似乎已经消失了。
    【解决方案2】:

    由于 executor 容器内存不足,请尝试将其添加到您的 spark 提交配置中:

    --conf 'spark.executor.memory=20g' \
    

    如果失败,请尝试更多。

    【讨论】:

    • 那么我必须减少执行者的数量,因为集群的内存限制。但是,EMR 内存统计数据(右上图)没有显示任何过载迹象,这不是很奇怪吗?
    • 哦,我现在明白了。总体而言,您有足够的内存,但似乎其中一个执行程序内存不足,是吗?
    • 可能需要将 fetchsize 降低到 500 或更少
    • 是的,没错,一个或多个由于某种原因内存不足。我不确定为什么会这样。要读入的数据总量(110GB 对 360GB)比我分配的内存要小得多……似乎由于某种原因,更少的执行程序(18 个而不是 30 个)获得了更多的数据。有什么想法吗?
    猜你喜欢
    • 2017-01-12
    • 2015-06-14
    • 1970-01-01
    • 2011-05-28
    • 2020-12-13
    • 1970-01-01
    • 2022-01-23
    • 2017-09-17
    • 2021-09-12
    相关资源
    最近更新 更多