本文详细介绍 Flink 中各种 Stream 类型之间的转换关系,帮助你理解数据流的处理流程。

一、概述

在 Flink 中,数据流会经过不同的转换,变成不同类型的 Stream。理解这些类型之间的转换关系,是写好 Flink 程序的基础。

核心流程

Source → DataStream → KeyedStream → WindowedStream → 聚合结果 → Sink

二、Stream 类型一览

Stream 类型说明如何得到
DataStream基础数据流Source 或转换操作
KeyedStream按 Key 分组的流dataStream.keyBy()
WindowedStream分组窗口流keyedStream.window()
AllWindowedStream全局窗口流dataStream.windowAll()
ConnectedStreams连接的两个流stream1.connect(stream2)
BroadcastConnectedStream广播连接流keyedStream.connect(broadcastStream)
JoinedStreamsJOIN 关联流stream1.join(stream2)
CoGroupedStreams协同分组流stream1.coGroup(stream2)

三、详细转换说明

1. DataStream(基础数据流)

这是最基础的流,所有数据处理的起点。

如何得到

// 方式1:从 Source 创建
DataStream<String> stream = env.addSource(new FlinkKafkaConsumer<>(...));

// 方式2:从集合创建
DataStream<Integer> stream = env.fromElements(1, 2, 3, 4, 5);

// 方式3:从其他 Stream 转换回来
DataStream<Order> stream = keyedStream.map(x -> x);

可以转换成

  • KeyedStream → 通过 keyBy()
  • ConnectedStreams → 通过 connect()
  • AllWindowedStream → 通过 windowAll()
  • JoinedStreams → 通过 join()
  • CoGroupedStreams → 通过 coGroup()
  • DataStreamSink → 通过 addSink()

2. KeyedStream(按 Key 分组的流)

按某个 Key 分组后的流,相同 Key 的数据会被发送到同一个 Task 处理。

如何得到

// 按用户ID分组
KeyedStream<Order, String> keyedStream = orderStream.keyBy(order -> order.getUserId());

// 按多个字段分组
KeyedStream<Order, Tuple2<String, String>> keyedStream = orderStream
    .keyBy(order -> Tuple2.of(order.getUserId(), order.getProductId()));

可以转换成

  • WindowedStream → 通过 window() / timeWindow() / countWindow()
  • DataStream → 通过 map() / flatMap() / process()
  • BroadcastConnectedStream → 通过 connect(broadcastStream)
  • CoGroupedStreams → 通过 coGroup()

代码示例

DataStream<Order> orderStream = env.addSource(...);

// DataStream → KeyedStream
KeyedStream<Order, String> keyedStream = orderStream
    .keyBy(order -> order.getUserId());

// KeyedStream → WindowedStream → DataStream
DataStream<OrderStats> result = keyedStream
    .timeWindow(Time.minutes(5))    // KeyedStream → WindowedStream
    .sum("amount");                  // WindowedStream → DataStream

result.addSink(...);

3. WindowedStream(分组窗口流)

对 KeyedStream 开窗后的流,用于分组聚合统计。

如何得到

// 滚动时间窗口
WindowedStream<Order, String, TimeWindow> windowedStream = keyedStream
    .timeWindow(Time.minutes(5));

// 滑动时间窗口
WindowedStream<Order, String, TimeWindow> windowedStream = keyedStream
    .timeWindow(Time.minutes(10), Time.minutes(5));  // 窗口10分钟,滑动5分钟

// 计数窗口
WindowedStream<Order, String, GlobalWindow> windowedStream = keyedStream
    .countWindow(100);  // 每100条数据一个窗口

// 自定义窗口
WindowedStream<Order, String, TimeWindow> windowedStream = keyedStream
    .window(TumblingEventTimeWindows.of(Time.hours(1)));

转换回 DataStream

// 方式1:内置聚合函数
DataStream<Order> result = windowedStream.sum("amount");
DataStream<Order> result = windowedStream.min("amount");
DataStream<Order> result = windowedStream.max("amount");

// 方式2:reduce
DataStream<Order> result = windowedStream.reduce((o1, o2) -> {
    o1.setAmount(o1.getAmount() + o2.getAmount());
    return o1;
});

// 方式3:aggregate(推荐,最灵活)
DataStream<OrderStats> result = windowedStream.aggregate(
    new AggregateFunction<Order, OrderAcc, OrderStats>() {
        @Override
        public OrderAcc createAccumulator() { return new OrderAcc(); }
        
        @Override
        public OrderAcc add(Order order, OrderAcc acc) {
            acc.count++;
            acc.totalAmount += order.getAmount();
            return acc;
        }
        
        @Override
        public OrderStats getResult(OrderAcc acc) {
            return new OrderStats(acc.count, acc.totalAmount);
        }
        
        @Override
        public OrderAcc merge(OrderAcc a, OrderAcc b) {
            a.count += b.count;
            a.totalAmount += b.totalAmount;
            return a;
        }
    }
);

// 方式4:apply(可以拿到窗口信息)
DataStream<String> result = windowedStream.apply(
    new WindowFunction<Order, String, String, TimeWindow>() {
        @Override
        public void apply(String key, TimeWindow window, 
                         Iterable<Order> input, Collector<String> out) {
            int count = 0;
            for (Order order : input) {
                count++;
            }
            out.collect("窗口: " + window + ", Key: " + key + ", 数量: " + count);
        }
    }
);

4. AllWindowedStream(全局窗口流)

对 DataStream 直接开窗(不需要 keyBy),用于全局聚合统计。

⚠️ 注意:并行度只能是 1!

如何得到

AllWindowedStream<Order, TimeWindow> allWindowedStream = orderStream
    .windowAll(TumblingEventTimeWindows.of(Time.minutes(5)));

// 或者
AllWindowedStream<Order, TimeWindow> allWindowedStream = orderStream
    .timeWindowAll(Time.minutes(5));

转换回 DataStream

// 全局统计:整个平台每5分钟的订单总数
DataStream<Long> totalCount = orderStream
    .timeWindowAll(Time.minutes(5))
    .apply(new AllWindowFunction<Order, Long, TimeWindow>() {
        @Override
        public void apply(TimeWindow window, Iterable<Order> values, 
                         Collector<Long> out) {
            long count = 0;
            for (Order order : values) {
                count++;
            }
            out.collect(count);
        }
    });

WindowedStream vs AllWindowedStream

对比项WindowedStreamAllWindowedStream
前提需要先 keyBy不需要 keyBy
并行度可以并行(多个Task)只能是 1
用途分组聚合(如每个用户的统计)全局聚合(如全平台的统计)

5. ConnectedStreams(连接的两个流)

连接两个流,可以共享状态。两个流的数据类型可以不同。

如何得到

// 订单流连接规则流
ConnectedStreams<Order, Rule> connectedStreams = orderStream.connect(ruleStream);

转换回 DataStream

// 方式1:process(最灵活,推荐)
DataStream<Order> result = connectedStreams.process(
    new CoProcessFunction<Order, Rule, Order>() {
        
        // 状态:保存当前规则
        private Rule currentRule;
        
        @Override
        public void processElement1(Order order, Context ctx, 
                                   Collector<Order> out) {
            // 处理订单流(第一个流)
            if (currentRule == null || order.getAmount() < currentRule.getThreshold()) {
                out.collect(order);  // 正常订单
            } else {
                // 触发风控
            }
        }
        
        @Override
        public void processElement2(Rule rule, Context ctx, 
                                   Collector<Order> out) {
            // 处理规则流(第二个流)
            this.currentRule = rule;
            System.out.println("规则更新: " + rule);
            // 注意:这里也可以 out.collect() 输出数据
        }
    }
);

// 方式2:map(两个流都必须输出)
DataStream<String> result = connectedStreams.map(
    new CoMapFunction<Order, Rule, String>() {
        @Override
        public String map1(Order order) {
            return "订单: " + order.getOrderId();
        }
        
        @Override
        public String map2(Rule rule) {
            return "规则: " + rule.getRuleId();
        }
    }
);

// 方式3:flatMap(可以选择性输出)
DataStream<Order> result = connectedStreams.flatMap(
    new CoFlatMapFunction<Order, Rule, Order>() {
        
        private Rule currentRule;
        
        @Override
        public void flatMap1(Order order, Collector<Order> out) {
            // 处理订单,可以输出 0 条或多条
            if (currentRule == null || order.getAmount() < currentRule.getThreshold()) {
                out.collect(order);
            }
        }
        
        @Override
        public void flatMap2(Rule rule, Collector<Order> out) {
            // 只更新规则,不输出
            this.currentRule = rule;
        }
    }
);

使用场景

  • 动态规则/配置更新
  • 两个流共享状态
  • 两个流的数据类型可以不同

6. BroadcastConnectedStream(广播连接流)

将一个流广播给另一个 KeyedStream 的所有并行 Task。

如何得到

// 1. 定义广播状态描述符
MapStateDescriptor<String, Rule> ruleStateDesc = 
    new MapStateDescriptor<>("rules", String.class, Rule.class);

// 2. 将规则流转为广播流
BroadcastStream<Rule> broadcastRules = ruleStream.broadcast(ruleStateDesc);

// 3. KeyedStream 连接广播流
BroadcastConnectedStream<Order, Rule> broadcastConnectedStream = 
    orderStream
        .keyBy(order -> order.getUserId())
        .connect(broadcastRules);

转换回 DataStream

DataStream<Order> result = broadcastConnectedStream.process(
    new KeyedBroadcastProcessFunction<String, Order, Rule, Order>() {
        
        @Override
        public void processElement(Order order, ReadOnlyContext ctx, 
                                  Collector<Order> out) {
            // 处理订单:读取广播状态
            Rule rule = ctx.getBroadcastState(ruleStateDesc).get("current");
            if (rule == null || order.getAmount() < rule.getThreshold()) {
                out.collect(order);
            }
        }
        
        @Override
        public void processBroadcastElement(Rule rule, Context ctx, 
                                           Collector<Order> out) {
            // 处理规则:更新广播状态(所有 Task 都会执行)
            ctx.getBroadcastState(ruleStateDesc).put("current", rule);
        }
    }
);

ConnectedStreams vs BroadcastConnectedStream

对比项ConnectedStreamsBroadcastConnectedStream
规则流分发只发到一个 Task广播到所有 Task
主流要求DataStream通常是 KeyedStream
状态管理普通状态广播状态(BroadcastState)
适用场景简单的两流关联规则/配置需要全局生效

7. JoinedStreams(JOIN 关联流)

两个流按 Key 在窗口内进行 INNER JOIN。

如何得到 & 转换回 DataStream

DataStream<OrderWithPayment> result = orderStream
    .join(paymentStream)
    .where(order -> order.getOrderId())      // 左流的 Key
    .equalTo(payment -> payment.getOrderId()) // 右流的 Key
    .window(TumblingEventTimeWindows.of(Time.minutes(10)))
    .apply(new JoinFunction<Order, Payment, OrderWithPayment>() {
        @Override
        public OrderWithPayment join(Order order, Payment payment) {
            return new OrderWithPayment(order, payment);
        }
    });

⚠️ 注意:join() 只能实现 INNER JOIN! 两边都有数据才会输出。


8. CoGroupedStreams(协同分组流)- 万能 JOIN ⭐

两个流按 Key 在窗口内分组,可以实现任意类型的 JOIN。

如何得到 & 转换回 DataStream

DataStream<String> result = orderStream
    .coGroup(paymentStream)
    .where(order -> order.getOrderId())
    .equalTo(payment -> payment.getOrderId())
    .window(TumblingEventTimeWindows.of(Time.minutes(10)))
    .apply(new CoGroupFunction<Order, Payment, String>() {
        @Override
        public void coGroup(Iterable<Order> orders, 
                           Iterable<Payment> payments, 
                           Collector<String> out) {
            
            List<Order> orderList = new ArrayList<>();
            orders.forEach(orderList::add);
            
            List<Payment> paymentList = new ArrayList<>();
            payments.forEach(paymentList::add);
            
            // 可以实现任意 JOIN 逻辑
            if (!orderList.isEmpty() && !paymentList.isEmpty()) {
                // INNER JOIN:两边都有
                out.collect("匹配成功: " + orderList.get(0).getOrderId());
            } else if (!orderList.isEmpty() && paymentList.isEmpty()) {
                // LEFT JOIN:只有左边(订单未支付)
                out.collect("未支付订单: " + orderList.get(0).getOrderId());
            } else if (orderList.isEmpty() && !paymentList.isEmpty()) {
                // RIGHT JOIN:只有右边(异常支付)
                out.collect("异常支付: " + paymentList.get(0).getOrderId());
            }
        }
    });

join() vs coGroup()

对比项join()coGroup()
JOIN 类型只能 INNER JOIN任意 JOIN(LEFT/RIGHT/FULL/INNER)
数据获取一对一匹配获取两边所有数据(Iterable)
灵活性
使用场景简单关联对账、复杂关联逻辑


五、常用转换代码模板

模板1:分组聚合(最常用)

// 每个用户每5分钟的订单金额
DataStream<Order> result = orderStream
    .keyBy(order -> order.getUserId())      // DataStream → KeyedStream
    .timeWindow(Time.minutes(5))            // KeyedStream → WindowedStream
    .sum("amount");                         // WindowedStream → DataStream

模板2:全局聚合

// 全平台每5分钟的订单总数(并行度=1)
DataStream<Long> result = orderStream
    .timeWindowAll(Time.minutes(5))         // DataStream → AllWindowedStream
    .apply((window, values, out) -> {       // AllWindowedStream → DataStream
        long count = 0;
        for (Order o : values) count++;
        out.collect(count);
    });

模板3:动态规则(connect)

// 订单流 + 规则流
DataStream<Order> result = orderStream
    .connect(ruleStream)                    // DataStream → ConnectedStreams
    .process(new CoProcessFunction<>() {    // ConnectedStreams → DataStream
        private Rule rule;
        
        public void processElement1(Order order, Context ctx, Collector<Order> out) {
            if (rule == null || order.getAmount() < rule.getThreshold()) {
                out.collect(order);
            }
        }
        
        public void processElement2(Rule rule, Context ctx, Collector<Order> out) {
            this.rule = rule;
        }
    });

模板4:广播规则(broadcast + connect)

// 规则广播给所有 Task
MapStateDescriptor<String, Rule> desc = new MapStateDescriptor<>("rules", String.class, Rule.class);

DataStream<Order> result = orderStream
    .keyBy(order -> order.getUserId())
    .connect(ruleStream.broadcast(desc))    // KeyedStream → BroadcastConnectedStream
    .process(new KeyedBroadcastProcessFunction<>() {  // → DataStream
        public void processElement(Order order, ReadOnlyContext ctx, Collector<Order> out) {
            Rule rule = ctx.getBroadcastState(desc).get("current");
            if (rule == null || order.getAmount() < rule.getThreshold()) {
                out.collect(order);
            }
        }
        
        public void processBroadcastElement(Rule rule, Context ctx, Collector<Order> out) {
            ctx.getBroadcastState(desc).put("current", rule);
        }
    });

模板5:两流 JOIN

// INNER JOIN
DataStream<OrderWithPayment> result = orderStream
    .join(paymentStream)
    .where(Order::getOrderId)
    .equalTo(Payment::getOrderId)
    .window(TumblingEventTimeWindows.of(Time.minutes(10)))
    .apply((order, payment) -> new OrderWithPayment(order, payment));

模板6:两流对账(coGroup)

// 可以实现 LEFT/RIGHT/FULL JOIN
DataStream<String> result = orderStream
    .coGroup(paymentStream)
    .where(Order::getOrderId)
    .equalTo(Payment::getOrderId)
    .window(TumblingEventTimeWindows.of(Time.hours(1)))
    .apply((orders, payments, out) -> {
        // 自定义 JOIN 逻辑
        List<Order> orderList = Lists.newArrayList(orders);
        List<Payment> paymentList = Lists.newArrayList(payments);
        
        if (!orderList.isEmpty() && paymentList.isEmpty()) {
            out.collect("未支付: " + orderList.get(0).getOrderId());
        }
    });

六、常见问题

Q1: KeyedStream 和 DataStream 有什么区别?

  • DataStream:普通数据流,数据随机分布在各个 Task
  • KeyedStream:按 Key 分组后的流,相同 Key 的数据在同一个 Task,可以使用 KeyedState

Q2: 什么时候用 connect,什么时候用 join?

  • connect:两个流共享状态,类型可以不同,不需要窗口
  • join:两个流按 Key 关联,需要窗口,只能 INNER JOIN

Q3: 为什么 AllWindowedStream 并行度只能是 1?

因为没有 keyBy,所有数据必须发到同一个 Task 才能做全局聚合。如果数据量大,应该用 keyBy().window() 先分组聚合,再 windowAll() 汇总。

Q4: coGroup 和 join 怎么选?

  • 简单 INNER JOIN:用 join(),代码更简洁
  • 需要 LEFT/RIGHT/FULL JOIN:用 coGroup()
  • 需要复杂关联逻辑(如对账):用 coGroup()

希望本文对你理解 Flink DataStream 类型转换有所帮助!

更多推荐