DanielLeens opened a new issue, #12097:
URL: https://github.com/apache/seatunnel/issues/12097

   ### Search before asking
   
   I reviewed #10288 (database-side sampled balanced sharding) and #10596 
(allow disabling sampling). This is a related, broader design discussion: 
reducing sampling scans alone does not bound the memory used by chunk 
generation, split assignment, reader queues, and checkpoint state. Please link 
or consolidate this issue if the community prefers one tracking issue.
   
   ### Description
   
   Could contributors help design a minimally invasive, reusable JDBC 
split-planning improvement that addresses both expensive sampling scans and 
unbounded in-memory split metadata for very large tables?
   
   We need to control database work, sample memory, and the number of 
simultaneously resident splits together, while preserving data coverage, 
reasonable work per split, and recovery semantics. No particular algorithm is 
prescribed here; alternative designs and database execution-plan evidence would 
be very welcome.
   
   #### Current upstream behavior
   
   The following findings were checked against `dev` commit 
[`ce4cde54c8abc7a6a9e3831110dbab029512ef7b`](https://github.com/apache/seatunnel/commit/ce4cde54c8abc7a6a9e3831110dbab029512ef7b).
 They are source-level observations, not the result of a new large-scale 
benchmark or an attached heap dump.
   
   1. **The default sampling path consumes the entire query result before 
retaining every Nth row.** `JdbcDialect.sampleDataFromColumn()` executes 
`SELECT split_column FROM table` (or from the configured subquery), iterates 
the complete ResultSet, retains every `samplingRate`-th value in an 
`ArrayList`, converts it to an array, and sorts it. Thus increasing the inverse 
sampling rate reduces retained samples but not the rows traversed and returned 
through JDBC. A covering index may change the physical access path, but it does 
not make this a bounded sample query. [Default 
implementation](https://github.com/apache/seatunnel/blob/ce4cde54c8abc7a6a9e3831110dbab029512ef7b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/JdbcDialect.java#L385),
 [Oracle 
override](https://github.com/apache/seatunnel/blob/ce4cde54c8abc7a6a9e3831110dbab029512ef7b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connect
 ors/seatunnel/jdbc/internal/dialect/oracle/OracleDialect.java#L333).
   
   2. **Chunk and split collections are eagerly materialized.** 
`DynamicChunkSplitter.createDynamicSplits()` obtains the complete 
`List<ChunkRange>` and builds a second `List<JdbcSourceSplit>`. Both sampled 
and iterative uneven splitting accumulate chunk ranges. This leaves memory 
proportional to the number of generated boundaries/splits even if the sampling 
query becomes cheaper. 
[Conversion](https://github.com/apache/seatunnel/blob/ce4cde54c8abc7a6a9e3831110dbab029512ef7b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/DynamicChunkSplitter.java#L81),
 [sampled and iterative 
splitting](https://github.com/apache/seatunnel/blob/ce4cde54c8abc7a6a9e3831110dbab029512ef7b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/DynamicChunkSplitter.java#L632).
   
   3. **Assignment moves the backlog to readers rather than bounding it.** 
`JdbcSourceSplitEnumerator.run()` generates all splits for a table, groups them 
into pending lists, and assigns each reader its complete list. 
`handleSplitRequest()` currently throws unsupported-operation. 
`JdbcSourceReader.addSplits()` appends all received splits to an uncapped 
deque, and `snapshotState()` copies that deque. The enumerator state also 
contains its pending split lists. 
[Enumerator](https://github.com/apache/seatunnel/blob/ce4cde54c8abc7a6a9e3831110dbab029512ef7b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/JdbcSourceSplitEnumerator.java#L73),
 
[Reader](https://github.com/apache/seatunnel/blob/ce4cde54c8abc7a6a9e3831110dbab029512ef7b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/JdbcSourceReader.java#L38),
 [state](https://github.com/apache/seatunnel/blob/ce4cde54c8abc7a6a9e38
 
31110dbab029512ef7b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/state/JdbcSourceState.java#L32).
   
   4. **Disabling sampling is already supported, but does not solve the full 
resource problem.** With `split.allow-sampling=false`, the uneven-data path 
uses iterative chunk-boundary queries and still builds the full chunk list. 
With sampling enabled, this path is chosen when the distribution heuristic 
considers the key uneven and the estimated split count exceeds 
`split.sample-sharding.threshold`. [Strategy 
selection](https://github.com/apache/seatunnel/blob/ce4cde54c8abc7a6a9e3831110dbab029512ef7b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/source/DynamicChunkSplitter.java#L455).
   
   #### Illustrative scale
   
   For a table with 100 billion rows, an **explicitly configured** 
`split.size=100000`, and inverse sampling rate 1000:
   
   - The estimated split count is 1 million; the actual count depends on key 
distribution and boundary deduplication.
   - The current full-result sampling path can retain approximately 100 million 
sample values and then sort them.
   - A plan with approximately 1 million distinct boundaries produces 
correspondingly large chunk/split lists, assignment payloads, reader queues, 
and recoverable pending state.
   
   These are arithmetic scaling examples, not measured object sizes or a claim 
that every such job necessarily OOMs. Actual memory depends on types, heap 
sizing, serialization, and consumption speed. The inspected upstream default 
`split.size` is **8096**, not 100000. 
[Defaults](https://github.com/apache/seatunnel/blob/ce4cde54c8abc7a6a9e3831110dbab029512ef7b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/config/JdbcSourceOptions.java#L49).
   
   #### Desired outcomes
   
   - **Bound planning memory across the job:** sample values, chunk metadata, 
enumerator backlog, assignment/serialization batches, reader prefetch, and 
checkpoint state. Per-table limits alone are insufficient for multi-table jobs. 
A task may process many small splits over its lifetime without keeping all of 
them resident simultaneously.
   - **Reduce and measure planning-side database work:** distinguish sampling, 
min/max, exact counts, statistics collection, and repeated boundary queries. A 
small result or `fetchSize` does not by itself prove low physical scan cost. 
The inspected Oracle row-count path can also execute `ANALYZE ... COMPUTE 
STATISTICS` unless `skip_analyze` is set, or use exact counting for query-based 
reads; these costs should be accounted for separately. [Oracle row-count 
path](https://github.com/apache/seatunnel/blob/ce4cde54c8abc7a6a9e3831110dbab029512ef7b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/oracle/OracleDialect.java#L224).
   - **Preserve manageable split sizes and skew handling:** simply lowering 
total split count may create very large or skewed splits. The current JDBC 
reader processes a whole split while holding the checkpoint lock, so larger 
splits also affect checkpoint latency and recovery work. A total-count cap and 
a live-buffer cap have different tradeoffs.
   - **Preserve correctness and recovery:** cover the intended query/filter 
scope without gaps or unintended overlap; handle NULLs, duplicate/hot keys, 
sparse ranges, empty/insufficient samples, and stale estimates. If planning 
becomes incremental, define ownership of planned, assigned, returned, and 
completed splits, checkpoint compatibility, and rescaling behavior. Restoring 
must not silently change already established boundaries through fresh sampling.
   - **Keep common logic shared and adaptations narrow:** editing the interface 
default alone will not cover dialects that override sampling. Database-native 
sampling/index behavior must be validated for supported database versions and 
query shapes, including databases without suitable indexes or native sampling.
   - **Keep the change reviewable:** prefer reuse of existing connector APIs 
and a minimal coherent implementation over an unrelated engine redesign; do not 
claim an end-to-end fix while one unbounded stage remains.
   
   #### Questions for contributors
   
   1. What is the smallest practical design that bounds live split state end to 
end: lazy split generation with bounded reader requests, a compact resumable 
range plan, a bounded total split count with an explicit split-size policy, or 
another approach?
   2. How should bounded samples be used to plan many reasonably balanced 
splits without building a full boundary list or falling back to a database 
query for every row-count chunk? Are coarse sampled ranges plus incremental 
refinement viable, and how should their database cost be limited?
   3. Which database-native block sampling, existing histogram/statistics, or 
indexed probing approaches have measured low scan cost? What should the 
explicit fallback be when those capabilities are unavailable? Reservoir 
sampling over the full ResultSet bounds memory but does not remove its 
scan/network cost; SQL-side random filters or `NTILE` should also be evaluated 
by actual plans rather than assumed efficient.
   4. Can the existing `sendSplitRequest()` / `handleSplitRequest()` interfaces 
support bounded JDBC prefetch with safe checkpoint/assignment ordering? How 
should synchronous requests, cancellation, and existing saved state be handled?
   
   #### Suggested validation
   
   - Compare baseline and candidate plans on dense, sparse, and highly skewed 
indexed keys, plus unsupported/no-index cases and filtered queries; report 
scanned rows/blocks, returned rows/bytes, SQL count, planning latency, and heap 
usage.
   - Use a synthetic million-split workload with slow or paused readers. Assert 
that generation stops at the configured budget and that no upstream list, 
downstream queue, in-flight payload, or completion-tracking collection grows 
with total processed split count.
   - Exercise multiple tables under a shared budget, wide metadata, empty 
samples, NULL/duplicate keys, stale statistics, timeout, and cancellation.
   - Inject failures around planning advancement, assignment, split return, and 
checkpoint completion. Verify restored coverage and existing delivery 
guarantees, including compatibility with old state and changes in reader 
parallelism.
   - Check both sampling-enabled and sampling-disabled paths. Improvements to 
one must not hide unbounded behavior in the other.
   
   ### Usage Scenario
   
   Large JDBC batch reads with billions to hundreds of billions of rows, uneven 
keys, multiple tables, and source databases where extra planning scans are 
expensive. The goal is a root-cause improvement that benefits JDBC users across 
dialects, rather than a driver-specific workaround. CDC snapshot splitting is 
related, but has a separate lifecycle and should be evaluated explicitly rather 
than assumed fixed by a JDBC change.
   
   ### Related issues
   
   - #10288: database-side sampled balanced sharding. This issue adds 
end-to-end live split memory and recovery constraints to that discussion.
   - #10596: disabling sampling. The current `split.allow-sampling` option is 
acknowledged; this is not another request for the switch.
   
   ### Are you willing to submit a PR?
   
   - [ ] Yes I am willing to submit a PR!
   
   Opening this for community design input and contributor proposals before 
choosing an implementation. Execution plans, production-scale measurements, and 
small proof-of-concept comparisons would be especially helpful.
   
   ### Code of Conduct
   
   - [x] I agree to follow this project's [Code of 
Conduct](https://www.apache.org/foundation/policies/conduct).
   


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