在 Apache Flink 中实现复杂的数据转换,通常离不开 DataStream API / Table API & SQL / 状态(State)/ 时间语义(Event Time / Watermark)/ 窗口 / 自定义函数 的组合使用。下面从 常见复杂场景 → 对应实现方式 → 示例代码 系统说明。
| 场景 | 推荐方式 |
|---|---|
| 结构化数据计算 | Table API / SQL |
| 流式事件驱动 | DataStream API |
| 状态复杂 | KeyedState / ValueState |
| 时序模式检测 | CEP |
| 简单 ETL | Map / FlatMap |
| 多流合并 | union / connect / join |
DataStream<Order> result = orderStream
.flatMap(new RichFlatMapFunction<Order, Order>() {
@Override
public void flatMap(Order value, Collector<Order> out) {
if (value.getAmount() > 100) {
value.setLevel("HIGH");
} else {
value.setLevel("LOW");
}
out.collect(value);
}
});
✅ 适合:字段补全、数据增强、条件拆分
DataStream<Result> result = stream
.keyBy(Order::getUserId)
.map(new RichMapFunction<Order, Result>() {
private ValueState<Integer> countState;
@Override
public void open(Configuration parameters) {
countState = getRuntimeContext().getState(
new ValueStateDescriptor<>("count", Integer.class)
);
}
@Override
public Result map(Order order) throws Exception {
Integer count = countState.value();
count = count == null ? 1 : count + 1;
countState.update(count);
return new Result(order.getUserId(), count);
}
});
✅ 适合:
DataStream<Summary> result = stream
.keyBy(Order::getShopId)
.window(TumblingEventTimeWindows.of(Time.minutes(5)))
.aggregate(new AggregateFunction<Order, Tuple2<Integer, Double>, Summary>() {
@Override
public Tuple2<Integer, Double> createAccumulator() {
return Tuple2.of(0, 0.0);
}
@Override
public Tuple2<Integer, Double> add(Order order, Tuple2<Integer, Double> acc) {
return Tuple2.of(acc.f0 + 1, acc.f1 + order.getAmount());
}
@Override
public Summary getResult(Tuple2<Integer, Double> acc) {
return new Summary(acc.f0, acc.f1);
}
@Override
public Tuple2<Integer, Double> merge(Tuple2<Integer, Double> a,
Tuple2<Integer, Double> b) {
return Tuple2.of(a.f0 + b.f0, a.f1 + b.f1);
}
});
✅ 适合:
DataStream<EnrichedOrder> result =
orderStream
.keyBy(Order::getUserId)
.connect(userStream.keyBy(User::getId))
.process(new CoProcessFunction<Order, User, EnrichedOrder>() {
private ValueState<User> userState;
@Override
public void open(Configuration parameters) {
userState = getRuntimeContext().getState(
new ValueStateDescriptor<>("user", User.class)
);
}
@Override
public void processElement1(Order order, Context ctx,
Collector<EnrichedOrder> out) {
User user = userState.value();
if (user != null) {
out.collect(new EnrichedOrder(order, user));
}
}
@Override
public void processElement2(User user, Context ctx,
Collector<EnrichedOrder> out) {
userState.update(user);
}
});
✅ 适合:
Pattern<LoginEvent, ?> pattern = Pattern
.<LoginEvent>begin("first").where(e -> !e.isSuccess())
.next("second").where(e -> !e.isSuccess())
.next("third").where(e -> !e.isSuccess());
DataStream<Alert> alerts = CEP.pattern(loginStream, pattern)
.select(patternMap -> {
LoginEvent e = patternMap.get("third").iterator().next();
return new Alert(e.getUserId(), "连续失败3次");
});
✅ 适合:
DataStream<Row> result = stream.map(json -> {
JSONObject obj = JSON.parseObject(json);
return Row.of(
obj.getString("order_id"),
obj.getJSONObject("user").getInteger("age"),
obj.getJSONArray("items").size()
);
});
或直接使用 Flink Table API JSON 函数
SELECT
userId,
COUNT(*) AS cnt,
SUM(amount) AS total
FROM orders
WHERE amount > 100
GROUP BY
userId,
TUMBLE(event_time, INTERVAL '5' MINUTE)
✅ 好处:
你可以直接贴你的数据流结构或业务需求,我给你写完整示例。
免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。