【问题标题】:SparkSQL PostgresQL Dataframe partitionsSparkSQL PostgresQL 数据框分区
【发布时间】:2015-09-29 06:18:57
【问题描述】:

我有一个非常简单的 SparkSQL 设置连接到 Postgres 数据库,我试图从表中获取一个 DataFrame,该 DataFrame 具有多个 X 分区(假设为 2)。代码如下:

Map<String, String> options = new HashMap<String, String>();
options.put("url", DB_URL);
options.put("driver", POSTGRES_DRIVER);
options.put("dbtable", "select ID, OTHER from TABLE limit 1000");
options.put("partitionColumn", "ID");
options.put("lowerBound", "100");
options.put("upperBound", "500");
options.put("numPartitions","2");
DataFrame housingDataFrame = sqlContext.read().format("jdbc").options(options).load();

由于某种原因,DataFrame 的一个分区几乎包含所有行。

据我所知lowerBound/upperBound 是用于微调的参数。在 SparkSQL 的文档(Spark 1.4.0 - spark-sql_2.11)中,它说它们用于定义步幅,而不是过滤/范围分区列。但这提出了几个问题:

  1. 步幅是 Spark 为每个执行程序(分区)查询数据库的频率(每个查询返回的元素数)?
  2. 如果不是,此参数的用途是什么,它们依赖于什么以及如何以稳定的方式平衡我的 DataFrame 分区(不要求所有分区包含相同数量的元素,只是存在一个平衡 -例如 2 个分区 100 个元素 55/45 、 60/40 甚至 65/35 都可以)

似乎无法找到这些问题的明确答案,想知道你们中的一些人是否可以为我澄清这一点,因为现在在处理 X 百万行时影响我的集群性能,所有繁重的工作都过去了给一个执行者。

干杯并感谢您的宝贵时间。

【问题讨论】:

    标签: postgresql apache-spark apache-spark-sql partition


    【解决方案1】:

    本质上,上下界和分区数用于计算每个并行任务的增量或拆分。

    假设表有分区列“year”,并且有 2006 年到 2016 年的数据。

    如果您将分区数定义为 10,下限为 2006 年,上限为 2016 年,您将让每个任务获取自己年份的数据 - 理想情况。

    即使您错误地指定了下限和/或上限,例如设置 lower = 0 和 upper = 2016,数据传输会有偏差,但是,您不会“丢失”或无法检索任何数据,因为:

    第一个任务将获取

    第二个任务将获取 0 到 2016/10 年之间的数据。

    第三个任务将获取 2016/10 和 2*2016/10 之间年份的数据。

    ...

    最后一个任务的 where 条件是 year->2016。

    T.

    【讨论】:

    • "最后一个任务的 where 条件为 year->2016。"你的意思是year &gt; 2016(比2016年大)还是year -&gt; 2016(直到2016年)。我认为您的意思是前者,但想澄清一下。
    【解决方案2】:

    下限确实用于分区列;参考这段代码(撰写本文时的当前版本):

    https://github.com/apache/spark/blob/40ed2af587cedadc6e5249031857a922b3b234ca/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/jdbc/JDBCRelation.scala

    函数columnPartition包含分区逻辑和使用下限/上限的代码。

    【讨论】:

      【解决方案3】:

      目前已确定下界和上界可以执行它们在之前的答案中所做的事情。对此的跟进将是如何在不查看最小最大值或您的数据严重倾斜的情况下跨分区平衡数据。

      如果您的数据库支持“散列”功能,它就可以解决问题。

      partitionColumn = "hash(column_name)%num_partitions"

      numPartitions = 10 // 随便你

      下界 = 0

      upperBound = numPartitions

      只要模运算返回 [0,numPartitions) 上的均匀分布,这将起作用

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2015-10-22
        • 2018-12-10
        • 2018-04-20
        • 2019-05-19
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多