如何基于控制流实时过滤Flink DataStream?顺序问题求解
我正在学习Apache Flink,遇到一个看似简单的场景却产生了疑惑,怀疑自己对Flink的核心用法存在根本性误解。该场景参考Flink官方文档的简单示例:我们有需要过滤掉的静态控制元素,同时还有持续流式传输的动态词。相关代码如下:
public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); DataStream<String> control = env .fromElements("DROP", "IGNORE") .keyBy(x -> x); DataStream<String> streamOfWords = env .fromElements("Apache", "DROP", "Flink", "IGNORE") .keyBy(x -> x); control .connect(streamOfWords) .flatMap(new ControlFunction()) .print(); env.execute(); }
public static class ControlFunction extends RichCoFlatMapFunction<String, String, String> { private ValueState<Boolean> blocked; @Override public void open(Configuration config) { blocked = getRuntimeContext() .getState(new ValueStateDescriptor<>("blocked", Boolean.class)); } @Override public void flatMap1(String control_value, Collector<String> out) throws Exception { blocked.update(Boolean.TRUE); } @Override public void flatMap2(String data_value, Collector<String> out) throws Exception { if (blocked.value() == null) { out.collect(data_value); } } }
文档中提到:
需要注意的是,你无法控制flatMap1和flatMap2回调函数的调用顺序。……在时序或顺序至关重要的场景中,你可能需要将事件缓存在Flink的托管状态中,直到应用准备好处理它们。
我的疑问是:若无法控制消费顺序,该如何实时过滤这些词?我是否误解了Apache Flink的使用场景与方式?是否应该采用非Flink API方式存储控制词,还是存在Flink原生的合适解决方案?
你没有误解Flink的核心用法,这个场景完全可以用Flink原生方案解决,不需要依赖外部存储。问题的核心在于处理流与控制流的时序对齐以及状态的正确使用,以下是具体方案:
1. 静态控制词:直接初始化状态
如果控制词是固定不变的,无需依赖控制流传输,直接在算子初始化阶段完成状态设置,从根源上避免时序问题。调整ControlFunction的open方法:
@Override public void open(Configuration config) { blocked = getRuntimeContext().getState(new ValueStateDescriptor<>("blocked", Boolean.class)); // 获取当前算子的key(对应每个词本身) String currentKey = getRuntimeContext().getCurrentKey(); // 静态控制词列表 String[] staticControlWords = {"DROP", "IGNORE"}; for (String word : staticControlWords) { if (word.equals(currentKey)) { blocked.update(true); break; } } }
这样数据流到达时,状态已经就绪,直接判断即可,无需等待控制流事件。
2. 动态控制词:使用广播流
如果控制词会动态更新(比如后续新增过滤词),广播流是更通用的方案:将控制流广播到所有下游算子,每个算子维护一份全局过滤列表,数据流到达时直接对比状态。
示例代码:
// 定义广播状态描述器 MapStateDescriptor<String, Boolean> broadcastStateDesc = new MapStateDescriptor<>( "blockedWords", String.class, Boolean.class ); // 控制流转为广播流 BroadcastStream<String> controlBroadcast = env .fromElements("DROP", "IGNORE") .broadcast(broadcastStateDesc); // 数据流连接广播流处理 DataStream<String> filteredStream = env .fromElements("Apache", "DROP", "Flink", "IGNORE") .connect(controlBroadcast) .process(new BroadcastProcessFunction<String, String, String>() { @Override public void processElement(String word, ReadOnlyContext ctx, Collector<String> out) throws Exception { // 读取广播状态判断是否过滤 ReadOnlyBroadcastState<String, Boolean> broadcastState = ctx.getBroadcastState(broadcastStateDesc); if (!broadcastState.contains(word)) { out.collect(word); } } @Override public void processBroadcastElement(String controlWord, Context ctx, Collector<String> out) throws Exception { // 更新全局广播状态 BroadcastState<String, Boolean> broadcastState = ctx.getBroadcastState(broadcastStateDesc); broadcastState.put(controlWord, true); } }); filteredStream.print(); env.execute();
广播流的优势在于控制词更新时,所有算子能实时同步状态,无需依赖keyBy,适合全局过滤场景。
3. 时序敏感场景:缓存未就绪事件
如果必须保留原有的keyBy+connect模式,且控制流是动态的,可以在数据流事件到达但控制状态未就绪时,先缓存事件,等控制流事件到达后再处理:
调整后的ControlFunction:
public static class ControlFunction extends RichCoFlatMapFunction<String, String, String> { private ValueState<Boolean> blocked; private ListState<String> cachedData; @Override public void open(Configuration config) { blocked = getRuntimeContext().getState(new ValueStateDescriptor<>("blocked", Boolean.class)); cachedData = getRuntimeContext().getListState(new ListStateDescriptor<>("cachedData", String.class)); } @Override public void flatMap1(String control_value, Collector<String> out) throws Exception { blocked.update(Boolean.TRUE); // 处理缓存的数据流事件 for (String data : cachedData.get()) { if (blocked.value() == null) { out.collect(data); } } cachedData.clear(); } @Override public void flatMap2(String data_value, Collector<String> out) throws Exception { if (blocked.value() == null) { // 控制状态未就绪,缓存事件 cachedData.add(data_value); } else { // 状态就绪,直接判断 if (!blocked.value()) { out.collect(data_value); } } } }
这种方式确保只有控制状态就绪后,才会处理对应的数据流事件,避免时序问题导致的错误过滤。
内容的提问来源于stack exchange,提问作者yokus

