programmerloverun commented on issue #12097:
URL: https://github.com/apache/seatunnel/issues/12097#issuecomment-5550865886

   @DanielLeens @davidzollo Here is my proposed solution. I would like to go 
over it with you and check whether there are any missing parts or unreasonable 
aspects.
   ## 1. Background
   
   The SeaTunnel JDBC Source, before starting parallel reads, first plans 
splits based on a split key: depending on the data distribution, it chooses 
uniform splitting, iterative non-uniform splitting, or sampling-based splitting.
   
   The current implementation has two key paths:
   
   1. **Sampling-based planning**: `sampleDataFromColumn()` queries the split 
column; the client iterates over the `ResultSet`, keeps samples at 
`inverse-sampling.rate`, then sorts them to determine boundaries. The default 
path essentially scans the entire split-column result set and keeps one row out 
of every N.
   2. **Split materialization and dispatch**: all splits are generated at once 
— the Enumerator splits the entire table and hands everything to the Reader in 
one go; the Reader buffers them in an unbounded queue.
   
   On small and medium tables this strategy is usually acceptable; but in batch 
sync scenarios with **tens to hundreds of billions of rows, skewed keys, 
multiple tables, and expensive planning scans against the source database**, 
planning cost and JVM memory scale linearly (or even super-linearly) with data 
size. The community already has a sampling switch 
([#10596](https://github.com/apache/seatunnel/issues/10596) / 
[#10604](https://github.com/apache/seatunnel/pull/10604)) and a database-side 
sampling proposal ([#10288](https://github.com/apache/seatunnel/issues/10288)), 
but neither covers the combined problem of "sampling scan + full 
materialization of split metadata".
   
   ## 2. Goal
   
   In ultra-large-table, skewed-key, multi-table JDBC Source scenarios, 
eliminate slow planning, memory bloat, and OOM risk caused by **full-scale 
sampling scans + fully materialized/buffered split metadata**, without 
sacrificing the correctness and recoverability of parallel reads.
   
   ## 3. Design
   
   ### 3.1 Reference Approaches
   
   1. **Apache Spark JDBC: arithmetic range splitting, with the number of 
splits bound to parallelism**
   
   Suppose a table has a numeric column `id`, roughly ranging from 1 to 1000, 
and you want to read this table with 4-way parallelism.
   
   You tell Spark four things:
   
   - Which column to split by: `id`
   - Minimum value: 1
   - Maximum value: 1000
   - Number of partitions: 4 (usually roughly equal to your desired parallelism)
   
   This generates 4 ranges (illustrative):
   
   | Split | Approximate condition | Reader |
   | --- | --- | --- |
   | 0   | `id < 251` (edge cases like very small values / `null` omitted) | 
Task 0 |
   | 1   | `id >= 251 AND id < 501` | Task 1 |
   | 2   | `id >= 501 AND id < 751` | Task 2 |
   | 3   | `id >= 751` | Task 3 |
   
   Each Task connects to the database on its own and queries only its own 
range. This is **arithmetic range splitting**: divide using the lower and upper 
bounds to obtain equal-width ranges.
   
   "**Number of splits bound to parallelism**" means: if you ask for 4 splits, 
you get roughly 4 splits — not 1 million splits just because the table has 10 
billion rows.
   
   2. **Alibaba DataX: database-side SAMPLE + number of splits close to 
Channel**
   
   > Channel: the number of concurrent channels in DataX, which can be 
understood as "how many read tasks run at the same time". For example, `channel 
= 4` means roughly 4-way parallel reads.
   > 
   > Database-side SAMPLE: sampling is done inside the database, instead of 
pulling the entire column back to DataX and keeping one row out of every N.
   > Oracle example: `SELECT id FROM t SAMPLE(1)` → the database returns only 
about 1% of the rows (the exact syntax varies by database).
   
   Suppose the table is very large, with 100 million rows.
   
   Instead of pulling 100 million `id`s, we only get back id samples at roughly 
**1%** of the original volume.
   
   Sort the samples and pick boundaries at a partition count close to the 
channel number — e.g. split into 4 ranges, yielding boundaries `a < b < c`.
   
   We then get four ranges: `(-∞, a) [a, b) [b, c) [c, +∞)`
   
   Each Channel/Task performs the full `SELECT` over its own range in parallel: 
the number of splits is bounded, and so is the planning traffic.
   
   3. **Apache Flink CDC: uniform step + index probing, avoiding full-column 
client-side sampling**
   
   Two strategies are used:
   
   **Strategy 1: the key distribution looks uniform — split directly with an 
arithmetic step**
   
   First query `MIN(id)` and `MAX(id)`, then estimate the row count and compute 
a distribution factor. If the distribution looks fairly uniform, split by step:
   
   ```java
   Suppose id ranges from 1 to 1000, and a chunk is roughly 250 rows.
   Split into:
     [1, 251), [251, 501), [501, 751), [751, +∞)
   ```
   
   **Strategy 2: the key distribution is non-uniform — repeatedly "probe" the 
database for the upper bound of the next chunk**
   
   When the value range is huge but the actual row count is small, or the data 
is crowded into one segment, equal-width splitting leads to:
   
   - many empty chunks, or
   - a few extremely heavy chunks.
   
   In this case, switch tactics: starting from the current lower bound, let the 
database tell you **what the upper bound is after roughly `chunkSize` rows 
further down**.
   
   ```java
   current lower bound = 1, chunkSize = 1000
   
   Ask the DB: with id >= 1, roughly 1000 rows further down, what is the 
largest id?
   DB answers: 3520
   → chunk [1, 3520)
   
   Ask again: with id >= 3520, roughly 1000 rows further down, what is the 
largest id?
   DB answers: 3605
   → chunk [3520, 3605)
   
   ……until MAX is covered
   ```
   
   ### 3.2 What We Do
   
   1. Following Spark/DataX, use a **split cap** to solve split explosion.
   2. Following Flink, use **uniform arithmetic splitting + boundary probing 
for non-uniform distributions** to replace full-column client-side sampling. 
SAMPLE and batched dispatch are not the main path in this phase. We only 
enhance the JDBC Source's generic splitting capability, keeping 
read-completeness and recoverability semantics while controlling memory and 
planning cost.
   
   <img width="2092" height="1212" alt="Image" 
src="https://github.com/user-attachments/assets/58d29f22-48c9-4c28-ad0f-6a08d0fa3e0a";
 />


-- 
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