Apache Flink 流处理基础教程
Mei Lin | 2026-09-02T01:08:03 | Database, Cloud
从零开始学习 Apache Flink 的流处理编程模型,涵盖 DataStream API、窗口计算、状态管理和 Exactly-Once 语义。
# Apache Flink 流处理基础教程 ## 为什么选 Flink 在实时数据处理领域,Flink 凭借以下特性脱颖而出: - 真正的流处理(非微批) - Exactly-Once 语义保证 - 毫秒级延迟 - 强大的窗口机制 - 统一的批流一体 API ## 核心概念 ``` Source (数据源) -> Transformation (算子: map/filter/keyBy/window) -> Sink (输出) ``` ## 环境搭建 ```xml org.apache.flink flink-streaming-java 1.19.0 org.apache.flink flink-connector-kafka 3.1.0-1.19 org.apache.flink flink-clients 1.19.0 ``` ## Hello World: 实时词频统计 ```java public class WordCount { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); DataStream text = env.socketTextStream("localhost", 9999); DataStream> counts = text .flatMap(new Tokenizer()) .keyBy(value -> value.f0) .sum(1); counts.print(); env.execute("Word Count"); } public static class Tokenizer implements FlatMapFunction> { @Override public void flatMap(String value, Collector> out) { for (String word : value.toLowerCase().split("\\W+")) { if (word.length() > 0) { out.collect(new Tuple2(word, 1)); } } } } } ``` ## 从 Kafka 消费 ```java KafkaSource source = KafkaSource.builder() .setBootstrapServers("localhost:9092") .setTopics("user-events") .setGroupId("flink-group") .setStartingOffsets(OffsetsInitializer.earliest()) .setValueOnlyDeserializer( new SimpleStringSchema()) .build(); DataStream events = env.fromSource( source, WatermarkStrategy.noWatermarks(), "Kafka Source"); ``` ## 窗口计算 ### 滚动窗口(Tumbling Window) 每 5 分钟统计一次: ```java DataStream> result = events .map(event -> parseEvent(event)) // (userId, 1L) .keyBy(t -> t.f0) .window(TumblingProcessingTimeWindows.of(Duration.ofMinutes(5))) .sum(1); ``` ### 滑动窗口(Sliding Window) 每分钟计算最近 5 分钟的数据: ```java DataStream> result = events .keyBy(t -> t.f0) .window(SlidingProcessingTimeWindows.of( Duration.ofMinutes(5), // 窗口大小 Duration.ofMinutes(1))) // 滑动步长 .sum(1); ``` ### 会话窗口(Session Window) 用户不活跃 10 分钟后关闭窗口: ```java DataStream> result = events .keyBy(t -> t.f0) .window(ProcessingTimeSessionWindows.withGap( Duration.ofMinutes(10))) .sum(1); ``` ## 事件时间与 Watermark 处理乱序事件的关键: ```java WatermarkStrategy strategy = WatermarkStrategy .forBoundedOutOfOrderness(Duration.ofSeconds(5)) .withTimestampAssigner((event, timestamp) -> event.getTimestamp()); DataStream stream = env .fromSource(source, strategy, "Source") .keyBy(Event::getUserId) .window(TumblingEventTimeWindows.of(Duration.ofMinutes(1))) .process(new MyWindowFunction()); ``` ## 状态管理 Flink 提供了强大的状态管理机制: ```java public class FraudDetector extends KeyedProcessFunction { // 声明状态 private ValueState flagState; private ValueState timerState; @Override public void open(Configuration config) { ValueStateDescriptor flagDescriptor = new ValueStateDescriptor("flag", Types.BOOLEAN); flagState = getRuntimeContext().getState(flagDescriptor); ValueStateDescriptor timerDescriptor = new ValueStateDescriptor("timer", Types.LONG); timerState = getRuntimeContext().getState(timerDescriptor); } @Override public void processElement(Transaction tx, Context ctx, Collector out) throws Exception { Boolean lastSmallTx = flagState.value(); if (lastSmallTx != null && tx.getAmount() > 500) { // 小额交易后紧跟大额交易 -> 可疑 out.collect(new Alert(tx.getAccountId())); } if (tx.getAmount() out) { // 超时清除标记 flagState.clear(); timerState.clear(); } } ``` ## Checkpoint 与容错 ```java // 启用 Checkpoint env.enableCheckpointing(60000); // 每 60 秒 env.getCheckpointConfig().setCheckpointingMode( CheckpointingMode.EXACTLY_ONCE); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000); env.getCheckpointConfig().setCheckpointStorage( "s3://my-bucket/flink-checkpoints"); ``` ## 部署模式 1. **Standalone**:开发测试用 2. **YARN**:Hadoop 生态集成 3. **Kubernetes**:推荐的生产方式,配合 Flink Operator ## 总结 Flink 是实时流处理的标杆框架。理解好 DataStream API、窗口机制和状态管理三个核心概念,就能应对大多数实时数据处理场景。在生产环境中,Checkpoint 和 Exactly-Once 语义保证了数据的正确性和系统的容错能力。