You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.26 23:55:23