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]
