twalthr opened a new pull request, #29419: URL: https://github.com/apache/flink/pull/29419
## What is the purpose of the change This implements broadcast state for process table functions (PTFs) as proposed in [FLIP-565](https://cwiki.apache.org/confluence/display/FLINK/FLIP-565%3A+Broadcast+State+for+Process+Table+Functions). A table argument with `ArgumentTrait.BROADCAST_SEMANTIC_TABLE` is sent to all virtual processors next to tables with row or set semantics. State entries declared with `@StateHint(StateKind.BROADCAST)` are backed by Flink's broadcast state, shared across all sets, and can only be modified while processing a broadcast table. This enables the broadcast state pattern (e.g. dynamic rules or configuration) in SQL and the Table API. `NOTIFY_STATEFUL_SETS` is out of scope and will be addressed separately. Support in `ProcessTableFunctionTestHarness` will follow in a separate PR. ## Brief change log - `flink-table-common` - New `StateKind { PER_SET, BROADCAST }` referenced by `StateHint#value()` - New `ArgumentTrait`/`StaticArgumentTrait.BROADCAST_SEMANTIC_TABLE` - `StateTypeStrategy` carries a broadcast flag; extraction rejects `ListView` and TTL for broadcast state - `SystemTypeInference` validates that at least one table with row or set semantics exists and rejects `PARTITION BY`/`ORDER BY` on broadcast tables. Similar to tables with set semantics, broadcast tables participate in event-time processing and require `on_time` if declared. - `flink-table-planner` - Broadcast tables are distributed via `BROADCAST_DISTRIBUTED` (`BroadcastPartitioner` in `StreamExecExchange`) - `StreamExecProcessTableFunction` creates a keyed or non-keyed multi-input transformation and passes broadcast information to the runtime - Code generation passes `null` for keyed state while processing a broadcast table - `flink-table-runtime` - `BroadcastStateAdapters` expose operator broadcast state as `MapView`, `ValueView`, and eager value state; writes are rejected while processing tables with row or set semantics - `BroadcastEvalCollector` and `BroadcastInternalTimeContext` reject emitting results and registering timers on the broadcast path - `Context#clearState`/`clearAllState` are scoped to the currently processed table - JavaDoc in `ProcessTableFunction`, `StateHint`, `StateKind`, `ArgumentTrait` and a new "Broadcast State" section in `ptfs.md` (en and zh) ## Verifying this change This change added tests and can be verified as follows: - `TypeInferenceExtractorTest`: extraction of broadcast state and table arguments, including invalid declarations (ListView, TTL, missing main table, incompatible traits) - `ProcessTableFunctionTest`: plan tests for set/row semantic main tables, multiple tables, `on_time`, updating inputs, and validation errors - `ProcessTableFunctionSemanticTests`: broadcast `MapView`, `ValueView`, and eager state with SQL and Table API, row semantic main tables, multiple broadcast tables, `on_time`, `clearAllState()`, as well as runtime errors for `collect()`, timers, and writes from the main path. End-of-input timers verify the final broadcast state content deterministically. - `ProcessTableFunctionRestoreTests`: restore of broadcast state from a savepoint - `BroadcastStateAdaptersTest`: read-only behavior of the state adapters ## 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 (`StateKind`, `StateHint`, `ArgumentTrait`, `StaticArgumentTrait`, `StateTypeStrategy`, as proposed in FLIP-565) - The serializers: no - The runtime per-record code paths (performance sensitive): yes (PTF runner and operators; the code paths of PTFs without broadcast tables remain unchanged) - Anything that affects deployment or recovery: JobManager (and its components), Checkpointing, Kubernetes/Yarn, ZooKeeper: no (PTF operators use operator broadcast state which is checkpointed as usual) - The S3 file system connector: no ## Documentation - Does this pull request introduce a new feature? yes - If yes, how is the feature documented? docs / JavaDocs --- ##### 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) -- 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]
