SEPURI-SAI-KRISHNA opened a new pull request, #29345:
URL: https://github.com/apache/flink/pull/29345
## What is the purpose of the change
FLIP-567 lists test snapshotting and restoration as part of the PTF test
harness, so that a test
setup can be built once and reused by several cases. This adds it:
`ProcessTableFunctionTestHarness.snapshot()` captures a test at a point in
time and the static
`ProcessTableFunctionTestHarness.restoreFromSnapshot(TestSnapshot)` builds a
new harness that
continues from there.
Two points settled with @autophagy on the ticket:
- Only the static `restoreFromSnapshot` is added. The
`Builder.restoreFromSnapshot` variant the FLIP
also shows was dropped as unintuitive, so the FLIP page wants that line
removing at some point.
- State and timers are deep copied when the snapshot is taken, so that
mutating them afterwards
cannot change the snapshot.
The snapshot holds the configuration the harness was built with, the state
of every partition, the
pending and the fired timers, the watermark of every table argument, and the
collected output. The
system clock the FLIP also mentions is not covered, because the harness has
no system clock yet
(`withInitialSystemClock` is not implemented).
This overlaps with #29130 (FLINK-40577) on `ProcessTableFunctionTestHarness`
and on `ptfs.md`.
Happy to rebase this one on top of it, whichever lands first.
## Brief change log
- `ProcessTableFunctionTestHarness`: `snapshot()`,
`restoreFromSnapshot(TestSnapshot)` and the
`TestSnapshot<OUT>` holder. The harness keeps a copy of its `Builder`, so
a snapshot can rebuild
an equivalently configured harness, and further configuration of the
original builder cannot
change what a snapshot restores.
- `StateConverter` gains `copyInternal`. State is stored in its internal
representation
(`ArrayData` / `MapData` / `RowData`), even for POJO state, since every
write goes through
`StateConverter.toInternal` first. The copy is therefore driven off the
state's internal
serializer (`InternalSerializers`) and never instantiates or touches a
user-facing state object.
`StructuredTypeStateConverter` now takes the state `DataType` instead of
only its conversion
class, so it can build that serializer.
- `TestHarnessStateManager`: `snapshotState()` / `restoreState(...)`, both
deep copying. A restore
replaces the live state rather than merging into it.
- `TestHarnessTimerManager`: `snapshot()` / `restore(...)` covering pending
timers, fired timers,
`watermarkByTable` and `globalWatermark`. `globalWatermark` matters on its
own: it is what
`TimeContext.currentWatermark()` reads, and it is what rejects a table
whose first watermark
arrives after the restore and would pull the global watermark backwards.
- `Timer.copy()`, since a `Timer` carries the mutable `fired` flag and a
`Row` key.
- Documentation: a new "Snapshotting and Restoring a Test" section in
`ptfs.md` (en and zh).
## Verifying this change
This change added tests and can be verified as follows:
- `ProcessTableFunctionTestHarnessSnapshotTest` (new, 13 tests) covers
snapshot and restore of
output, POJO / `Row` / `ListView` / `MapView` state, pending and fired
timers and watermarks;
that a snapshot is unaffected by later elements, state changes, timer
firings and builder
configuration; that a snapshot can be restored more than once into
independent harnesses; that a
restore uses a fresh function instance; and that state cleared before the
snapshot does not come
back from the builder's initial state.
- `StateConverterTest` (new, 5 tests) covers `copyInternal` for all four
converters, including that
mutating a nested element, map value or row after the copy is not visible
through it.
- `TestHarnessTimerManagerTest` (+6 tests) covers the timer manager's
snapshot and restore directly.
- Every production change was checked red/green by reverting it in isolation
and confirming a named
test fails: no value copy in the converters (4 failures in
`StateConverterTest`), the per-key
state map shared instead of copied, `Timer.copy()` returning `this`,
`globalWatermark` not
restored, `watermarkByTable` not restored, pending timer sets shared, the
output list shared,
the builder not copied, and each of state / timers / output / function
output not restored.
- `mvn verify` on `flink-table/flink-table-test-utils` passes: 139 tests,
plus spotless and
checkstyle.
- The example in the documentation was run as written before being added.
## Does this pull request potentially affect one of the following parts:
- Dependencies (does it add or upgrade a dependency): no
- The public API, i.e., is any changed class annotated with
`@Public(Evolving)`: yes - two new
methods and one new nested class on `ProcessTableFunctionTestHarness`,
which is
`@PublicEvolving`, and a package-private method on `Timer`. All additive.
- The serializers: no
- The runtime per-record code paths (performance sensitive): no
- Anything that affects deployment or recovery: JobManager (and its
components), Checkpointing, Kubernetes/Yarn, ZooKeeper: no
- The S3 file system connector: no
## Documentation
- Does this pull request introduce a new feature? yes
- If yes, how is the feature documented? docs and JavaDocs
---
##### Was generative AI tooling used to co-author this PR?
- [X] Yes (please specify the tool below)
Generated-by: Claude Code (Opus 5)
--
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]