leaves12138 commented on code in PR #10036:
URL: https://github.com/apache/paimon/pull/10036#discussion_r4059240870
##########
paimon-python/pypaimon/read/streaming_table_scan.py:
##########
@@ -410,6 +412,13 @@ def _try_native_plan(self, expected_snapshot_id: int,
def _create_changelog_plan(self, snapshot: Snapshot) -> Plan:
"""Read from changelog_manifest_list
(changelog-producer=input/full-compaction/lookup)."""
+ plan = self._try_native_plan(
+ snapshot.id,
+ incremental_range=(snapshot.id - 1, snapshot.id),
+ incremental_mode='changelog',
+ )
Review Comment:
[P1] Preserve overwrite changelog frames before accepting the native plan
This reuses batch incremental-changelog semantics for a continuous-streaming
frame, but the two paths disagree on OVERWRITE.
`ChangelogFollowUpScanner.should_scan()` accepts every snapshot containing a
changelog manifest. Such snapshots are real: writing to a
`changelog-producer=input` PK table through
`new_batch_write_builder().overwrite()` produces an OVERWRITE snapshot with
physical changelog files. In contrast, Rust `IncrementalScan::plan_changelog()`
explicitly skips OVERWRITE snapshots. The native call consequently returns an
empty plan with the expected snapshot ID, so it is accepted here rather than
triggering fallback, and streaming advances beyond unread changes.
Reproduced with this head and Rust #894 (65e8400), using only ordinary
writes and the public streaming builder:
1. Snapshot 1: insert `(k=1, v='old')`.
2. Snapshot 2: overwrite with `(k=2, v='replacement')`.
3. Snapshot 3: insert `(k=3, v='later')`.
4. Resume streaming from snapshot 2 with row kinds enabled.
On this PR's base (1aa0adb), all four Python/native planner-reader
combinations return `+I replacement` in frame 2, then `+I later` in frame 3. On
this head, enabling native planning produces an empty frame 2 and advances
`next_snapshot_id` to 3, regardless of which reader is selected. Disabling
native planning still returns the replacement event.
Please preserve the existing streaming behavior, e.g. keep these OVERWRITE
frames on the Python fallback or expose a single-snapshot native changelog plan
that does not apply the batch commit-kind exclusion. The batch API's existing
OVERWRITE policy should not be changed indiscriminately. Add an actual
overwrite/resume regression alongside the APPEND-only streaming test.
--
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]