Flink DataStream 类型转换
本文详细介绍 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) |
| JoinedStreams | JOIN 关联流 | 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:
| 对比项 | WindowedStream | AllWindowedStream |
|---|---|---|
| 前提 | 需要先 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:
| 对比项 | ConnectedStreams | BroadcastConnectedStream |
|---|---|---|
| 规则流分发 | 只发到一个 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:普通数据流,数据随机分布在各个 TaskKeyedStream:按 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 类型转换有所帮助!
更多推荐


所有评论(0)