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]

Reply via email to