snuyanzin commented on code in PR #29419:
URL: https://github.com/apache/flink/pull/29419#discussion_r4210195672
##########
flink-table/flink-table-common/src/main/java/org/apache/flink/table/functions/ProcessTableFunction.java:
##########
@@ -378,6 +389,63 @@
* }
* }</pre>
*
+ * <h2>Broadcast State</h2>
+ *
+ * <p>Broadcast state enables patterns such as a rule engine where rules are
dynamically updated at
+ * runtime, or dynamic configuration that influences the processing of the
main table(s).
+ *
+ * <p>Broadcast state is tightly coupled to a broadcast semantic table. A
table argument declared
+ * with {@link ArgumentTrait#BROADCAST_SEMANTIC_TABLE} acts as a "side" or
"control" input. Every
+ * row of a broadcast table is sent to all virtual processors, regardless of
any PARTITION BY
+ * clause. A PTF can store this broadcast information in state entries
declared with
+ * {@code @StateHint(StateKind.BROADCAST)}. This broadcast state is shared
across all sets and can
+ * be read when processing rows from the main table(s).
+ *
+ * <p>The following rules apply to broadcast tables and broadcast state:
+ *
+ * <ul>
+ * <li>At least one table argument with row or set semantics must be
declared next to broadcast
+ * tables. Multiple broadcast tables and multiple broadcast state
entries are supported.
+ * <li>Broadcast state can only be modified while processing a broadcast
row. This includes
+ * clearing broadcast state. When processing rows of the main table(s),
broadcast state is
+ * read-only.
+ * <li>While processing a broadcast row, there is no key context. State
entries that are scoped to
+ * a set are passed as null, results can not be emitted via {@code
collect()}, and timers
+ * cannot be registered or cleared. {@link Context#clearAllState()} only
clears broadcast
+ * state in this case.
+ * <li>It is the responsibility of the PTF implementer to maintain identical
broadcast state
+ * across all virtual processors, i.e. broadcast state should only be
updated
+ * deterministically based on the broadcast rows.
+ * </ul>
+ *
+ * <p>Note: The system decides which input row is streamed through the virtual
processor next. A
+ * change to a broadcast state entry has no effect on rows of the main
table(s) that have been
+ * processed before.
+ *
+ * <pre>{@code
+ * // Function that filters sentences using a dynamically updated list of bad
words
+ * class RuleFunction extends ProcessTableFunction<String> {
+ * public void eval(
+ * @StateHint(StateKind.BROADCAST) MapView<String, Boolean> badWords,
+ * @ArgumentHint(ROW_SEMANTIC_TABLE) Row data,
+ * @ArgumentHint(BROADCAST_SEMANTIC_TABLE) Row rules) throws Exception {
+ * // Write access to broadcast state for the broadcast table
+ * if (rules != null) {
+ * badWords.put(rules.getFieldAs("word"), true);
+ * return;
+ * }
+ * // Read access to broadcast state for the main table
+ * String sentence = data.getFieldAs("sentence");
+ * for (String word : sentence.split(" ")) {
Review Comment:
the example in docs also uses `toLowerCase()`
is it intentional?
https://github.com/apache/flink/pull/29419/changes#diff-3454533d1ebf67d5630b4f39257c7ec9b5b172f12f7191cb26a35644a66ce55bR1434
--
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]