尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Flink Broadcast State 模式:用广播流实现动态规则下的实时模式匹配

Flink Broadcast State 模式:用广播流实现动态规则下的实时模式匹配 Flink Broadcast State 模式用广播流实现动态规则下的实时模式匹配【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink本篇技术指南以 Apache Flink 官方文档《Broadcast State 模式》为主体完整讲解 broadcast state广播状态的编程模型与实战用法如何把一个规则流广播到所有下游算子实例、如何用MapStateDescriptor描述广播状态、如何通过connect()BroadcastProcessFunction/KeyedBroadcastProcessFunction实现动态规则 实时匹配并深入剖析其底层约束无跨 task 通讯、checkpoint 写放大、仅内存存储等。读完本文你将能够独立编写一个规则可实时热更新的流式模式识别作业并清楚哪些场景适合、哪些场景严禁使用 broadcast state。有关有状态流处理的基础概念请参阅 有状态流处理。提供的 API这里我们用一个贯穿全文的例子来展现 broadcast state 提供的接口。假设存在一个序列序列中的元素是具有不同颜色与形状的图形我们希望在序列里相同颜色的图形中寻找满足一定顺序模式的图形对比如在红色的图形里有一个长方形跟着一个三角形。同时我们希望寻找的模式也会随着时间而改变——这正是 broadcast state 的核心价值模式作为低吞吐的规则流广播到所有实例图形作为高吞吐的数据流在每个实例本地完成匹配规则更新无需重启作业。在这个例子中我们定义两个流图形Item流元素具有颜色Color和形状Shape两个属性规则Rule流代表希望寻找的模式例如first 长方形, second 三角形。第一步按颜色分区图形流在图形流中我们需要首先使用颜色将流进行分区keyBy这能确保相同颜色的图形会流转到相同的物理机上从而把跨实例的状态协作转化为单实例内的局部匹配。// 将图形使用颜色进行划分 KeyedStreamItem, Color colorPartitionedStream itemStream .keyBy(new KeySelectorItem, Color(){...});# 将图形使用颜色进行划分 color_partitioned_stream item_stream.key_by(lambda item: ...)第二步广播规则流并声明广播状态对于规则流它应该被广播到所有的下游 task 中下游 task 应当存储这些规则并根据它寻找满足规则的图形对。下面这段代码会完成i) 将规则广播给所有下游 task ii) 使用MapStateDescriptor来描述并创建 broadcast state 在下游的存储结构。// 一个 map descriptor它描述了用于存储规则名称与规则本身的 map 存储结构 MapStateDescriptorString, Rule ruleStateDescriptor new MapStateDescriptor( RulesBroadcastState, BasicTypeInfo.STRING_TYPE_INFO, TypeInformation.of(new TypeHintRule() {})); // 广播流广播规则并且创建 broadcast state BroadcastStreamRule ruleBroadcastStream ruleStream .broadcast(ruleStateDescriptor);# 一个 map descriptor它描述了用于存储规则名称与规则本身的 map 存储结构 rule_state_descriptor MapStateDescriptor(RuleBroadcastState, Types.STRING(), Types.PICKLED_BYTE_ARRAY()) # 广播流广播规则并且创建 broadcast state rule_broadcast_stream rule_stream.broadcast(rule_state_descriptor)从实现上看DataStream#broadcast(MapStateDescriptor...)会先把流的连接方式设置为BroadcastPartitioner每个元素被发送到下游的每一个并行实例再包装成BroadcastStream见 DataStream.java。因此MapStateDescriptor的名字如RulesBroadcastState就是 broadcast state 在运行时注册表中的唯一标识getBroadcastState()时必须使用同一个 descriptor 才能取到同一份状态。注意 descriptor 需要同时给出 key 与 value 的TypeInformation因为 broadcast state 本质上是分布式环境下的一个 MapK→V序列化/反序列化完全依赖这些类型信息MapStateDescriptor的完整定义见 MapStateDescriptor.java。第三步连接两个流并编写匹配逻辑最终为了使用规则来筛选图形序列我们需要将两个流关联起来完成我们的模式识别逻辑。为了关联一个非广播流keyed 或者 non-keyed与一个广播流BroadcastStream我们可以调用非广播流的方法connect()并将BroadcastStream当做参数传入。该方法的实现见 DataStream.java返回BroadcastConnectedStream其类型方法process()接收一个特殊的CoProcessFunction来书写我们的模式识别逻辑。具体传入process()的是哪个类型取决于非广播流的类型如果流是一个keyed流那就是KeyedBroadcastProcessFunction类型如果流是一个non-keyed流那就是BroadcastProcessFunction类型。注意connect()方法需要由非广播流来进行调用BroadcastStream作为参数传入。反过来调用在 API 层面是不成立的。在我们的例子中图形流是一个 keyed stream所以我们书写的代码如下DataStreamString output colorPartitionedStream .connect(ruleBroadcastStream) .process( // KeyedBroadcastProcessFunction 中的类型参数表示 // 1. key stream 中的 key 类型 // 2. 非广播流中的元素类型 // 3. 广播流中的元素类型 // 4. 结果的类型在这里是 string new KeyedBroadcastProcessFunctionColor, Item, Rule, String() { // 模式匹配逻辑 } );class MyKeyedBroadcastProcessFunction(KeyedBroadcastProcessFunction): # 模式匹配逻辑 ... output color_partitioned_stream \ .connect(rule_broadcast_stream) \ .process(MyKeyedBroadcastProcessFunction())BroadcastProcessFunction 和 KeyedBroadcastProcessFunction在传入的BroadcastProcessFunction或KeyedBroadcastProcessFunction中我们需要实现两个方法processBroadcastElement()方法负责处理广播流中的元素processElement()负责处理非广播流中的元素。两个子类型定义如下public abstract class BroadcastProcessFunctionIN1, IN2, OUT extends BaseBroadcastProcessFunction { public abstract void processElement(IN1 value, ReadOnlyContext ctx, CollectorOUT out) throws Exception; public abstract void processBroadcastElement(IN2 value, Context ctx, CollectorOUT out) throws Exception; }public abstract class KeyedBroadcastProcessFunctionKS, IN1, IN2, OUT { public abstract void processElement(IN1 value, ReadOnlyContext ctx, CollectorOUT out) throws Exception; public abstract void processBroadcastElement(IN2 value, Context ctx, CollectorOUT out) throws Exception; public void onTimer(long timestamp, OnTimerContext ctx, CollectorOUT out) throws Exception; }class BroadcastProcessFunction(BaseBroadcastProcessFunction, Generic[IN1, IN2, OUT]): abstractmethod def process_element(value: IN1, ctx: ReadOnlyContext): pass abstractmethod def process_broadcast_element(value: IN2, ctx: Context): passclass KeyedBroadcastProcessFunction(BaseBrodcastProcessFunction, Generic[KEY, IN1, IN2, OUT]): abstractmethod def process_element(value: IN1, ctx: ReadOnlyContext): pass abstractmethod def process_broadcast_element(value: IN2, ctx: Context): pass def on_timer(timestamp: int, ctx: OnTimerContext): pass源码层面两个抽象类均继承自BaseBroadcastProcessFunction见 BaseBroadcastProcessFunction.java并在其内部定义了三种 ContextBaseContext提供timestamp()、output(OutputTag, value)、currentProcessingTime()、currentWatermark()等基础能力Context广播侧在BaseContext之上提供getBroadcastState(MapStateDescriptor)返回可读写的BroadcastState见 BroadcastState.java接口含put/putAll/remove/iterator/entries等写与遍历方法ReadOnlyContext非广播侧getBroadcastState()返回ReadOnlyBroadcastState只读视图见 ReadOnlyBroadcastState.java仅含get/contains/immutableEntries等只读方法。这两个类及完整的方法签名分别定义于 BroadcastProcessFunction.java 与 KeyedBroadcastProcessFunction.java类型参数、Context 能力与上文的摘要一一对应可作为实现时的权威参考。需要注意的是processBroadcastElement()负责处理广播流的元素而processElement()负责处理另一个流的元素。两个方法的第二个参数Context不同但均有以下方法得到广播流的存储状态ctx.getBroadcastState(MapStateDescriptorK, V stateDescriptor)查询元素的时间戳ctx.timestamp()查询目前的 Watermarkctx.currentWatermark()目前的处理时间processing timectx.currentProcessingTime()产生旁路输出ctx.output(OutputTagX outputTag, X value)1. 得到广播流的存储状态ctx.get_broadcast_state(stateDescriptor: MapStateDescriptor) 2. 查询元素的时间戳ctx.timestamp() 3. 查询目前的Watermarkctx.current_watermark() 4. 目前的处理时间(processing time)ctx.current_processing_time() 5. 产生旁路输出yield output_tag, value在getBroadcastState()方法中传入的stateDescriptor应该与调用.broadcast(ruleStateDescriptor)的参数相同。为什么广播侧可写、非广播侧只读这两个方法的区别在于对 broadcast state 的访问权限不同处理广播流元素这端是具有读写权限的而对于处理非广播流元素这端是只读的。这样做的原因是Flink 中不存在跨 task 通讯。为了保证 broadcast state 在所有的并发实例中一致我们在处理广播流元素的时候给予写权限——广播流中的每个元素在所有 task 中均可见并且要求对这些元素处理是一致的那么最终所有 task 得到的 broadcast state 就是一致的。注意processBroadcastElement()的实现必须在所有的并发实例中具有确定性的结果。例如不能在其中使用Math.random()或读取系统时钟作为写入 broadcast state 的值否则各实例状态会发散恢复或扩容时可能产生不可预测的结果。KeyedBroadcastProcessFunction 的独有能力同时KeyedBroadcastProcessFunction在 Keyed Stream 上工作所以它提供了一些BroadcastProcessFunction没有的功能processElement()的参数ReadOnlyContext提供了方法能够访问 Flink 的定时器服务timerService()见 KeyedBroadcastProcessFunction.java可以注册事件定时器event-time timer或者处理时间的定时器processing-time timer。当定时器触发时会调用onTimer()方法并提供OnTimerContext——它具备ReadOnlyContext的全部功能并且额外提供查询当前触发的是一个事件还是处理时间的定时器timeDomain()查询定时器关联的 keygetCurrentKey()。processBroadcastElement()方法中的参数Context会提供方法applyToKeyedState(StateDescriptorS, VS stateDescriptor, KeyedStateFunctionKS, S function)。这个方法使用一个KeyedStateFunction能够对stateDescriptor对应的 state 中所有 key 的存储状态进行某些操作。目前 PyFlink 不支持apply_to_keyed_state。注意注册一个定时器只能在KeyedBroadcastProcessFunction的processElement()方法中进行。在processBroadcastElement()方法中不能注册定时器因为广播的元素中并没有关联的 key。完整的模式匹配实现回到我们当前的例子中KeyedBroadcastProcessFunction应该实现如下new KeyedBroadcastProcessFunctionColor, Item, Rule, String() { // 存储部分匹配的结果即匹配了一个元素正在等待第二个元素 // 我们用一个数组来存储因为同时可能有很多第一个元素正在等待 private final MapStateDescriptorString, ListItem mapStateDesc new MapStateDescriptor( items, BasicTypeInfo.STRING_TYPE_INFO, new ListTypeInfo(Item.class)); // 与之前的 ruleStateDescriptor 相同 private final MapStateDescriptorString, Rule ruleStateDescriptor new MapStateDescriptor( RulesBroadcastState, BasicTypeInfo.STRING_TYPE_INFO, TypeInformation.of(new TypeHintRule() {})); Override public void processBroadcastElement(Rule value, Context ctx, CollectorString out) throws Exception { ctx.getBroadcastState(ruleStateDescriptor).put(value.name, value); } Override public void processElement(Item value, ReadOnlyContext ctx, CollectorString out) throws Exception { final MapStateString, ListItem state getRuntimeContext().getMapState(mapStateDesc); final Shape shape value.getShape(); for (Map.EntryString, Rule entry : ctx.getBroadcastState(ruleStateDescriptor).immutableEntries()) { final String ruleName entry.getKey(); final Rule rule entry.getValue(); ListItem stored state.get(ruleName); if (stored null) { stored new ArrayList(); } if (shape rule.second !stored.isEmpty()) { for (Item i : stored) { out.collect(MATCH: i - value); } stored.clear(); } // 不需要额外的 else{} 段来考虑 rule.first rule.second 的情况 if (shape.equals(rule.first)) { stored.add(value); } if (stored.isEmpty()) { state.remove(ruleName); } else { state.put(ruleName, stored); } } } }class MyKeyedBroadcastProcessFunction(KeyedBroadcastProcessFunction): def __init__(self): self._map_state_desc MapStateDescriptor(item, Types.STRING(), Types.LIST(Types.PICKLED_BYTE_ARRAY())) self._rule_state_desc MapStateDescriptor(RulesBroadcastState, Types.STRING(), Types.PICKLED_BYTE_ARRAY()) self._map_state None def open(self, ctx: RuntimeContext): self._map_state ctx.get_map_state(self._map_state_desc) def process_broadcast_element(value: Rule, ctx: KeyedBroadcastProcessFunction.Context): ctx.get_broadcast_state(self._rule_state_desc).put(value.name, value) def process_element(value: Item, ctx: KeyedBroadcastProcessFunction.ReadOnlyContext): shape value.get_shape() for rule_name, rule in ctx.get_broadcast_state(self._rule_state_desc).items(): stored self._map_state.get(rule_name) if stored is None: stored [] if shape rule.second and len(stored) 0: for i in stored: yield MATCH: {} - {}.format(i, value) stored [] if shape rule.first: stored.append(value) if len(stored) 0: self._map_state.remove(rule_name) else: self._map_state.put(rule_name, stored)这段代码的匹配逻辑值得逐段解读两份状态各司其职itemskeyedMapState以规则名为 key值是等待配对的元素列表保存每个 key颜色分区下的部分匹配结果RulesBroadcastStatebroadcast state保存当前生效的全部规则广播侧只做一件事把新规则put进 broadcast state。规则的新增、删除、更新都通过这一入口完成实现规则热更新非广播侧遍历全部规则对每个到达的Item用immutableEntries()遍历当前所有规则只读视图避免并发修改问题若当前形状等于某规则的second且已有等待中的第一个元素则输出MATCH并清空已匹配队列若形状等于first则将其加入等待队列空队列及时清理stored为空时state.remove(ruleName)避免 keyed state 无限膨胀。该逻辑在官方测试中也有直接对应CoBroadcastWithKeyedOperatorTest覆盖了广播侧写入 state、非广播侧读取校验、applyToKeyedState对全部 key 批量操作、广播侧/非广播侧output()旁路输出、以及基于timerService()注册 event-time 定时器等场景见 CoBroadcastWithKeyedOperatorTest.java非 keyed 版本对应的测试见 CoBroadcastWithNonKeyedOperatorTest.java。需要为函数编写单元测试时还可以借助 ProcessFunctionTestHarnesses.java 中的forBroadcastProcessFunction系列方法快速搭建双输入测试环境。重要注意事项这里有一些 broadcast state 的重要注意事项在使用它时需要时刻清楚没有跨 task 通讯如上所述这就是为什么只有在(Keyed)-BroadcastProcessFunction中处理广播流元素的方法里可以更改 broadcast state 的内容。同时用户需要保证所有 task 对于 broadcast state 的处理方式是一致的否则会造成不同 task 读取 broadcast state 时内容不一致的情况最终导致结果不一致。这一点同时被写进了 BroadcastState.java 的接口契约注释中the user has to guarantee that all task instances store the same elements in this type of state。broadcast state 在不同的 task 的事件顺序可能是不同的虽然广播流中元素的广播过程能够保证所有的下游 task 全部能够收到但在不同 task 中元素的到达顺序可能不同。所以 broadcast state 的更新不能依赖于流中元素到达的顺序。典型反例规则流里先发规则 A再发规则 B你不能假设所有实例都按 A→B 的顺序应用——正确做法是让每条规则自身携带生效语义如版本号、时间戳、覆盖式 key使状态收敛不依赖到达次序。所有的 task 均会对 broadcast state 进行 checkpoint虽然所有 task 中的 broadcast state 是一致的但当 checkpoint 来临时所有 task 均会对 broadcast state 做 checkpoint。这个设计是为了防止在作业恢复后读文件造成的文件热点。当然这种方式会造成 checkpoint 一定程度的写放大放大倍数为 p并行度。Flink 会保证在恢复状态/改变并发的时候数据没有重复且没有缺失。在作业恢复时如果与之前具有相同或更小的并发度所有的 task 读取之前已经 checkpoint 过的 state在增大并发的情况下task 会读取本身的 state多出来的并发p_new-p_old会使用轮询调度round-robin算法读取之前 task 的 state。这一恢复语义同样记录在 BroadcastState.java 的注释中也是广播侧写入必须确定性的深层原因——若各实例状态不一致恢复或扩容时的分区再分配将是不可预测的。不使用 RocksDB state backendbroadcast state 在运行时保存在内存中需要保证内存充足。这一特性同样适用于所有其他 Operator State。因此规则的规模必须控制在可接受的内存预算内例如几千到几万条级别的规则通常没有问题具体取决于单条规则的大小与并行度不要把 broadcast state 当作存放海量字典数据的存储。适用场景小结综合上述 API 与约束broadcast state 最典型的适用场景包括动态规则/配置实时下发规则变化频繁、希望免重启生效如反欺诈、风控规则、动态阈值低吞吐控制流 高吞吐数据流广播侧吞吐远低于数据侧如黑白名单更新流 用户行为流跨实例共享一份只读/准实时视图所有并行实例需要看到同一份数据如全局特征、模型版本。反之如果规则数据量巨大、更新极为频繁、或对一致性要求苛刻到需要依赖到达顺序则应考虑其他方案如维表 join 缓存刷新、外部配置中心轮询等。结语Broadcast State 是 Flink 状态编程中少有的跨实例共享状态方案通过broadcast(MapStateDescriptor...)把规则流变成每个算子实例都能接收并写入的广播状态再配合connect()与KeyedBroadcastProcessFunction或BroadcastProcessFunction完成低延迟的动态模式匹配。它的使用边界同样清晰——无跨 task 通讯决定了广播侧可写、数据侧只读的权限模型checkpoint 全量落盘带来 p 倍写放大纯内存存储要求使用者精打细算规则规模。把握住这些要点你就能在动态规则类场景中稳定、可控地使用这一模式。{{ top }}【免费下载链接】flink项目地址: https://gitcode.com/gh_mirrors/fli/flink创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表