Flink Broadcast State使用疑问:本地存储替代及broadcast方法差异
1. 是否需要使用Broadcast State,能否直接存储到算子本地数据结构
- 如果你完全不需要作业容错(不开Checkpoint,作业挂了直接丢弃所有状态重启),可以临时用本地变量
sortedTimestamps存储时间戳,但只要是生产环境需要故障恢复、开了Checkpoint的场景,必须使用Broadcast State,不能用本地变量,原因如下:- 本地变量属于JVM堆内存数据,作业崩溃重启、扩缩容、Flink版本升级重启时,本地变量的数据会直接丢失,重启后需要等下一次MySQL查询触发才能拿到新的时间戳,这期间的业务逻辑会完全失效,甚至出现空指针异常。你当前代码里的
sortedTimestamps是普通成员变量,作业重启后会是null,直接调用sortedTimestamps(0)就会直接抛出空指针,这是典型的本地变量风险。 - 就算你手动在
snapshotState里对本地变量做快照,普通的Operator State默认的分配策略是分片下发,无法保证重启后每个并行实例都拿到完整的全量时间戳列表,而Broadcast State天生设计为所有并行实例持有全量相同的状态,Checkpoint和恢复逻辑Flink已经帮你封装好,不需要额外处理。 - 本地变量无法保证跨并行实例的状态一致性,后续如果调整广播源并行度,很容易出现不同并行实例持有的时间戳不一致的问题,Broadcast State可以保证状态一致性语义。
- 本地变量属于JVM堆内存数据,作业崩溃重启、扩缩容、Flink版本升级重启时,本地变量的数据会直接丢失,重启后需要等下一次MySQL查询触发才能拿到新的时间戳,这期间的业务逻辑会完全失效,甚至出现空指针异常。你当前代码里的
- 对应你代码里的疑问:需要先定义广播状态的描述符(比如
MapStateDescriptor),在构建广播流时传入该描述符,在processBroadcastElement中通过ctx.getBroadcastState(stateDesc)拿到广播状态写入时间戳,processElement中同样从广播状态读取时间戳做判断,不要用本地的sortedTimestamps变量。
2. .broadcast(stateDesc)和.broadcast()的区别
broadcast():普通的广播数据操作,仅将上游的每条数据复制多份,下发给下游所有并行实例,本身不关联任何状态、也不会做持久化,数据下发后就丢弃,适合不需要持久化广播规则、不需要容错的临时广播场景。broadcast(stateDesc):关联了状态描述符的广播操作,会生成专用的BroadcastStream,广播的每条数据会写入到下游所有并行实例的广播状态中,状态会跟随Checkpoint持久化,作业故障重启时自动恢复全量状态到每个并行实例,适合需要持久化配置、规则的生产场景,就是你当前需求匹配的用法。
内容的提问来源于stack exchange,提问作者objectt
相关产品推荐
相关产品推荐

