一、概念

Flink 的广播状态(Broadcast State) 是用于在 流处理(DataStream API)中,将一条控制流或规则流广播到所有并行实例,并在每个算子实例中维护一致的共享状态的一种机制。

二、主要应用场景

规则流 + 数据流,规则需要动态更新,且所有并行实例都要用同一套规则。(具体见七)

  • 动态规则下发
  • 实时风控规则更新
  • 动态维表更新(小维表广播 + 大事件流 join)
  • 配置变更实时生效

三、为什么需要广播变量?

在分布式流处理中:

  • 主数据流(如用户行为流)会被 keyBy 分区
  • 控制流(如规则更新流)需要被 所有并行任务看到

如果直接 keyBy:

  • 规则可能只到达某个分区
  • 其他分区拿不到规则

广播状态的作用:
👉 把规则流广播给所有并行算子实例
👉 每个实例都维护一份本地规则副本


四、核心API结构

1️⃣ 定义广播状态描述符

MapStateDescriptor — 广播状态的描述符,本质是一个 Map 结构:

MapStateDescriptor<String, Rule> ruleStateDescriptor =
    new MapStateDescriptor<>(
        "RulesBroadcastState",
        BasicTypeInfo.STRING_TYPE_INFO,
        TypeInformation.of(Rule.class)
    );

2️⃣ 广播规则流

BroadcastStream — 通过 broadcast() 将普通流转为广播流:

BroadcastStream<Rule> broadcastStream =
    ruleStream.broadcast(ruleStateDescriptor);

3️⃣ 使用方式

KeyedStream + BroadcastStream → KeyedBroadcastProcessFunction

DataStream<String> result = keyedDataStream
    .connect(broadcastRuleStream)
    .process(new KeyedBroadcastProcessFunction<String, Event, Rule, String>() {

        // 处理广播流中的规则数据(只在一个并行实例上执行,但会同步到所有实例)
        @Override
        public void processBroadcastElement(Rule rule, Context ctx, Collector<String> out) throws Exception {
            BroadcastState<String, Rule> broadcastState = ctx.getBroadcastState(ruleStateDescriptor);
            broadcastState.put(rule.getId(), rule);  // 更新规则
        }

        // 处理数据流中的每条数据(可读取广播状态,但不能修改)
        @Override
        public void processElement(Event event, ReadOnlyContext ctx, Collector<String> out) throws Exception {
            ReadOnlyBroadcastState<String, Rule> broadcastState = ctx.getBroadcastState(ruleStateDescriptor);
            Rule rule = broadcastState.get(event.getRuleId());
            if (rule != null && rule.matches(event)) {
                out.collect("matched: " + event);
            }
        }
    });

4️⃣ 关键限制与注意事项

限制 说明
只读约束 processElement 中只能读广播状态,不能写,防止各并行实例状态不一致
广播流并行度 广播流的上游算子并行度必须为 1,否则无法保证所有实例收到相同数据
状态大小 广播状态存在每个 TaskManager 上,数据量不能太大,否则内存压力大
顺序不保证 广播元素到达各并行实例的时间可能不同,业务上要考虑时序问题
Checkpoint 广播状态支持 Checkpoint,但每个并行实例都会保存一份,存储开销是 N 倍

5️⃣ 状态访问权限对比

processBroadcastElement()  →  可读写 BroadcastState
processElement()           →  只读 ReadOnlyBroadcastState

这个设计是刻意的 — 广播状态的修改必须统一在 processBroadcastElement 中完成,保证所有并行实例的状态一致性。


五、BroadcastProcessFunction 详解

有两种常用类型:

类型 说明
BroadcastProcessFunction 非 keyBy 主流
KeyedBroadcastProcessFunction 主流已 keyBy

通常用第二种。

代码示例

public class MyFunction extends KeyedBroadcastProcessFunction<
        String,        // 主流 key 类型
        Event,         // 主流数据类型
        Rule,          // 广播流数据类型
        Result> {      // 输出类型

    // 处理主数据流
    @Override
    public void processElement(
            Event event,
            ReadOnlyContext ctx,
            Collector<Result> out) throws Exception {

        ReadOnlyBroadcastState<String, Rule> rules =
                ctx.getBroadcastState(ruleStateDescriptor);

        Rule rule = rules.get(event.getRuleId());

        if (rule != null) {
            // 根据规则处理数据
        }
    }

    // 处理广播规则流
    @Override
    public void processBroadcastElement(
            Rule rule,
            Context ctx,
            Collector<Result> out) throws Exception {

        BroadcastState<String, Rule> state =
                ctx.getBroadcastState(ruleStateDescriptor);

        state.put(rule.getId(), rule);
    }
}


六、广播状态的核心特点

1️⃣ 每个并行实例都有一份完整副本

  • 并不是共享内存
  • 每个 task 自己维护一份

优点:

  • 本地读取,性能高
  • 不需要跨网络通信

2️⃣ 只允许在广播流中写入

流类型 是否可写广播状态
主流 ❌ 只读
广播流 ✅ 可写

这是为了保证一致性。

3️⃣ 支持 checkpoint

广播状态:

  • 会参与 checkpoint
  • 支持 exactly-once
  • 故障恢复时规则状态也会恢复

4️⃣ 不能 TTL

Broadcast State 不支持 State TTL。

如果规则需要过期:

  • 需要自己实现过期逻辑

    • 规则自带过期时间 + 懒清理(最常用)
    • 规则流发“撤销/删除事件”(推荐,语义最清晰)
    • 广播侧“自建定时器”做周期清理(没有 TTL 时的工程解法)

七、经典应用场景

1️⃣ 实时规则下发(最典型场景)

例如:

  • 实时风控规则
  • 实时营销活动规则
  • 标签圈选规则
  • 动态黑名单
  • AB 实验配置
  • 动态阈值参数

例如在实时营销系统中:

  • 主流:用户行为数据流
  • 广播流:营销规则流(从 MySQL binlog / Kafka 获取)

规则更新后,需要立即影响后续数据计算。

2️⃣ 小表与大流 Join

当:

  • 小表数据量较小(如几万~几十万条)
  • 需要高频匹配
  • 不适合频繁访问外部存储

可以将小表作为广播流,下发到所有 TaskManager 本地状态中。

3️⃣ 维表实时更新

例如:

  • 用户等级规则
  • 风控策略表
  • 配置参数表

相比普通 MapState,BroadcastState 的特点是:

每个并行实例都会完整持有一份规则副本。


八、Broadcast State vs 其他状态

对比项 Broadcast State Keyed State
是否 keyBy 不需要 需要
是否每个实例一份 按 key 分片
写权限 仅广播流 主流
使用场景 规则、配置 用户数据


九、如何保持数据一致性

1️⃣ Checkpoint 机制保证状态一致

Broadcast State 属于:

Flink Managed State

它和普通 Keyed State 一样:

  • 会参与 checkpoint
  • 恢复时可回滚到一致状态
  • 支持 Exactly-Once 语义

只要开启 checkpoint:

env.enableCheckpointing(5000);

就能保证:

  • 主流数据
  • 广播规则
  • 状态更新

在同一一致性点对齐。

2️⃣ 广播流顺序一致性

Flink 保证:

所有并行实例接收到广播流的顺序完全一致。

原因:

  • 广播流是单流复制
  • 每个 SubTask 接收顺序相同
  • 状态更新逻辑一致

前提是:

广播流必须是单分区或保证全局顺序。

3️⃣ 规则更新幂等设计(业务层保证)

Broadcast State 只能保证:

技术一致性

但不能保证:

规则逻辑正确

因此建议:

  • 规则加 version 字段
  • 新规则覆盖旧规则
  • 使用 upsert 模式
  • 禁止删除后再新增

例如:

{
  "rule_id": 1001,
  "version": 5,
  "type": "update"
}

处理逻辑:

  • 只接受 version 更大的规则
  • 避免乱序覆盖

4️⃣ 规则和主流的时序一致性问题

需要注意一个关键问题:

广播规则更新和主流数据到达可能存在时序差异。

可能出现:

  • 规则已更新
  • 部分数据按旧规则处理
  • 部分数据按新规则处理

解决方式:

方法一:规则带生效时间

规则增加:

{
  "rule_id": 1,
  "effective_time": 1700000000
}

在主流中判断:

if (eventTime >= rule.effectiveTime) {
   使用新规则
}

方法二:双流对齐(高级做法)

  • 使用 eventTime
  • 规则流也打 watermark
  • 利用定时器对齐时间线

适合高精度场景。

5️⃣ 故障恢复一致性

Flink 恢复时:

  • 主流回滚
  • 广播状态回滚
  • Kafka offset 回滚

恢复到 checkpoint 对齐点。

因此:

规则和数据处理状态保持一致。


十、Broadcast State 常见坑

❌ 1. 规则太大

广播状态不适合:

  • 百万级大表
  • 高频全量更新

否则:

  • 状态膨胀
  • checkpoint 变慢
  • OOM

❌ 2. 没开启 checkpoint

没有 checkpoint:

  • 无法保证一致性
  • 故障后状态丢失

❌ 3. 广播流多分区乱序

如果规则来自 Kafka 多分区:

  • 可能乱序
  • 建议设置 1 分区

❌ 4. 状态 TTL 误删

Broadcast State 不支持 TTL 自动清理
需要手动管理。


十一、完整的示例

最后再来个完整示例吧,我们实现一下动态规则匹配

public class DynamicRuleExample {

    // 规则 POJO
    public static class Rule {
        public String ruleId;
        public String field;
        public String value;

        public boolean matches(Event event) {
            return value.equals(event.getField(field));
        }
    }

    // 事件 POJO
    public static class Event {
        public String userId;
        public Map<String, String> fields;

        public String getField(String key) {
            return fields.getOrDefault(key, "");
        }
    }

    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        // 开启 Checkpoint
        env.enableCheckpointing(5000);

        // 模拟事件流
        DataStream<Event> eventStream = env.addSource(new EventSource());

        // 模拟规则流(并行度必须为 1)
        DataStream<Rule> ruleStream = env.addSource(new RuleSource()).setParallelism(1);

        // 定义广播状态描述符
        MapStateDescriptor<String, Rule> ruleStateDesc = new MapStateDescriptor<>(
            "rules",
            BasicTypeInfo.STRING_TYPE_INFO,
            TypeInformation.of(new TypeHint<Rule>() {})
        );

        // 广播规则流
        BroadcastStream<Rule> broadcastRules = ruleStream.broadcast(ruleStateDesc);

        // 事件流按 userId 分组后与广播流连接
        DataStream<String> result = eventStream
            .keyBy(e -> e.userId)
            .connect(broadcastRules)
            .process(new KeyedBroadcastProcessFunction<String, Event, Rule, String>() {

                @Override
                public void processBroadcastElement(
                        Rule rule,
                        Context ctx,
                        Collector<String> out) throws Exception {

                    BroadcastState<String, Rule> state = ctx.getBroadcastState(ruleStateDesc);

                    if (rule == null) {
                        // 支持删除规则
                        state.remove(rule.ruleId);
                    } else {
                        state.put(rule.ruleId, rule);
                    }

                    // 可以通过 applyToKeyedState 访问 keyed state(高级用法)
                    // ctx.applyToKeyedState(keyedStateDesc, (key, keyedState) -> { ... });
                }

                @Override
                public void processElement(
                        Event event,
                        ReadOnlyContext ctx,
                        Collector<String> out) throws Exception {

                    ReadOnlyBroadcastState<String, Rule> state = ctx.getBroadcastState(ruleStateDesc);

                    // 遍历所有规则进行匹配
                    for (Map.Entry<String, Rule> entry : state.immutableEntries()) {
                        Rule rule = entry.getValue();
                        if (rule.matches(event)) {
                            out.collect(String.format(
                                "user=%s matched rule=%s", event.userId, rule.ruleId
                            ));
                        }
                    }
                }
            });

        result.print();
        env.execute("Dynamic Rule Matching");
    }
}

在上面的代码中,有一部分TODO的内容:applyToKeyedState
processBroadcastElement 中可以通过 applyToKeyedState 访问所有 key 的 keyed state,适合规则变更时需要清理或重置历史状态的场景

@Override
public void processBroadcastElement(Rule newRule, Context ctx, Collector<String> out) throws Exception {
    ctx.getBroadcastState(ruleStateDesc).put(newRule.ruleId, newRule);

    // 规则更新时,重置所有用户的匹配计数
    ctx.applyToKeyedState(matchCountDesc, (userId, matchCountState) -> {
        matchCountState.update(0);
    });
}


核心思路就一句话:规则/配置走广播流,数据走主流,两者 connect 后在 processBroadcastElement 里统一维护状态,在 processElement 里只读使用。这个读写分离的设计保证了分布式环境下状态的一致性。

Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐