Flink 1.17.2中RichCoFlatMapFunction使用异常问题求助
问题分析与解决方案
核心原因
你遇到的问题本质是基于集合的数据源(fromElements)元素处理顺序不确定,同时你采用按key关联的Connected Streams,导致业务流(streamOfWords)的元素先于对应key的控制流(control)元素到达算子实例,此时blackList状态还未被初始化,过滤逻辑自然失效。
解决方案
针对这个场景,推荐两种方案,分别适用于测试场景和生产场景:
方案1:生产场景 - 使用广播流(Broadcast Stream)
因为黑名单是全局生效的控制规则,不需要按key分区,广播流可以将控制消息发送到所有下游算子实例,确保所有业务流元素都能读取到最新的黑名单状态。
修改后的代码如下:
import org.apache.flink.api.common.state.MapStateDescriptor; import org.apache.flink.streaming.api.datastream.BroadcastStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.functions.co.BroadcastProcessFunction; import org.apache.flink.util.Collector; public class ControlStreamBroadcast { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 定义广播状态描述器,存储黑名单规则 MapStateDescriptor<String, Boolean> blackListDescriptor = new MapStateDescriptor<>("blacklist", String.class, Boolean.class); // 控制流:广播黑名单规则 BroadcastStream<String> controlBroadcast = env .fromElements("DROP", "IGNORE") .broadcast(blackListDescriptor); // 业务流 env.fromElements("Apache", "DROP", "Flink", "IGNORE") .connect(controlBroadcast) .process(new BroadcastProcessFunction<String, String, String>() { @Override public void processElement(String value, ReadOnlyContext ctx, Collector<String> out) throws Exception { // 读取广播状态,判断是否在黑名单中 Boolean isBlacklisted = ctx.getBroadcastState(blackListDescriptor).get(value); if (isBlacklisted == null || !isBlacklisted) { out.collect(value); } } @Override public void processBroadcastElement(String value, Context ctx, Collector<String> out) throws Exception { // 更新广播状态,将黑名单标记为true ctx.getBroadcastState(blackListDescriptor).put(value, true); } }) .print(); env.execute(ControlStreamBroadcast.class.getName()); } }
方案2:测试场景 - 强制控制流先处理
如果只是测试Connected Streams的逻辑,可以通过设置事件时间语义+调整时间戳的方式,让控制流元素先被处理:
import org.apache.flink.api.common.eventtime.WatermarkStrategy; import org.apache.flink.api.common.state.ValueState; import org.apache.flink.api.common.state.ValueStateDescriptor; import org.apache.flink.api.common.typeinfo.Types; import org.apache.flink.configuration.Configuration; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.functions.co.RichCoFlatMapFunction; import org.apache.flink.util.Collector; import org.slf4j.Logger; import org.slf4j.LoggerFactory; public class ControlStreamOrdered { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // 使用事件时间语义 env.setStreamTimeCharacteristic(org.apache.flink.streaming.api.TimeCharacteristic.EventTime); DataStream<String> control = env .fromElements("DROP", "IGNORE") // 给控制流元素分配更早的时间戳(0) .assignTimestampsAndWatermarks(WatermarkStrategy.<String>forMonotonousTimestamps() .withTimestampAssigner((element, recordTimestamp) -> 0L)) .keyBy(x -> x); DataStream<String> streamOfWords = env .fromElements("Apache", "DROP", "Flink", "IGNORE") // 给业务流元素分配更晚的时间戳(1) .assignTimestampsAndWatermarks(WatermarkStrategy.<String>forMonotonousTimestamps() .withTimestampAssigner((element, recordTimestamp) -> 1L)) .keyBy(x -> x); control.connect(streamOfWords) .flatMap(new ControlFunction()) .print(); env.execute(ControlStreamOrdered.class.getName()); } public static class ControlFunction extends RichCoFlatMapFunction<String, String, String> { private final Logger LOGGER = LoggerFactory.getLogger(ControlFunction.class); private ValueState<Boolean> blackList; @Override public void open(Configuration parameters) throws Exception { ValueStateDescriptor<Boolean> desc = new ValueStateDescriptor<>("blackList", Types.BOOLEAN); blackList = getRuntimeContext().getState(desc); } @Override public void flatMap1(String value, Collector<String> out) throws Exception { blackList.update(true); LOGGER.info(String.format("[flatMap1] value=%s, state=%s", value, blackList.value())); } @Override public void flatMap2(String value, Collector<String> out) throws Exception { LOGGER.info(String.format("[flatMap2] value=%s, state=%s", value, blackList.value())); if (blackList.value() == null || !blackList.value()) { out.collect(value); } } } }
说明
- 广播流方案更适合生产场景,因为黑名单通常是全局规则,不需要按key隔离,且能动态更新规则。
- 事件时间调整的方案仅适合测试,生产环境中不建议依赖这种强制顺序的方式,因为真实数据流的时间戳是业务生成的,无法人为控制。
内容的提问来源于stack exchange,提问作者J.Y. Wang
相关产品推荐
相关产品推荐

