partitionColumn,lowerBound,upperBound,numPartitions参数是什么意思?

Bha*_*rla 24 jdbc apache-spark apache-spark-sql

虽然通过在星火JDBC连接获取来自SQL Server的数据,我发现我可以设置一些并行的参数,如partitionColumn,lowerBound,upperBound,和numPartitions.我已经通过spark文档,但无法理解它.

谁能解释一下这些参数的含义?

小智 25

很简单:

  • partitionColumn 是一个应该用于确定分区的列.
  • lowerBoundupperBound确定要获取的值的范围.完整数据集将使用与以下查询对应的行:

    SELECT * FROM table WHERE partitionColumn BETWEEN lowerBound AND upperBound
    
    Run Code Online (Sandbox Code Playgroud)
  • numPartitions确定要创建的分区数.范围lowerBound和之间的范围upperBound分为numPartitions步幅等于:

    upperBound / numPartitions - lowerBound / numPartitions
    
    Run Code Online (Sandbox Code Playgroud)

    例如,如果:

    • lowerBound:0
    • upperBound:1000
    • numPartitions:10

    Stride等于100,分区对应于以下查询:

    • SELECT * FROM table WHERE partitionColumn BETWEEN 0 AND 100
    • SELECT * FROM table WHERE partitionColumn BETWEEN 100 AND 200
    • ...
    • SELECT * FROM table WHERE partitionColumn BETWEEN 900 AND 1000

  • 在 [Spark 文档](https://spark.apache.org/docs/latest/sql-data-sources-jdbc.html) 中,它说: **注意,lowerBound 和 upperBound 仅用于决定分区步长,而不是用于过滤表中的行。因此表中的所有行都将被分区并返回。此选项仅适用于读取。**这意味着将获取*整个*表,而不仅仅是 lowerBound 和 upperBound 之间的部分。 (13认同)
  • 答案不准确,因为在某些数据库中“BETWEEN”既包含下限又包含上限。实际实现分别使用 `>=` 和 `<`: [Spark doc](https://github.com/apache/spark/blob/17edfec59de8d8680f7450b4d07c147c086c105a/sql/core/src/main/scala/org/apache/spark /sql/execution/datasources/jdbc/JDBCRelation.scala#L85-L97) (5认同)
  • 这个答案是完全错误的,意味着上限和下限值过滤了正在读取的数据集。 (4认同)
  • 见安德里亚的回答。第一个和最后一个 SELECT 在他的答案中是正确的,但在这个答案中不正确 (3认同)

And*_*rea 13

实际上上面的列表遗漏了一些东西,特别是第一个和最后一个查询.

如果没有它们,您将丢失一些数据(之前lowerBound和之后的数据upperBound).从示例中不清楚,因为下限是0.

完整列表应该是:

SELECT * FROM table WHERE partitionColumn < 100

SELECT * FROM table WHERE partitionColumn BETWEEN 0 AND 100  
SELECT * FROM table WHERE partitionColumn BETWEEN 100 AND 200  
Run Code Online (Sandbox Code Playgroud)

...

SELECT * FROM table WHERE partitionColumn > 9000
Run Code Online (Sandbox Code Playgroud)

  • 这对于 JdbcRDD 是 100% 准确的([参见代码](https://github.com/apache/spark/blob/17edfec59de8d8680f7450b4d07c147c086c105a/sql/core/src/main/scala/org/apache/spark/sql/execution/数据源/jdbc/JDBCRelation.scala#L85-L97))。特别是,如果您将“upperBound”设置得太低,则一个执行程序将比其他执行程序执行更多的工作,并且可能会耗尽内存。 (3认同)

小智 8

创建分区不会由于过滤而导致数据丢失。的upperBoundlowerbound随着numPartitions仅仅定义分区如何被创建。的upperBoundlowerbound不限定用于partitionColumn的值的范围内(过滤器),以被取出。

For a given input of lowerBound (l), upperBound (u) and numPartitions (n) 
The partitions are created as follows:

stride, s= (u-l)/n

**SELECT * FROM table WHERE partitionColumn < l+s or partitionColumn is null**
SELECT * FROM table WHERE partitionColumn >= l+s AND <2s  
SELECT * FROM table WHERE partitionColumn >= l+2s AND <3s
...
**SELECT * FROM table WHERE partitionColumn >= l+(n-1)s**
Run Code Online (Sandbox Code Playgroud)

例如upperBound = 500lowerBound = 0numPartitions = 5。分区将根据以下查询:

SELECT * FROM table WHERE partitionColumn < 100 or partitionColumn is null
SELECT * FROM table WHERE partitionColumn >= 100 AND <200 
SELECT * FROM table WHERE partitionColumn >= 200 AND <300
SELECT * FROM table WHERE partitionColumn >= 300 AND <400
...
SELECT * FROM table WHERE partitionColumn >= 400
Run Code Online (Sandbox Code Playgroud)

根据的实际值范围partitionColumn,每个分区的结果大小会有所不同。


Hem*_*wda 6

想要添加到经过验证的答案中,因为

没有它们,您可能会丢失一些数据,从而产生误导。

从文档中, 请注意,lowerBound和upperBound仅用于确定分区跨度,而不是用于过滤表中的行。因此,表中的所有行都将被分区并返回。此选项仅适用于阅读。

这表示您的表格有1100行,而您指定

lowerBound 0

upperBound 1000和

numPartitions:10,则不会丢失1000至1100行。您最终将得到一些分区比预期更多的行(跨步值为100)。

  • 你知道 Spark 如何处理剩下的 100 行吗?例如,这是否意味着您的 10 个分区将有 110 行? (2认同)