【问题标题】:Parallel reading and processing from Postgres using (Py)Spark使用 (Pyspark) 从 Postgres 并行读取和处理
【发布时间】:2021-03-28 04:26:49
【问题描述】:

我有一个关于从 Postgres 数据库中读取大量数据并使用 spark 并行处理它的问题。假设我在 Postgres 中有一个表,我想使用 JDBC 读入 Spark。假设它有以下列:

  • id (bigint)
  • 日期(日期时间)
  • 许多其他列(不同类型)

目前 Postgres 表没有分区。我想并行转换大量数据,并最终将转换后的数据存储在其他地方。

问题:我们如何优化从 Postgres 中并行读取数据?

文档 (https://spark.apache.org/docs/latest/sql-data-sources-jdbc.html) 建议使用 partitionColum 并行处理查询。另外需要设置lowerBoundupperBound。据我了解,就我而言,我可以将iddate 列用于partitionColumn。但是,这里的问题是如何在对其中一列进行分区时设置lowerBoundupperBound 值。我注意到如果设置不当,我的情况会出现数据偏差。对于 Spark 中的处理,我不关心自然分区。我只需要尽可能快地转换所有数据,因此我认为优化非偏斜分区是首选。

我已经为此提出了一个解决方案,但我不确定这样做是否真的有意义。本质上,它是将 id 散列到分区中。我的解决方案是在id 列上使用mod() 并指定分区数。那么dbtable 中的字段将类似于:

"(SELECT *, mod(id, <<num-parallel-queries>>) as part FROM <<schema>>.<<table>>) as t"

然后我使用partitionColum="part"lowerBound=0upperBound=&lt;&lt;num-parallel-queries&gt;&gt; 作为 Spark 读取 JDBC 作业的选项。

请让我知道这是否有意义!

【问题讨论】:

    标签: python postgresql apache-spark jdbc


    【解决方案1】:

    按主键列“分区”是个好主意。

    要获得大小相等的分区,请使用表统计信息:

    SELECT histogram_bounds::text::bigint[]
    FROM pg_stats
    WHERE tablename = 'mytable'
      AND attname = 'id';
    

    如果 default_statistics_target 的默认值为 100,这将是一个包含 101 个值的数组,用于分隔 0 到 100 的百分位数。您可以使用它对表进行均匀分区。

    例如:如果数组看起来像 {42,10001,23066,35723,49756,...,999960} 并且您需要 50 个分区,则第一个将是 id id

    【讨论】:

      猜你喜欢
      • 2015-03-21
      • 1970-01-01
      • 2020-10-28
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多