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]

Reply via email to