【发布时间】: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 并行处理查询。另外需要设置lowerBound和upperBound。据我了解,就我而言,我可以将id 和date 列用于partitionColumn。但是,这里的问题是如何在对其中一列进行分区时设置lowerBound 和upperBound 值。我注意到如果设置不当,我的情况会出现数据偏差。对于 Spark 中的处理,我不关心自然分区。我只需要尽可能快地转换所有数据,因此我认为优化非偏斜分区是首选。
我已经为此提出了一个解决方案,但我不确定这样做是否真的有意义。本质上,它是将 id 散列到分区中。我的解决方案是在id 列上使用mod() 并指定分区数。那么dbtable 中的字段将类似于:
"(SELECT *, mod(id, <<num-parallel-queries>>) as part FROM <<schema>>.<<table>>) as t"
然后我使用partitionColum="part"、lowerBound=0 和upperBound=<<num-parallel-queries>> 作为 Spark 读取 JDBC 作业的选项。
请让我知道这是否有意义!
【问题讨论】:
标签: python postgresql apache-spark jdbc