Skip to content

[Client] Support bounded log scanning - #4545

Open
yangshangqing95 wants to merge 1 commit into
apache:mainfrom
yangshangqing95:feat/client/bounded-log-scan
Open

yangshangqing95 wants to merge 1 commit into
apache:mainfrom
yangshangqing95:feat/client/bounded-log-scan

Conversation

@yangshangqing95

@yangshangqing95 yangshangqing95 commented Oct 4, 2026 •

Copy link
Copy Markdown

Summary

The runtime change is intentionally small, a significant portion of this PR is test coverage for boundary conditions, compatibility, remote reads, retention, and Arrow resource ownership.

Purpose

Linked issue: close #4544
This PR adds a bounded range scan capability to LogScanner, allowing callers to subscribe to a bucket within an exclusive offset range:

[startingOffset, stoppingOffset)

Design

             returned records
        ┌──────────────────────┐
        ▼                      ▼
log: ... start ............... stop ........
        [                      )
                               │
                               └─ finishedBuckets
                                  reported once

Reports completion once per bounded subscription through finishedBuckets(). 
A result reporting completion may also contain the final records, 
which callers must consume before completing the read.

Each subscribed bucket keeps its own stopping offset. The client:

  • never exposes records at or beyond stoppingOffset;
  • caps logical consumed progress at the stopping offset;
  • reports completion once through finishedBuckets();
  • supports empty ranges, including start >= stop;
  • supports both explicit starting offsets and EARLIEST_OFFSET;
  • applies the same bounded semantics to row and Arrow scans.

For an explicit starting offset, the client can enforce the range directly.

For EARLIEST_OFFSET, the client additionally needs to know the physical offset to which the server resolved the EARLIEST request. This PR therefore adds the optional resolved_earliest_offset field to the fetch response.

Server changes are limited to reporting the resolved EARLIEST offset. Enforcing the stopping offset on the server and migrating connectors are outside this PR.

Client                                  Server

subscribeBounded(EARLIEST, stop)
    │
    │ Fetch(offset = EARLIEST)
    ├───────────────────────────────────►
    │
    │       resolve EARLIEST -> physical offset
    │
    │ response(resolved_earliest_offset)
    ◄───────────────────────────────────┤
    │
    │ apply [resolved offset, stop)
    │
    ▼
bounded result / completion

Why not use log_start_offset?

log_start_offset may eventually be a natural source for this information, but using it in this PR would significantly broaden the scope.

Its complete client/server semantics are not currently established across the relevant fetch paths, and the server still does not consistently populate it with the physical log start offset. Defining and enabling that behavior deserves a separate and more detailed design discussion because it can affect existing fetch, retention, remote-log, and compatibility semantics.

There is also a difference:

  • log_start_offset describes a property/boundary of the log;
  • resolved_earliest_offset describes how a particular EARLIEST fetch request was resolved.

Keeping the latter as an optional request specific response field also allows a new client to distinguish:

For a successful fetch response to an EARLIEST_OFFSET request:

field absent    -> the response does not provide the required resolved EARLIEST offset
field = 0       -> EARLIEST was legitimately resolved to offset 0
field = N       -> EARLIEST was resolved to offset N

This lets bounded EARLIEST scans fail explicitly against an older server instead of silently making an incorrect range decision.

If log_start_offset gets well defined and consistently implemented in the future, and is sufficient for bounded EARLIEST resolution, the implementation can be consolidated in a follow-up. That broader protocol/storage design is intentionally not part of this PR.

Brief change log

  • Add additive subscribeBounded(...) APIs to LogScanner.
  • Add optional resolved_earliest_offset to fetch responses for EARLIEST resolution.
  • Preserve existing error delivery and resource cleanup semantics while ensuring already consumed bounded progress is delivered before a later error.
  • Keep existing unbounded public APIs and their existing semantics unchanged.

Tests

Added/updated unit and integration coverage for:

  • explicit bounded ranges and exclusive stopping offsets;
  • empty ranges: start == stop, start > stop, and stop == 0;
  • one shot completion reporting;
  • row and Arrow bounded scans;
  • Arrow batch truncation and off-heap resource ownership;
  • filtered empty responses before and at the stopping offset;
  • EARLIEST_OFFSET resolution before, at, and beyond the stopping offset;
  • retention where the physical earliest offset has advanced beyond the bounded range;
  • local and remote EARLIEST resolution propagation;
  • remote segment pruning at the stopping offset;
  • stale fetch responses after re-subscription;
  • new-client / old-server behavior for bounded EARLIEST;
  • error delivery after bounded records, progress, or completion;
  • preservation of existing unbounded error and queue behavior;
  • legacy LogScanner implementations using the existing subscription APIs;
  • RPC serialization/deserialization of the new optional field.

API and Format

This PR adds new API surface but does not change the existing unbounded LogScanner APIs.

The bounded subscription methods are added as default interface methods so existing LogScanner implementations remain compatible. Implementations that do not support bounded scans continue to work and reject the new operation with UnsupportedOperationException.

Existing unbounded Log Table public API behavior is intentionally preserved.

The fetch protocol adds:

optional int64 resolved_earliest_offset = 11;

This is an additive optional protobuf field:

  • old clients can ignore it when talking to a new server;
  • new clients do not require it for existing unbounded scans or bounded scans with an explicit starting offset;
  • bounded scans starting from EARLIEST_OFFSET require the server to report the resolved starting offset unless the range can be completed locally, such as when stoppingOffset == 0. If the required metadata is absent, polling fails with UnsupportedOperationException.

Documentation

The new bounded subscription semantics are documented in the LogScanner API, including:

  • exclusive [startingOffset, stoppingOffset) semantics;
  • supported starting offsets;
  • empty-range behavior;
  • completion reporting;
  • the compatibility requirement for bounded scans starting from EARLIEST_OFFSET.

@yangshangqing95

Copy link
Copy Markdown
Author

Hi @fresh-borzoni @polyzos , would appreciate it if you could take a look at this PR when you have time, thx.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

[client] Support bounded Log Table scanning

1 participant