使用 Apache Flink 进行流式数据计算,通常包括 环境搭建 → 数据源接入 → 流式计算逻辑 → 数据输出 → 运行与监控 几个步骤。下面我从 概念 + 示例 + 常见场景 三个层面,系统性地讲一下。
流式计算是指 数据持续不断到达,系统实时处理并产生结果,而不是等数据全部就绪再计算。
Flink 的特点:
DataStream API(主流)
↓
Source → Transformation → Sink
核心抽象:
StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();
常用配置:
env.enableCheckpointing(5000); // 每5秒一次 checkpoint
env.setParallelism(2);
DataStream<String> stream = env
.socketTextStream("localhost", 9999);
Kafka Source 示例:
KafkaSource<String> source = KafkaSource
.<String>builder()
.setBootstrapServers("localhost:9092")
.setTopics("input-topic")
.setGroupId("flink-group")
.setValueOnlyDeserializer(new SimpleStringSchema())
.build();
DataStream<String> stream = env.fromSource(
source,
WatermarkStrategy.noWatermarks(),
"Kafka Source"
);
DataStream<Tuple2<String, Integer>> result =
stream
.flatMap((String line, Collector<Tuple2<String, Integer>> out) -> {
for (String word : line.split(" ")) {
out.collect(new Tuple2<>(word, 1));
}
})
.returns(Types.TUPLE(Types.STRING, Types.INT))
.keyBy(value -> value.f0)
.sum(1);
滚动窗口(Tumbling Window)
DataStream<WordCount> result =
stream
.map(line -> new WordCount(line, 1))
.keyBy(WordCount::getWord)
.window(TumblingProcessingTimeWindows.of(Time.seconds(10)))
.sum("count");
滑动窗口(Sliding Window)
.window(SlidingProcessingTimeWindows.of(Time.seconds(30), Time.seconds(10)))
WatermarkStrategy<String> strategy =
WatermarkStrategy
.<String>forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withTimestampAssigner((event, timestamp) -> extractTimestamp(event));
DataStream<String> stream = env
.fromSource(kafkaSource, strategy, "kafka");
stream.sinkTo(kafkaSink);
result.print();
result.addSink(new MySQLSink());
env.execute("Flink Streaming Job");
⚠️ Flink 是懒执行,只有 execute 才会运行
public class StreamWordCount {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<String> text =
env.socketTextStream("localhost", 9999);
DataStream<Tuple2<String, Integer>> counts =
text.flatMap((String line, Collector<Tuple2<String, Integer>> out) -> {
for (String word : line.split(" ")) {
out.collect(new Tuple2<>(word, 1));
}
})
.returns(Types.TUPLE(Types.STRING, Types.INT))
.keyBy(t -> t.f0)
.sum(1);
counts.print();
env.execute("Stream WordCount");
}
}
| 场景 | 核心用法 |
|---|---|
| 实时 ETL | map / filter |
| 实时统计 | keyBy + window |
| 实时去重 | state / keyBy |
| 实时 TopN | KeyedProcessFunction |
| 实时对账 | CoProcessFunction |
| CDC 同步 | Flink CDC |
| 复杂事件处理 | CEP |
env.enableCheckpointing(5000);
如果你愿意,我可以:
✅ 给你一个 Kafka + Flink CDC 的完整示例
✅ 手把手教你 实时 PV / UV 统计
✅ 解释 窗口触发、Watermark、迟到数据
✅ 对比 Flink vs Spark Streaming
你可以直接告诉我你现在的 使用场景 或 技术栈。
免责声明:本站发布的内容(图片、视频和文字)以原创、转载和分享为主,文章观点不代表本网站立场,如果涉及侵权请联系站长邮箱:is@yisu.com进行举报,并提供相关证据,一经查实,将立刻删除涉嫌侵权内容。