Bha*_*rla 24 jdbc apache-spark apache-spark-sql
虽然通过在星火JDBC连接获取来自SQL Server的数据,我发现我可以设置一些并行的参数,如partitionColumn,lowerBound,upperBound,和numPartitions.我已经通过spark文档,但无法理解它.
谁能解释一下这些参数的含义?
小智 25
很简单:
partitionColumn 是一个应该用于确定分区的列.lowerBound并upperBound确定要获取的值的范围.完整数据集将使用与以下查询对应的行:
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:0upperBound:1000numPartitions:10
Stride等于100,分区对应于以下查询:
SELECT * FROM table WHERE partitionColumn BETWEEN 0 AND 100SELECT * FROM table WHERE partitionColumn BETWEEN 100 AND 200...SELECT * FROM table WHERE partitionColumn BETWEEN 900 AND 1000And*_*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)
小智 8
创建分区不会由于过滤而导致数据丢失。的upperBound,lowerbound随着numPartitions仅仅定义分区如何被创建。的upperBound和lowerbound不限定用于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 = 500,lowerBound = 0和numPartitions = 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,每个分区的结果大小会有所不同。
想要添加到经过验证的答案中,因为
没有它们,您可能会丢失一些数据,从而产生误导。
从文档中, 请注意,lowerBound和upperBound仅用于确定分区跨度,而不是用于过滤表中的行。因此,表中的所有行都将被分区并返回。此选项仅适用于阅读。
这表示您的表格有1100行,而您指定
lowerBound 0
upperBound 1000和
numPartitions:10,则不会丢失1000至1100行。您最终将得到一些分区比预期更多的行(跨步值为100)。
| 归档时间: |
|
| 查看次数: |
10482 次 |
| 最近记录: |