programmerloverun commented on PR #12292:
URL: https://github.com/apache/seatunnel/pull/12292#issuecomment-5652334918
# 1. Background
When SeaTunnel JDBC Source performs parallel reads, it first plans Splits
for the table based on the partition column, and then assigns the generated
Splits to different Readers for parallel execution.
The current Split planning process mainly consists of two parts:
1. Determine the data boundaries of each Split.
2. Generate all Splits and assign them to Readers.
For tables of normal size, this mechanism usually works well. However, for
very large tables containing tens or hundreds of billions of rows, skewed
partition columns, or multi-table synchronization scenarios, the current
implementation can introduce significant planning overhead and memory pressure.
## 1.1 Excessive Large-Table Scanning Caused by Client-Side Sampling
When the partition column is unevenly distributed, the current
implementation may use `sampleDataFromColumn()` to sample the partition column.
The current sampling process is essentially:
```text
Query the partition column from the database
↓
Iterate through the ResultSet on the client
↓
Keep one sample every N rows
↓
Sort the samples
↓
Calculate Split boundaries
```
Although only a small number of samples are eventually retained, the JDBC
client still needs to iterate through a large amount of data.
For example, suppose a table contains 10 billion rows:
```text
10 billion split keys
↓
Returned through JDBC ResultSet
↓
Scanned row by row on the client
↓
Keep one sample every 10,000 rows
```
Even if only 1 million samples are retained in the end, the cost of database
queries, network transmission, and client-side `ResultSet` traversal remains
very high.
Therefore, the main problem with the current sampling strategy is not the
number of retained samples, but rather:
> A large amount of raw data must be scanned in order to obtain a relatively
small number of samples.
---
## 1.2 Eager Split Materialization Causes Memory Usage to Grow with Data Size
After Split planning is completed, the current implementation generates all
Splits for the entire table at once.
The overall process is roughly:
```text
Entire Table
↓
Calculate all ChunkRanges
↓
List<ChunkRange>
↓
Convert all ranges to JdbcSourceSplit
↓
List<JdbcSourceSplit>
↓
Enumerator
↓
Assign all Splits to Readers
↓
Reader Queue
```
Assume:
```text
split.size = 100000
Table size = 100 billion rows
```
This may theoretically produce:
```text
1 million Splits
```
These Splits may occupy JVM memory simultaneously at multiple stages:
```text
ChunkRange List
+
JdbcSourceSplit List
+
Enumerator Pending Splits
+
Reader Pending Splits
+
Checkpoint State
```
As a result, the memory consumed by Split metadata increases with the total
number of Splits.
For very large tables or jobs reading multiple tables concurrently, this can
easily lead to:
- Continuous growth of Coordinator memory usage
- Increased Reader JVM memory usage
- Frequent Full GC
- Oversized Checkpoint State
- OOM in extreme cases
---
## 1.3 Simply Reducing the Number of Splits Does Not Fully Solve the Problem
A straightforward solution is to follow the Spark JDBC approach:
```text
Number of Splits ≈ Parallelism
```
For example:
```text
parallelism = 4
```
Only generate:
```text
Split 0
Split 1
Split 2
Split 3
```
This can indeed reduce the memory consumed by Split metadata.
However, in a data-skew scenario such as:
```text
Split 0 1 million rows
Split 1 1 million rows
Split 2 1 million rows
Split 3 5 billion rows
```
the last Reader becomes a severe straggler.
Therefore, this optimization should not simply solve the problem by reducing
the total number of Splits.
Instead, the goal is to:
> Preserve a reasonable Split granularity while avoiding generating and
storing all Splits at once.
---
# 2. Goals
This optimization mainly addresses two problems:
```text
High Split boundary planning cost
+
Eager materialization of large amounts of Split metadata
```
The overall goals are described below.
## 2.1 Reduce Split Planning Cost for Large Tables
Avoid scanning the entire partition column on the client merely to calculate
Split boundaries.
Different strategies are used depending on the data distribution:
```text
Uniformly distributed data
↓
Arithmetic Range Splitting
Unevenly distributed data
↓
Index Probing to find the next boundary
```
This prevents Split planning cost from growing linearly with the total
amount of table data.
---
## 2.2 Avoid Generating All Splits at Once
Change the current behavior from:
```text
Generate all Splits at once
```
to:
```text
Generate Splits on demand
```
The Enumerator no longer keeps the Split list for the entire table.
Instead, it only maintains:
```text
Current Split Generator State
+
A small number of generated but unassigned Splits
+
A small number of assigned but unfinished Splits
```
This reduces the memory complexity of Split metadata from:
```text
O(totalSplitCount)
```
to approximately:
```text
O(parallelism)
```
---
## 2.3 Prevent Readers from Accumulating Large Numbers of Splits
Currently, a Reader may receive a large number of Splits at once and store
them in a local queue.
After the optimization, each Reader only holds a small number of Splits:
```text
Reader finishes the current Split
↓
Request the next Split from Enumerator
↓
Enumerator dynamically generates / assigns a Split
```
This prevents Reader memory usage from increasing with the total number of
Splits in the job.
---
## 2.4 Preserve Read Correctness and Recovery Semantics
The optimization mainly changes:
```text
How Splits are generated
When Splits are generated
How Splits are assigned
```
It does not change the data-reading semantics of each individual Split.
The implementation must still guarantee that:
- All data is completely covered.
- No data ranges are missed between Splits.
- Split boundaries do not cause duplicate reads.
- Splits can be re-executed after Reader failures.
- Split planning progress can be restored from checkpoints.
---
# 3. Overall Architecture
The optimization consists of two main parts:
```text
Split Boundary Planning
+
Split Lifecycle Management
```
`Split Boundary Planning` addresses:
> How Split boundaries can be calculated efficiently.
`Split Lifecycle Management` addresses:
> How Splits are generated, stored, and assigned without accumulating all
Splits in memory at once.
<img width="1226" height="1283" alt="image"
src="https://github.com/user-attachments/assets/3eb12026-c273-4620-8418-7cb21f4c1495"
/>
---
## 3.1 Split Boundary Generation
The system first determines whether the partition column is relatively
evenly distributed based on the available statistics.
### Uniformly Distributed Data
For relatively uniform data, the existing arithmetic range splitting
strategy can continue to be used.
For example:
```text
MIN(id) = 1
MAX(id) = 1000
Expected step corresponding to Split Size = 250
```
The following Splits can be generated:
```text
Split 1 : [1, 251)
Split 2 : [251, 501)
Split 3 : [501, 751)
Split 4 : [751, +∞)
```
This approach does not require scanning the actual table data.
Only the following information is required:
```text
MIN
MAX
RowCount
```
Subsequent Splits can then be calculated arithmetically.
---
## 3.2 Unevenly Distributed Data
For partition columns with significant data skew, Split boundaries are no
longer calculated through client-side full-column sampling.
Instead, an Index Probing strategy similar to Flink CDC is used.
Assume:
```text
currentBoundary = 1000
splitSize = 100000
```
The database is queried to:
```text
Start from id >= 1000
↓
Find the id around the 100,000th row
```
For example, suppose the database returns:
```text
nextBoundary = 3520
```
Then the following Split is generated:
```text
[1000, 3520)
```
The next probing operation starts from:
```text
3520
```
The overall process becomes:
```text
current = MIN
│
▼
Query the key around splitSize rows ahead
│
▼
nextBoundary
│
▼
Generate:
[current, nextBoundary)
│
▼
current = nextBoundary
│
▼
Generate the next Split
```
This eliminates the need to:
```text
SELECT the entire partition column
↓
Traverse the entire JDBC ResultSet
↓
Perform client-side sampling
↓
Sort samples on the client
```
Each query only determines the next boundary required for the current Split.
---
## 3.3 Lazy Split Generation
One of the core changes is to replace the current:
```java
List<ChunkRange>
```
model with a model similar to:
```java
SplitGenerator.next()
```
### Current Model
```text
generate()
↓
Split1
Split2
Split3
Split4
...
Split1000000
↓
All stored in memory
```
### Optimized Model
```text
next()
↓
Split1
next()
↓
Split2
next()
↓
Split3
```
The Splitter only needs to maintain the current generation state, for
example:
```text
minValue
maxValue
currentBoundary
chunkSize
finished
```
It no longer needs to store hundreds of thousands or millions of Splits that
have not yet been executed.
---
## 3.4 Split Pull / Request Model
The Split assignment model is also changed from:
```text
Enumerator Push All
```
to:
```text
Reader Pull Split
```
### Current Model
```text
Enumerator
│
├── Split1
├── Split2
├── Split3
├── ...
└── Split100000
│
▼
Reader Queue
```
### Optimized Model
```text
Reader
│
│ requestSplit
▼
Enumerator
│
│ splitter.next()
▼
Split
│
▼
Reader
```
After a Reader finishes its current Split, it requests another Split.
Therefore, even if:
```text
The entire job eventually executes 1 million Splits
```
the number of Split objects actually present in JVM memory during execution
may only be:
```text
Number of Readers × Small Buffer
```
For example:
```text
parallelism = 16
Splits simultaneously in memory ≈ 16 ~ 32
```
instead of:
```text
1 million
```
---
## 3.5 Checkpoint Recovery
With Lazy Split Generation, checkpoints no longer need to store all
unexecuted Splits.
Only the following information needs to be persisted:
```text
Split Generator State
+
Assigned but unfinished Splits
+
A small number of Pending Splits
```
For example, suppose Split planning has reached:
```text
[1,100)
[100,300)
[300,700)
currentBoundary = 700
```
The checkpoint can store:
```text
currentBoundary = 700
maxBoundary = 10000000
```
After recovery:
```text
Continue generating Splits from 700
```
There is no need to pre-store:
```text
Split4
Split5
Split6
...
Split1000000
```
This significantly reduces the size of the Checkpoint State.
---
# 4. Scope of Changes
This optimization mainly affects three areas of JDBC Source:
```text
Split Planning
Enumerator
Reader
```
In principle, the actual JDBC data-reading logic should remain unchanged.
---
## 4.1 DynamicChunkSplitter
`DynamicChunkSplitter` is the core component of this optimization.
Its current responsibilities are roughly:
```text
Calculate data distribution
↓
Calculate all ChunkRanges
↓
Generate all Splits
```
After the optimization:
```text
Initialize data distribution information
↓
Maintain the current planning state
↓
Generate the next Split on demand
```
The following capabilities may need to be added or adjusted:
```java
open()
hasNext()
nextSplit()
snapshotState()
restoreState()
```
Internally, the Splitter maintains:
```text
minValue
maxValue
currentBoundary
splitSize
distributionFactor
finished
```
instead of:
```java
List<ChunkRange>
```
---
## 4.2 Split Logic for Uniformly Distributed Data
The existing arithmetic splitting algorithm for uniformly distributed data
can be retained.
The main change is not:
```text
How to calculate a Split
```
but:
```text
When to calculate a Split
```
Currently:
```text
Calculate all Splits in a single loop
```
After the optimization:
```text
Call nextSplit()
↓
Calculate the next boundary
based on currentBoundary
↓
Return one Split
```
For example:
```text
current = 1
step = 250
```
First call:
```text
[1,251)
```
Second call:
```text
[251,501)
```
Continue until:
```text
MAX
```
is reached.
---
## 4.3 Split Logic for Unevenly Distributed Data
For unevenly distributed data, Index Probing is preferred to determine the
next boundary.
This gradually replaces the current default client-side full-column sampling
strategy.
### Current Process
```text
Query split column
↓
Client scans a large ResultSet
↓
Inverse sampling
↓
Sort samples
↓
Calculate all Splits at once
```
### Optimized Process
```text
currentBoundary
↓
Execute Boundary Query
↓
nextBoundary
↓
Return one Split
```
The Boundary Query should use the partition column index whenever possible.
Different databases use different pagination syntax, such as `LIMIT`,
`OFFSET`, or `FETCH`. Therefore, SQL generation should preferably reuse
existing JDBC Dialect capabilities or be moved into the JDBC Dialect layer.
---
## 4.4 JdbcSourceEnumerator
The Enumerator no longer obtains all Splits for the entire table in advance.
The current logic is roughly:
```text
prepareSplits()
↓
List<JdbcSourceSplit>
↓
Put all Splits into pending state
↓
Assign Splits to Readers
```
After the optimization:
```text
Reader Request
↓
Check Pending Splits
↓
No available Pending Split
↓
splitter.nextSplit()
↓
Assign to Reader
```
The Enumerator only needs to maintain a small amount of state:
```text
DynamicChunkSplitter State
Pending Splits
Assigned But Not Finished Splits
```
This prevents Enumerator memory usage from increasing with the total number
of Splits.
---
## 4.5 handleSplitRequest
Currently, Split assignment between Reader and Enumerator mainly relies on
eager distribution.
This optimization should make actual use of:
```java
handleSplitRequest()
```
so that Readers can actively request new Splits.
The process becomes:
```text
Reader starts
↓
requestSplit
↓
Enumerator
↓
nextSplit
↓
assignSplit
↓
Reader executes the Split
↓
Split completed
↓
requestSplit
```
If all Splits have already been generated:
```java
splitter.hasNext() == false
```
the Enumerator sends:
```text
NoMoreSplits
```
to the Reader.
---
## 4.6 JdbcSourceReader
Readers should no longer receive and retain a large number of Splits.
The Reader queue should remain bounded.
Ideally, each Reader should only hold:
```text
Current Split
+
At most a small number of prefetched Splits
```
For example:
```text
1 ~ 2 Splits
```
After the current Split is completed, the Reader:
```text
Notifies / requests the Enumerator
```
to obtain the next Split.
As a result, Reader memory usage is no longer significantly affected by:
```text
The total number of Splits in the job
```
---
## 4.7 Enumerator State / Checkpoint
Enumerator State needs to include the Split Generator state.
This mainly includes:
```text
Current table
Current partition column
minValue
maxValue
currentBoundary
splitSize
Current splitting strategy
Whether Split planning has finished
```
It should also store:
```text
Assigned but unfinished Splits
```
After failure recovery:
1. Restore unfinished Splits.
2. Restore the Split Generator State.
3. Continue generating new Splits from the last saved boundary.
This ensures that:
```text
No Splits are lost
No data is missed
```
---
## 4.8 JDBC Dialect
Index Probing requires database-specific SQL for determining the next Split
boundary.
Logically, the query needs to:
```text
Start from currentBoundary
Order by the partition key
Find the row around splitSize
Return the corresponding key
```
Different databases may use:
```text
LIMIT
OFFSET
FETCH NEXT
ROWNUM
```
Therefore, we may consider adding an API such as:
```java
buildNextChunkBoundaryQuery(...)
```
to the JDBC Dialect, or reuse the existing pagination capabilities.
The first stage should prioritize support for common JDBC Source databases.
Database-native sampling features such as:
```text
SAMPLE
TABLESAMPLE
```
are not part of the core scope of this optimization.
They can be introduced later as database-specific optimizations.
---
## 4.9 Configuration Compatibility
The existing:
```properties
split.size
```
semantics should remain unchanged as much as possible.
`split.size` should continue to represent:
> The expected amount of data processed by one Split.
This optimization does not recommend simply introducing:
```text
max.total.split.count
```
and forcing the entire job to generate only a fixed number of Splits.
Doing so may result in excessively large individual Splits, which can cause:
- Data skew
- Reader stragglers
- Increased amount of repeated reading after failures
- Excessively long execution time for a single Split
If the number of Splits held in memory needs to be controlled, the limit
should instead apply to:
```text
max.pending.splits
```
which means:
> The maximum number of generated
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]