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]
