spuru9 opened a new pull request, #29352:
URL: https://github.com/apache/flink/pull/29352
## What is the purpose of the change
`ProcessTableFunctionSemanticTests[process-stateful-multi-input-with-timeout]`
is flaky on master, release-2.3 and release-2.2 (see FLINK-40571).
The two inputs of the PTF are processed in arrival order, not event-time
order (`on_time` does not affect the processing order). `TimedJoinFunction`
keeps a single score per key, so if the city row arrived after several scores,
the earlier scores were overwritten and only one join row was emitted at the
city row's timestamp, e.g. `+I[Bob, Bob, 6 score in city London,
1970-01-01T00:00:00Z]`. The pinned output only matched the city-first
interleaving.
Neither `ORDER BY` (sort buffers flush per input, so cross-input order is
still not guaranteed) nor a materialized assertion as in FLINK-40568 (the
output is append-only) fixes this, so the test data is restructured instead.
## Brief change log
- Give the program its own score source and extend `TIMED_CITY_SOURCE`
(only used by this program), so that the output is independent of the
interleaving:
- a key with a city has a single score with the same timestamp as its
city (Bob, Charly, Dave): the join is emitted with the same row and timestamp
regardless of which side arrives first
- keys without a city may have multiple scores (Alice, Frank): they come
from a single input, so their order is preserved; Alice also covers overwriting
the score and replacing the named timer
- a city without a score (Eve) emits nothing
- Timeouts stay deterministic as timers only fire once the watermark of
both inputs passed them, i.e. after both inputs have finished.
## Verifying this change
This change is a test-only fix and is covered by the modified test itself.
Verified locally with a temporary test (not committed) that runs the exact
data of the program:
- normal interleaving: 100/100 passed
- all cities forced to arrive after all scores
(`source.sleep-after-elements`): 5/5 passed
- all scores forced to arrive after all cities: 5/5 passed
- with the old data and cities forced to arrive late, the test fails 5/5
with the same error as in CI
- `ProcessTableFunctionSemanticTests`: 10/10 runs passed
## 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)`: no
- 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? no
- If yes, how is the feature documented? not applicable
---
##### Was generative AI tooling used to co-author this PR?
- [X] Yes (please specify the tool below)
Generated-by: Claude Code (claude-opus-5-5)
🤖 Generated with [Claude Code](https://claude.com/claude-code)
--
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]