一、基础转换算子(最常用)

这类算子用于对数据流进行基础的格式转换、过滤、映射,是处理数据的第一步。

1. map:一对一转换

作用:将数据流中的每个元素转换为另一个元素(输入 1 个,输出 1 个)。场景:字段提取、格式转换(如字符串转对象)。

java

运行

// 示例:提取订单金额并转为Double类型
DataStream<Double> amountStream = orderStream
    .map(line -> {
        // 按逗号拆分每行数据
        String[] fields = line.split(",");
        // 提取第3个字段(金额)并转为Double
        return Double.parseDouble(fields[2]);
    });
// 输出:100.0, 200.0, 150.0, 300.0
2. flatMap:一对多转换

作用:将一个元素转换为 0 个、1 个或多个元素(输入 1 个,输出多个)。场景:数据拆分(如一行拆多行)、脏数据过滤。

java

运行

// 示例:拆分订单信息为Tuple2(用户ID, 金额),并过滤支付失败的订单
DataStream<Tuple2<String, Double>> userAmountStream = orderStream
    .flatMap(new FlatMapFunction<String, Tuple2<String, Double>>() {
        @Override
        public void flatMap(String line, Collector<Tuple2<String, Double>> out) throws Exception {
            String[] fields = line.split(",");
            String userId = fields[1];
            double amount = Double.parseDouble(fields[2]);
            String status = fields[3];
            
            // 只保留支付成功的订单
            if ("pay_success".equals(status)) {
                out.collect(Tuple2.of(userId, amount));
            }
        }
    });
// 输出:(user_01,100.0), (user_01,150.0), (user_03,300.0)
3. filter:数据过滤

作用:根据条件筛选出符合要求的元素。场景:脏数据过滤、业务规则筛选(如只保留大额订单)。

java

运行

// 示例:过滤出金额大于200的成功订单
DataStream<String> highAmountStream = orderStream
    .filter(line -> {
        String[] fields = line.split(",");
        double amount = Double.parseDouble(fields[2]);
        String status = fields[3];
        // 条件:支付成功 且 金额>200
        return "pay_success".equals(status) && amount > 200;
    });
// 输出:order_004,user_03,300,pay_success

二、聚合算子(核心统计)

聚合算子需结合keyBy使用(先分组,再聚合),是实时统计的核心。

1. keyBy:数据分组

作用:按指定字段将数据流分组(类似 SQL 的 GROUP BY),是聚合的前提。注意keyBy返回KeyedStream,只能在 KeyedStream 上执行聚合。

java

运行

// 示例:按用户ID分组
KeyedStream<Tuple2<String, Double>, String> keyedStream = userAmountStream
    // 按Tuple2的第一个字段(用户ID)分组
    .keyBy(tuple -> tuple.f0);
2. sum/avg/max/min:基础聚合

作用:对分组后的数据进行求和、平均值、最大值、最小值计算。场景:实时统计用户累计消费、订单最大金额等。

java

运行

// 示例:统计每个用户的累计消费金额
DataStream<Tuple2<String, Double>> sumStream = keyedStream
    // 对Tuple2的第二个字段(金额)求和
    .sum(1);
// 输出:(user_01,100.0) → (user_01,250.0) → (user_03,300.0)

// 示例:统计每个用户的平均消费金额
DataStream<Tuple2<String, Double>> avgStream = keyedStream
    .avg(1);
// 输出:(user_01,100.0) → (user_01,125.0) → (user_03,300.0)
3. reduce:自定义聚合

作用:自定义聚合逻辑(比 sum/avg 更灵活),支持增量聚合。场景:复杂统计(如累计金额 + 订单数)。

java

运行

// 示例:统计每个用户的累计金额和订单数(Tuple3:用户ID, 累计金额, 订单数)
DataStream<Tuple3<String, Double, Integer>> reduceStream = keyedStream
    .reduce((t1, t2) -> {
        // t1:历史聚合结果;t2:新到来的元素
        String userId = t1.f0;
        double totalAmount = t1.f1 + t2.f1; // 累计金额
        int orderCount = 1 + (t1.f2 == null ? 0 : t1.f2); // 订单数
        return Tuple3.of(userId, totalAmount, orderCount);
    }, () -> Tuple3.of("", 0.0, 0)); // 初始值

// 输出:(user_01,100.0,1) → (user_01,250.0,2) → (user_03,300.0,1)

三、窗口算子(实时统计核心)

Flink 是流式计算,窗口用于将无限流切分为有限的 “批次” 进行统计,结合keyBy使用。

1. 滚动窗口(Tumbling Window)

作用:窗口大小固定,无重叠(如每 5 分钟统计一次)。场景:固定周期统计(如每小时用户消费总额)。

java

运行

// 示例:5秒滚动窗口,统计每个用户的消费总额
DataStream<Tuple2<String, Double>> tumblingWindowStream = keyedStream
    // 5秒滚动窗口
    .window(TumblingProcessingTimeWindows.of(Time.seconds(5)))
    // 对金额求和
    .sum(1);
2. 滑动窗口(Sliding Window)

作用:窗口大小固定,有重叠(如每 2 分钟统计最近 5 分钟的数据)。场景:高频统计(如实时监控,每 10 秒统计最近 1 分钟的订单量)。

java

运行

// 示例:滑动窗口(窗口5秒,滑动2秒)
DataStream<Tuple2<String, Double>> slidingWindowStream = keyedStream
    .window(SlidingProcessingTimeWindows.of(Time.seconds(5), Time.seconds(2)))
    .sum(1);
3. 会话窗口(Session Window)

作用:按用户会话划分窗口(如用户连续操作 30 秒内为一个会话)。场景:用户行为分析(如统计用户一次会话内的消费金额)。

java

运行

// 示例:会话窗口(超时时间3秒,无操作3秒则窗口关闭)
DataStream<Tuple2<String, Double>> sessionWindowStream = keyedStream
    .window(ProcessingTimeSessionWindows.withGap(Time.seconds(3)))
    .sum(1);

四、连接 / 拆分算子

1. union:合并同类型数据流

作用:将多个同类型的数据流合并为一个(字段结构必须完全一致)。场景:合并多来源的同类型数据(如多个省份的订单流)。

java

运行

// 模拟第二个订单流
DataStream<Tuple2<String, Double>> orderStream2 = env.fromElements(
    Tuple2.of("user_02", 180.0),
    Tuple2.of("user_03", 50.0)
);
// 合并两个数据流
DataStream<Tuple2<String, Double>> unionStream = userAmountStream.union(orderStream2);
2. connect:连接不同类型数据流

作用:连接两个不同类型的数据流,支持自定义处理逻辑。场景:关联补充数据(如订单流 + 用户信息流)。

java

运行

// 模拟用户信息流(用户ID, 用户名)
DataStream<Tuple2<String, String>> userInfoStream = env.fromElements(
    Tuple2.of("user_01", "张三"),
    Tuple2.of("user_02", "李四")
);
// 连接订单流和用户信息流
ConnectedStreams<Tuple2<String, Double>, Tuple2<String, String>> connectedStreams = 
    userAmountStream.connect(userInfoStream)
    // 按用户ID分组(两个流的关联键)
    .keyBy(t1 -> t1.f0, t2 -> t2.f0);

// 处理连接后的数据流(关联用户名和消费金额)
DataStream<String> resultStream = connectedStreams
    .map(new CoMapFunction<Tuple2<String, Double>, Tuple2<String, String>, String>() {
        @Override
        public String map1(Tuple2<String, Double> order) throws Exception {
            // 处理订单流(暂时无用户名,先返回默认值)
            return order.f0 + ",未知用户," + order.f1;
        }

        @Override
        public String map2(Tuple2<String, String> user) throws Exception {
            // 处理用户信息流(暂时无金额,先返回默认值)
            return user.f0 + "," + user.f1 + ",0.0";
        }
    });

总结

  1. 基础转换map(一对一)、flatMap(一对多)、filter(过滤)是数据预处理的核心,几乎所有 Flink 任务都会用到;
  2. 聚合统计:先keyBy分组,再用sum/avg/reduce做聚合,是实时统计的基础;
  3. 窗口核心:滚动窗口(无重叠)、滑动窗口(有重叠)、会话窗口(按会话)是流式统计的关键,需结合业务场景选择。

更多推荐