05. Apache Flink 实时计算详解

深入 Apache Flink 实时计算引擎:DataStream API 编程模型、Checkpoint 容错机制、Watermark 事件时间与窗口计算、以及 Exactly-Once 语义实现原理。

Apache Flink 是业界领先的分布式流处理引擎,其基于事件时间的窗口计算和精确的 Checkpoint 机制使其成为实时数据处理的首选方案。本文从 DataStream API 出发,深入讲解 Flink 的核心原理与生产实践。

1. DataStream API 编程模型

              Flink Client
         提交 Job、生成 JobGraph
                   |
           JobManager (主节点)
  +-------------+ +-------------+ +----------+
  | Dispatcher  | | JobMaster   | | Resource |
  | (接收提交)   | | (调度执行)   | | Manager  |
  +-------------+ +-------------+ +----------+
                   |
          TaskManager (工作节点) x N
  +------------+ +------------+ +------------+
  | Slot 1     | | Slot 2     | | Slot ...   |
  | (Task)     | | (Task)     | | (Task)     |
  +------------+ +------------+ +------------+
         |              |              |
    Netty 网络传输 (Backpressure 背压)

1.2 DataStream 基础操作

import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
import org.apache.flink.streaming.api.windowing.time.Time;

public class StreamingJob {
    public static void main(String[] args) throws Exception {
        final StreamExecutionEnvironment env =
            StreamExecutionEnvironment.getExecutionEnvironment();

        // 1. 读取数据流(Source)
        DataStream<String> stream = env
            .fromSource(
                kafkaSource,
                WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(5)),
                "Kafka Source"
            );

        // 2. 转换(Transformation)
        DataStream<OrderEvent> events = stream
            .map(new JsonToOrderMapper())
            .filter(e -> e.amount > 0);

        // 3. KeyBy + 窗口聚合
        DataStream<Tuple2<String, Double>> result = events
            .keyBy(e -> e.category)
            .window(TumblingEventTimeWindows.of(Time.minutes(1)))
            .aggregate(new SumAggregate());

        // 4. 输出(Sink)
        result.addSink(new FlinkJedisPoolConfig.Builder()
            .setHost("redis")
            .build());

        // 执行
        env.execute("RealTimeAnalytics");
    }
}

1.3 转换算子分类

算子类型示例并行度数据流特征
Map/Filter/FlatMap.map(), .filter()1:1One-to-one
KeyBy.keyBy(x -> x.userId)重分区Redistribution
Reduce/Aggregate.reduce(), .aggregate()同 Key 内Keyed Stream
Window.window(…), .timeWindow(…)同 Key 内窗口内聚合
Union/Connect.union(), .connect()合并流Multi-stream
ProcessFunction.process()灵活底层 API
// KeyedProcessFunction:底层 API,可访问状态和定时器
class OrderTimeoutAlert extends KeyedProcessFunction<String, OrderEvent, Alert> {
    private ValueState<Long> timerState;

    @Override
    public void processElement(OrderEvent event, Context ctx, Collector<Alert> out) {
        long timeout = event.createTime + 15 * 60 * 1000; // 15分钟超时
        ctx.timerService().registerEventTimeTimer(timeout);
        timerState.update(timeout);
    }

    @Override
    public void onTimer(long timestamp, OnTimerContext ctx, Collector<Alert> out) {
        out.collect(new Alert(ctx.getCurrentKey(), "ORDER_TIMEOUT"));
    }
}

2. Checkpoint 容错机制

2.1 Checkpoint 原理

Checkpoint 是 Flink 实现 Exactly-Once 的核心机制,基于 Chandy-Lamport 分布式快照算法

JobManager (Checkpoint Coordinator)
    |
    |  Trigger Checkpoint (Barrier)
    v
Source -> Map -> KeyBy -> Window -> Sink
   |                             |
   |  <--- Barrier 注入 -----------|
   |                             |
   |  Source Snapshot(offset)     |
   |  -> 向下游广播 Barrier        |
   |                             |
   |       Barrier 对齐 (Alignment)
   |                             |
   |  各算子保存 State 到 State Backend
   |                             |
   |  Sink 预提交 (Pre-Commit)    |
   |                             |
   |  JobManager 确认完成 → 通知 Checkpoint 成功

Barrier 对齐 vs 非对齐

模式原理延迟适用场景
Aligned等待所有输入通道 Barrier 到达较低低反压场景
UnalignedBarrier 超越数据先行更低高反压、要求低延迟
// Checkpoint 配置
env.enableCheckpointing(60000); // 60s 间隔
env.getCheckpointConfig().setCheckpointingMode(
    CheckpointingMode.EXACTLY_ONCE
);
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000);
env.getCheckpointConfig().setCheckpointTimeout(600000);
env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);
env.getCheckpointConfig().enableUnalignedCheckpoints();

// State Backend
env.setStateBackend(new HashMapStateBackend());
env.getCheckpointConfig().setCheckpointStorage("hdfs:///checkpoints");

2.2 State 状态管理

State 类型说明适用场景
ValueState单值状态计数器、最新值
ListState列表状态缓存事件集合
MapStateMap 结构Key-Value 查找
ReducingState用于 Reduce 的单个值增量聚合
AggregatingState用于 Aggregate 的单个值复杂增量聚合
// ValueState 示例:去重计数
class UniqueUserCount extends RichFlatMapFunction<Event, Metrics> {
    private ValueState<HashSet<String>> userState;

    @Override
    public void open(Configuration parameters) {
        StateTtlConfig ttl = StateTtlConfig
            .newBuilder(Time.hours(24))
            .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
            .setStateVisibility(StateTtlConfig.StateVisibility.ReturnExpiredIfNotCleanedUp)
            .cleanupFullSnapshot()
            .build();

        userState = getRuntimeContext().getState(
            new ValueStateDescriptor<>("users", Types.SET(Types.STRING))
        );
        // 启用 TTL
        ((ValueStateDescriptor<HashSet<String>>) userState.getStateDescriptor()).enableTimeToLive(ttl);
    }

    @Override
    public void flatMap(Event event, Collector<Metrics> out) throws Exception {
        HashSet<String> users = userState.value();
        if (users == null) users = new HashSet<>();
        users.add(event.userId);
        userState.update(users);
        out.collect(new Metrics(event.windowStart, users.size()));
    }
}

2.3 State Backend 对比

Backend存储介质容量快照方式适用
MemoryStateBackendJVM Heap同步本地测试
FsStateBackendJVM Heap + 异步文件异步小状态生产
RocksDBStateBackendRocksDB(本地磁盘)大(可溢出)增量异步大状态生产
// RocksDB 增量 Checkpoint(推荐生产配置)
EmbeddedRocksDBStateBackend rocksDb = new EmbeddedRocksDBStateBackend(true);
env.setStateBackend(rocksDb);
env.getCheckpointConfig().setCheckpointStorage("hdfs:///flink-checkpoints");

3. Watermark 与事件时间

3.1 三种时间语义

时间类型来源特点适用场景
Event Time数据自带时间戳最准确、处理乱序日志分析、订单统计
Ingestion Time进入 Flink 的时间无需提取时间戳近似有序场景
Processing Time算子处理的时间最低延迟实时监控、告警
// Event Time 配置
env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);

// Watermark 策略
stream.assignTimestampsAndWatermarks(
    WatermarkStrategy
        .<Event>forBoundedOutOfOrderness(Duration.ofSeconds(30))
        .withTimestampAssigner((event, timestamp) -> event.eventTime)
        .withIdleness(Duration.ofMinutes(5)) // 空闲源处理
);

3.2 Watermark 传播与生成策略

事件流 (Event Time):
[1s] [2s] [5s] [3s] [8s] [12s] [15s]   ← 乱序到达

With BoundedOutOfOrderness(maxDelay=3s):
Watermark = maxEventTime - maxDelay

Event     Watermark    说明
1s        -2s          (无)
2s        -1s          (无)
5s         2s          ← 时间推进到 2s,窗口 [0,5) 触发
3s         2s          ← 迟到(但 < 3s 延迟),可处理
8s         5s          ← 窗口 [5,10) 触发
12s        9s          ← 窗口 [10,15) 未触发
15s       12s          ← 窗口 [10,15) 触发
// 自定义 Watermark 生成器
class PunctuatedWatermarkStrategy implements WatermarkStrategy<Event> {
    @Override
    public WatermarkGenerator<Event> createWatermarkGenerator(WatermarkGeneratorSupplier.Context ctx) {
        return new WatermarkGenerator<Event>() {
            private long maxTimestamp = Long.MIN_VALUE;

            @Override
            public void onEvent(Event event, long eventTimestamp, WatermarkOutput output) {
                maxTimestamp = Math.max(maxTimestamp, eventTimestamp);
                // 特定事件触发 Watermark 推进
                if (event.isEndOfBatch) {
                    output.emitWatermark(new Watermark(maxTimestamp));
                }
            }

            @Override
            public void onPeriodicEmit(WatermarkOutput output) {
                output.emitWatermark(new Watermark(maxTimestamp - 5000));
            }
        };
    }
}

3.3 迟到数据处理

stream
    .keyBy(e -> e.userId)
    .window(TumblingEventTimeWindows.of(Time.minutes(5)))
    .allowedLateness(Time.minutes(2))    // 允许 2 分钟迟到
    .sideOutputLateData(lateTag)         // 超时的进侧输出流
    .aggregate(new CountAggregate());

// 处理侧输出迟到数据
DataStream<Event> lateStream = result.getSideOutput(lateTag);
lateStream.addSink(new LateDataSink()); // 写入死信队列/冷存储

4. 窗口计算

4.1 窗口类型对比

窗口类型触发条件重叠性适用场景
Tumbling固定时间间隔不重叠每五分钟统计
Sliding固定时间间隔 + 滑动步长重叠过去一小时的每分钟统计
Session活动间隙(Gap)动态用户行为会话分析
Global全局单一窗口-需要自定义触发器
// Tumbling Window:每 5 分钟统计一次
stream.keyBy(e -> e.category)
    .window(TumblingEventTimeWindows.of(Time.minutes(5)))
    .aggregate(new AverageAggregate());

// Sliding Window:每 1 分钟输出过去 5 分钟的统计
stream.keyBy(e -> e.category)
    .window(SlidingEventTimeWindows.of(Time.minutes(5), Time.minutes(1)))
    .aggregate(new SumAggregate());

// Session Window:10 分钟无活动则关闭会话
stream.keyBy(e -> e.userId)
    .window(EventTimeSessionWindows.withGap(Time.minutes(10)))
    .aggregate(new SessionAggregate());

4.2 窗口函数

函数说明输出
reduce()两两归约同类型
aggregate()增量聚合(输入/累加器/输出可不同)任意类型
apply()全量窗口函数(获取窗口内所有数据)任意类型
process()ProcessWindowFunction(带 Context)任意类型
// AggregateFunction:增量聚合,内存友好
class AverageAggregate implements AggregateFunction<Event, Tuple2<Long, Long>, Double> {
    @Override
    public Tuple2<Long, Long> createAccumulator() {
        return Tuple2.of(0L, 0L);
    }

    @Override
    public Tuple2<Long, Long> add(Event event, Tuple2<Long, Long> acc) {
        return Tuple2.of(acc.f0 + event.value, acc.f1 + 1);
    }

    @Override
    public Double getResult(Tuple2<Long, Long> acc) {
        return acc.f0 / (double) acc.f1;
    }

    @Override
    public Tuple2<Long, Long> merge(Tuple2<Long, Long> a, Tuple2<Long, Long> b) {
        return Tuple2.of(a.f0 + b.f0, a.f1 + b.f1);
    }
}

// ProcessWindowFunction:获取窗口元数据(全量聚合,内存消耗大)
class TopNProcess extends ProcessWindowFunction<Event, Result, String, TimeWindow> {
    @Override
    public void process(String key, Context ctx, Iterable<Event> events, Collector<Result> out) {
        long windowStart = ctx.window().getStart();
        List<Event> sorted = StreamSupport.stream(events.spliterator(), false)
            .sorted(Comparator.comparing(e -> -e.value))
            .limit(10)
            .collect(Collectors.toList());
        out.collect(new Result(key, windowStart, sorted));
    }
}

5. Exactly-Once 语义

5.1 端到端 Exactly-Once

Source (Kafka)         Flink Transformations               Sink (Kafka)
  +--------+           +-------+  +-------+  +-------+      +---------+
  | Offset | --------> | Map   |->| KeyBy |->|Window | ---> | Producer|
  | Committed           |       |  |       |  |       |      | Txn     |
  | (2PC)               | State |  | State |  | State |      | (2PC)   |
  +--------+           +-------+  +-------+  +-------+      +---------+
       |                                                        |
       |<---------------- Two-Phase Commit --------------------->|
       |                                                        |
  Pre-commit → Snapshot State → Commit Kafka Txn → Commit Offset

5.2 两阶段提交实现

阶段SourceSink
Pre-commit保存 offset 到状态预提交事务(Kafka txn.begin)
Checkpoint状态快照至 State Backend状态包含事务 ID
Commit随 Checkpoint 完成确认提交事务(Kafka txn.commit)
Abort-Checkpoint 失败时回滚事务
// Kafka Producer 配置:Exactly-Once Sink
Properties props = new Properties();
props.put("bootstrap.servers", "kafka:9092");
props.put("transactional.id", "flink-producer-1"); // 事务 ID
props.put("enable.idempotence", "true");
props.put("acks", "all");

FlinkKafkaProducer<String> sink = new FlinkKafkaProducer<>(
    "output-topic",
    new SimpleStringSchema(),
    props,
    FlinkKafkaProducer.Semantic.EXACTLY_ONCE
);

stream.addSink(sink);

5.3 Exactly-Once 与 At-Least-Once 对比

语义数据保证Checkpoint 开销Sink 要求性能
Exactly-Once无重复、无丢失高(Barrier 对齐/非对齐)需支持事务较低
At-Least-Once不丢失、可能重复低(无需 Barrier 对齐)幂等或去重较高
// At-Least-Once 配置(允许重复,性能更高)
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.AT_LEAST_ONCE);

// Sink 配合幂等写入(如使用业务主键)
// 即使重复投递,最终结果一致

6. 生产实战配置

# YARN Session 模式
./bin/yarn-session.sh -n 4 -jm 2048 -tm 8192 -s 4

# Per-Job 模式
./bin/flink run -m yarn-cluster \
    -yn 4 -yjm 2048 -ytm 8192 \
    -ys 4 -p 16 \
    ./job.jar

# Kubernetes Native
./bin/flink run-application \
    --target kubernetes-application \
    -Dkubernetes.cluster-id=flink-cluster \
    -Dkubernetes.container.image=flink-job:latest \
    ./job.jar

6.2 关键调参

参数建议值说明
execution.checkpointing.interval30s - 10min平衡容错与性能
execution.checkpointing.max-concurrent-checkpoints1防止过多 Checkpoint 干扰
state.backend.incrementaltrueRocksDB 增量快照
state.backend.rocksdb.memory.managedtrue托管 RocksDB 内存
table.exec.mini-batch.enabledtrueSQL 微批优化
table.exec.mini-batch.allow-latency1-5s微批延迟
parallelism.defaultKafka Partition 数的倍数避免数据倾斜

6.3 反压处理

// 监控反压
env.getConfig().setAutoWatermarkInterval(200);

// 反压常见原因与解决
// 1. 算子处理慢 → 增加并行度或优化逻辑
// 2. Sink 吞吐量不足 → 批量写入、异步化
// 3. 数据倾斜 → Key 加盐重新分区
// 4. GC 频繁 → 调整内存配置,使用 G1GC

总结

场景推荐方案
通用实时 ETLDataStream API + Event Time
窗口聚合Tumbling/Sliding + AggregateFunction
大状态容错RocksDBStateBackend + Incremental Checkpoint
Exactly-Once SinkKafka Producer with Transaction
迟到数据处理allowedLateness + sideOutputLateData
高反压场景Unaligned Checkpoint + 背压监控
SQL 类实时分析Flink SQL + Window TVF

7.1 三种 State Backend 详解

Flink 将算子状态持久化到 State Backend,直接决定作业的容量上限与容错恢复效率。生产环境应根据状态大小、延迟要求和硬件条件做选择。

特性MemoryStateBackendFsStateBackendEmbeddedRocksDBStateBackend
存储介质JVM HeapJVM Heap + 异步快照到文件系统本地 RocksDB(磁盘)
最大状态量小(受 TaskManager 堆内存限制)中(受堆内存限制,但快照落地快)大(可超出内存,磁盘可扩展)
快照方式同步快照到 JobManager 内存异步快照到文件系统(HDFS/S3)增量异步快照到文件系统
状态访问延迟极低(内存访问)极低(内存访问)毫秒级(磁盘/内存缓存混合)
适用场景本地开发、极小状态演示中小状态、低延迟敏感作业大状态、海量 Key、生产首选

7.2 增量 Checkpoint 配置

RocksDB 支持增量 Checkpoint,仅将新增 SST 文件与变更写入远端存储,可显著降低快照时间与网络负载。

// 生产推荐配置:RocksDB + 增量 Checkpoint
Configuration conf = new Configuration();

// 启用增量 Checkpoint
conf.setBoolean(StateBackendOptions.ASYNC_SNAPSHOTS, true);

EmbeddedRocksDBStateBackend rocksDb = new EmbeddedRocksDBStateBackend(true);
env.setStateBackend(rocksDb);
env.getCheckpointConfig().setCheckpointStorage("hdfs:///flink-checkpoints");

// RocksDB 内存调优
env.getConfig().set(ConfigConstants.ROCKSDB_MEMORY_MANAGED_KEY, "true");
env.getConfig().set(ConfigConstants.ROCKSDB_FIXED_PER_SLOT_MEMORY_SIZE, "256mb");

// 增量 Checkpoint 高级选项
env.getCheckpointConfig().enableUnalignedCheckpoints();
env.getCheckpointConfig().setAlignmentTimeout(Duration.ofSeconds(30));

7.3 State TTL 与清理策略

大状态场景必须配置 TTL,避免历史 Key 无限膨胀导致 OOM 或磁盘耗尽。

StateTtlConfig ttlConfig = StateTtlConfig
    .newBuilder(Time.hours(48))
    .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
    .setStateVisibility(StateTtlConfig.StateVisibility.ReturnExpiredIfNotCleanedUp)
    .cleanupIncrementally(10, true)
    .build();

ValueStateDescriptor<Long> descriptor = new ValueStateDescriptor<>("counter", Types.LONG);
descriptor.enableTimeToLive(ttlConfig);

8.1 Two-Phase Commit 原理

Flink 端到端 Exactly-Once 依赖两阶段提交(2PC)协调 Source 与 Sink 的事务边界。

Coordinator (JobManager)
  |
  +-- Phase 1: Pre-commit ---+
  |   Source 保存 offset 到状态     |
  |   Sink 开启事务并写入数据         |
  |   所有算子完成快照                |
  +---------------------------+
  |
  +-- Phase 2: Commit -------+
  |   Checkpoint 成功确认后         |
  |   Sink 提交事务(Kafka commit)   |
  |   Source 提交 offset            |
  +---------------------------+
  |
  +-- Abort (失败时) -----------+
     Sink 回滚事务,数据可重放

8.2 Kafka 事务性 Sink 配置

Properties producerProps = new Properties();
producerProps.put("bootstrap.servers", "kafka:9092");
producerProps.put("transactional.id", "flink-txn-sink-" + subtaskIndex);
producerProps.put("enable.idempotence", "true");
producerProps.put("acks", "all");
producerProps.put("retries", Integer.MAX_VALUE);
producerProps.put("max.in.flight.requests.per.connection", "5");

FlinkKafkaProducer<String> kafkaSink = new FlinkKafkaProducer<>(
    "output-events",
    new JsonSerializationSchema(),
    producerProps,
    FlinkKafkaProducer.Semantic.EXACTLY_ONCE
);

stream.addSink(kafkaSink);

8.3 端到端 Exactly-Once 完整配置

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

// Checkpoint 基础配置
env.enableCheckpointing(60000, CheckpointingMode.EXACTLY_ONCE);
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000);
env.getCheckpointConfig().setCheckpointTimeout(600000);
env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);
env.getCheckpointConfig().setExternalizedCheckpointCleanup(
    ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION
);

// 超时未对齐回退到 Aligned
env.getCheckpointConfig().enableUnalignedCheckpoints();
env.getCheckpointConfig().setAlignmentTimeout(Duration.ofSeconds(30));

// Source: Kafka 自动提交关闭,依赖 Checkpoint 保存 offset
Properties consumerProps = new Properties();
consumerProps.put("bootstrap.servers", "kafka:9092");
consumerProps.put("group.id", "flink-exactly-once-group");
consumerProps.put("auto.offset.reset", "earliest");
consumerProps.put("enable.auto.commit", "false"); // 必须关闭

FlinkKafkaConsumer<Event> kafkaSource = new FlinkKafkaConsumer<>(
    "input-events",
    new EventDeserializationSchema(),
    consumerProps
);
kafkaSource.setCommitOffsetsOnCheckpoints(true);

DataStream<Event> source = env.addSource(kafkaSource);
source.addSink(kafkaSink);

env.execute("ExactlyOncePipeline");

9.1 Pattern API 定义

CEP(Complex Event Processing)允许在流上定义事件模式,匹配复杂业务规则,如欺诈检测、异常告警。

Pattern<LoginEvent, LoginEvent> loginFailPattern = Pattern
    .<LoginEvent>begin("first")
        .where(evt -> evt.status.equals("FAIL"))
    .next("second")
        .where(evt -> evt.status.equals("FAIL"))
    .next("third")
        .where(evt -> evt.status.equals("FAIL"))
    .within(Time.minutes(5));

PatternStream<LoginEvent> patternStream = CEP.pattern(
    loginStream.keyBy(e -> e.userId),
    loginFailPattern
);

DataStream<Alert> alerts = patternStream.process(
    new PatternProcessFunction<LoginEvent, Alert>() {
        @Override
        public void processMatch(
            Map<String, List<LoginEvent>> match,
            Context ctx,
            Collector<Alert> out
        ) {
            LoginEvent first = match.get("first").get(0);
            out.collect(new Alert(
                first.userId,
                "TRIPLE_LOGIN_FAIL",
                ctx.timestamp()
            ));
        }
    }
);

9.2 超时与部分匹配处理

Pattern<PaymentEvent, PaymentEvent> paymentPattern = Pattern
    .<PaymentEvent>begin("create")
        .where(evt -> evt.type.equals("CREATE"))
    .followedBy("pay")
        .where(evt -> evt.type.equals("PAY"))
    .within(Time.minutes(10));

PatternStream<PaymentEvent> ps = CEP.pattern(
    paymentStream.keyBy(e -> e.orderId),
    paymentPattern
);

// 处理超时:订单创建后 10 分钟未支付
OutputTag<String> timeoutTag = new OutputTag<String>("timeout") {};

SingleOutputStreamOperator<OrderResult> result = ps.process(
    new PatternTimeoutHandler<OrderResult, String>()
);

// 超时侧输出流
DataStream<String> timeoutStream = result.getSideOutput(timeoutTag);
timeoutStream.addSink(new OrderTimeoutSink());

9.3 CEP 在金融风控场景的应用

// 场景:刷券套利检测 —— 同一设备在短时间内高频切换用户领取优惠券
Pattern<CouponEvent, CouponEvent> arbitragePattern = Pattern
    .<CouponEvent>begin("claim1")
        .where(evt -> evt.action.equals("CLAIM"))
    .next("claim2")
        .where(evt -> evt.action.equals("CLAIM"))
        .where(new SimpleCondition<CouponEvent>() {
            @Override
            public boolean filter(CouponEvent evt) {
                return !evt.userId.equals(prevUserId);
            }
        })
    .next("claim3")
        .where(evt -> evt.action.equals("CLAIM"))
    .within(Time.seconds(30));

// 匹配到模式后写入风控事件队列,触发人工审核
patternStream.process(new FraudAlertHandler()).addSink(new RiskControlSink());

10.1 流批一体 SQL 基础

Flink SQL 是流批一体查询的统一入口,同一条 SQL 既可以跑在流模式也可以跑在批模式。

-- 注册 Kafka Source 表
CREATE TABLE user_events (
    user_id STRING,
    event_type STRING,
    event_time TIMESTAMP(3),
    amount DECIMAL(10, 2),
    WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
    'connector' = 'kafka',
    'topic' = 'user-events',
    'properties.bootstrap.servers' = 'kafka:9092',
    'format' = 'json',
    'scan.startup.mode' = 'latest-offset'
);

-- 注册 Upsert Kafka Sink 表
CREATE TABLE event_summary (
    event_type STRING PRIMARY KEY NOT ENFORCED,
    total_amount DECIMAL(18, 2),
    event_count BIGINT
) WITH (
    'connector' = 'upsert-kafka',
    'topic' = 'event-summary',
    'properties.bootstrap.servers' = 'kafka:9092',
    'key.format' = 'json',
    'value.format' = 'json'
);

-- 流式聚合写入 Sink
INSERT INTO event_summary
SELECT event_type, SUM(amount) AS total_amount, COUNT(*) AS event_count
FROM user_events
GROUP BY event_type;

10.2 Temporal Table Join

Temporal Table Join 用于在流上关联维度表的最新版本,常用于关联 slowly-changing 维度。

-- 流表
CREATE TABLE orders (
    order_id STRING,
    currency STRING,
    amount DECIMAL(10, 2),
    order_time TIMESTAMP(3),
    WATERMARK FOR order_time AS order_time - INTERVAL '10' SECOND
) WITH ('connector' = 'kafka', ...);

-- 维度表(MySQL CDC 或 Lookup)
CREATE TABLE currency_rates (
    currency STRING,
    rate DECIMAL(10, 6),
    update_time TIMESTAMP(3),
    PRIMARY KEY (currency) NOT ENFORCED,
    WATERMARK FOR update_time AS update_time - INTERVAL '5' SECOND
) WITH (
    'connector' = 'jdbc',
    'url' = 'jdbc:mysql://mysql:3306/warehouse',
    'table-name' = 'currency_rates',
    'lookup.cache.max-rows' = '1000',
    'lookup.cache.ttl' = '10 min'
);

-- Temporal Join:关联该订单时间点的最新汇率
SELECT
    o.order_id,
    o.amount,
    o.amount * r.rate AS amount_usd
FROM orders AS o
LEFT JOIN currency_rates FOR SYSTEM_TIME AS OF o.order_time AS r
ON o.currency = r.currency;

10.3 Window TVF(Table-Valued Function)

Flink 1.13+ 推荐使用标准 SQL 的 WINDOW TVF 语法表达窗口聚合。

-- TUMBLE Window TVF
SELECT
    window_start,
    window_end,
    user_id,
    COUNT(*) AS event_count,
    SUM(amount) AS total_amount
FROM TABLE(
    TUMBLE(TABLE user_events, DESCRIPTOR(event_time), INTERVAL '10' MINUTES)
)
GROUP BY window_start, window_end, user_id;

-- CUMULATE(累积窗口):每 1 分钟输出一次过去 1 小时的累积结果
SELECT
    window_start,
    window_end,
    event_type,
    COUNT(*) AS cnt
FROM TABLE(
    CUMULATE(
        TABLE user_events,
        DESCRIPTOR(event_time),
        INTERVAL '1' MINUTES,
        INTERVAL '1' HOUR
    )
)
GROUP BY window_start, window_end, event_type;

10.4 SQL Client 使用

# 启动 Flink SQL Client
./bin/sql-client.sh embedded

# 或在会话中提交 SQL 文件
./bin/sql-client.sh -f /path/to/job.sql

# 设置执行模式
SET 'execution.runtime-mode' = 'streaming';
SET 'sql-client.execution.result-mode' = 'table';

# 查看执行计划
EXPLAIN PLAN FOR SELECT ... ;

11.1 Native Kubernetes 部署模式

Flink 1.10+ 支持 Native Kubernetes 部署,JobManager 直接与 K8s API Server 交互动态申请 TaskManager Pod。

apiVersion: flink.apache.org/v1beta1
kind: FlinkDeployment
metadata:
  name: realtime-etl
spec:
  image: flink:1.18-scala_2.12
  flinkVersion: v1.18
  jobManager:
    resource:
      memory: "2Gi"
      cpu: 1
  taskManager:
    resource:
      memory: "4Gi"
      cpu: 2
    replicas: 3
  job:
    jarURI: local:///opt/flink/examples/streaming/StateMachineExample.jar
    parallelism: 6
    upgradeMode: stateful
    state: running

11.2 TaskManager Pod 模板

apiVersion: v1
kind: Pod
metadata:
  name: task-manager-template
spec:
  containers:
    - name: flink-main-container
      resources:
        requests:
          memory: "4Gi"
          cpu: "2"
        limits:
          memory: "4Gi"
          cpu: "2"
      env:
        - name: FLINK_TM_HEAP_SIZE
          value: "3072m"
        - name: JVM_ARGS
          value: "-XX:+UseG1GC -XX:MaxGCPauseMillis=100"
      volumeMounts:
        - name: checkpoint-volume
          mountPath: /checkpoints
  volumes:
    - name: checkpoint-volume
      persistentVolumeClaim:
        claimName: flink-checkpoint-pvc

11.3 自动扩缩容配置

# Flink Autoscaler(基于自定义指标或 CPU / 背压)
apiVersion: autoscaling/v2
kind: HorizontalPodAutoscaler
metadata:
  name: flink-taskmanager-hpa
spec:
  scaleTargetRef:
    apiVersion: apps/v1
    kind: Deployment
    name: flink-taskmanager
  minReplicas: 2
  maxReplicas: 20
  metrics:
    - type: Resource
      resource:
        name: cpu
        target:
          type: Utilization
          averageUtilization: 70
    - type: Pods
      pods:
        metric:
          name: flink_taskmanager_job_task_backPressuredTimeMsPerSecond
        target:
          type: AverageValue
          averageValue: "200"
  behavior:
    scaleUp:
      stabilizationWindowSeconds: 60
      policies:
        - type: Percent
          value: 100
          periodSeconds: 60

12.1 背压诊断(Web UI Backpressure)

Flink Web UI 提供 Backpressure tabs,用颜色标识各 Subtask 的背压状态:

颜色状态含义处理建议
OK(绿色)无背压数据顺畅流动保持现状
LOW(黄色)轻度背压偶尔阻塞观察并优化慢算子
HIGH(红色)严重背压持续阻塞立即排查下游瓶颈

诊断流程:打开 Web UI -> Job -> Backpressure -> 逐层向上游定位瓶颈算子。

12.2 反压处理与网络缓冲区调优

// 网络缓冲区配置(flink-conf.yaml)
// taskmanager.memory.network.fraction: 0.15
// taskmanager.memory.network.max: 256mb
// taskmanager.memory.network.min: 128mb

// 在代码中动态调整缓冲区
Configuration netConf = new Configuration();
netConf.set(TaskManagerOptions.NETWORK_MEMORY_FRACTION, 0.15);
netConf.set(TaskManagerOptions.NETWORK_MEMORY_MAX, MemorySize.parse("256mb"));
netConf.set(TaskManagerOptions.NETWORK_MEMORY_MIN, MemorySize.parse("128mb"));

// 启用 Credit-based 流量控制(默认已开启,Flink 1.5+)
// 解决反压的另一种思路:增大下游并行度或批量写入

12.3 对象重用以减少 GC

// 开启对象链接重用(谨慎使用,仅当确认不会修改输入对象时)
env.getConfig().enableObjectReuse();

// 配合 ObjectReuse 的算子示例:避免在 FlatMap 中创建新对象
class OptimizedMapper extends RichFlatMapFunction<String, Event> {
    private final Event reuseEvent = new Event();

    @Override
    public void flatMap(String line, Collector<Event> out) {
        // 复用同一对象实例,减少 heap 分配
        reuseEvent.userId = parseUserId(line);
        reuseEvent.timestamp = parseTimestamp(line);
        reuseEvent.value = parseValue(line);
        out.collect(reuseEvent);
    }
}

// 设置 G1GC 降低 STW
// env.getConfig().addJVMOptions("-XX:+UseG1GC -XX:MaxGCPauseMillis=50");

13.1 Pulsar Source & Sink

Apache Pulsar 与 Flink 结合可构建企业级实时链路,Pulsar 提供分层存储与多租户隔离。

// Pulsar Source
PulsarSource<String> pulsarSource = PulsarSource.builder()
    .setServiceUrl("pulsar://pulsar:6650")
    .setAdminUrl("http://pulsar:8080")
    .setStartCursor(StartCursor.earliest())
    .setTopics("persistent://public/default/events")
    .setDeserializationSchema(new SimpleStringSchema())
    .setSubscriptionName("flink-consumer-sub")
    .setSubscriptionType(SubscriptionType.Key_Shared)
    .build();

DataStream<String> stream = env.fromSource(
    pulsarSource,
    WatermarkStrategy.noWatermarks(),
    "Pulsar Source"
);

// Pulsar Sink
PulsarSink<String> pulsarSink = PulsarSink.builder()
    .setServiceUrl("pulsar://pulsar:6650")
    .setAdminUrl("http://pulsar:8080")
    .setProducerConfig(
        PulsarSinkOptions.PULSAR_BATCHING_ENABLED,
        Boolean.TRUE
    )
    .setTopics("persistent://public/default/output")
    .setSerializationSchema(new SimpleStringSchema())
    .build();

stream.sinkTo(pulsarSink);

13.2 Key-Shared 订阅模式

Key-Shared 订阅保证同一 Key 的消息路由到同一消费者,天然与 Flink 的 KeyBy 语义对齐。

// Key-Shared 订阅确保顺序性与并行扩展同时满足
PulsarSource<Event> keyedSource = PulsarSource.builder()
    .setServiceUrl("pulsar://pulsar:6650")
    .setTopicsList(Arrays.asList("order-events"))
    .setDeserializationSchema(new EventSchema())
    .setSubscriptionName("flink-key-shared-sub")
    .setSubscriptionType(SubscriptionType.Key_Shared)
    // 允许粘性 Key 散列范围重新分配
    .setConfig(PulsarSourceOptions.PULSAR_ALLOW_TOPIC_LISTENER, true)
    .build();

13.3 Schema 演进

// Pulsar Schema 自动注册与演进
Schema<Event> eventSchema = Schema.AVRO(Event.class);

PulsarSource<Event> sourceWithSchema = PulsarSource.builder()
    .setServiceUrl("pulsar://pulsar:6650")
    .setTopics("events-topic")
    .setSchema(eventSchema)
    .setSubscriptionName("flink-schema-sub")
    .build();

// Schema 兼容性策略在 Pulsar 侧配置:FULL / BACKWARD / FORWARD

14. 生产故障案例

14.1 Checkpoint 超时排查

现象:Checkpoint 持续超时,状态越来越大。

排查步骤

  1. 查看 Web UI Checkpoint 页面,观察 Sync DurationAsync Duration
  2. Async Duration 过长 → 状态过大,确认是否开启增量 Checkpoint。
  3. Alignment Duration 过长 → 存在背压,启用 Unaligned Checkpoint。
  4. 检查 State Backend 网络带宽(HDFS/S3 上传速度)。
// 诊断性配置:输出 Checkpoint 统计到日志
env.getCheckpointConfig().enableCheckpointingIntervalLogging();

// 若 Checkpoint 超时,临时增大 timeout 并启用增量
checkpointConfig.setCheckpointTimeout(900000);
((EmbeddedRocksDBStateBackend) env.getStateBackend()).setIncrementalRestorePath(...);

14.2 OOM 调优

现象:TaskManager 频繁 OOMKilled 或 Full GC。

原因排查方式解决方案
状态过大未溢出Web UI 状态大小切换到 RocksDBStateBackend
网络缓冲区过高GC 日志分析调整 network.memory.fraction
窗口堆积检查窗口数量与水印推进设置 allowedLateness,清理过期窗口
用户代码内存泄漏Heap Dump 分析修复代码,避免无限集合增长

14.3 数据倾斜处理

// 方案 1:两阶段聚合(Local + Global)
DataStream<Result> preAggregated = stream
    .map(e -> Tuple2.of(randomPrefix(e.userId, 10), e.amount))
    .keyBy(t -> t.f0)
    .window(TumblingEventTimeWindows.of(Time.minutes(1)))
    .aggregate(new PreAggregate())
    .map(t -> Tuple2.of(removePrefix(t.f0), t.f1))
    .keyBy(t -> t.f0)
    .window(TumblingEventTimeWindows.of(Time.minutes(1)))
    .aggregate(new FinalAggregate());

// 方案 2:自定义分区器打散热点 Key
stream.partitionCustom(new HotKeyPartitioner(), e -> e.userId);

15. FAQ

Q1: Flink 的 Checkpoint 和 Savepoint 有什么区别?
Checkpoint 是自动周期触发、用于作业容错恢复的内部快照;Savepoint 是手动触发、用于版本升级或迁移的外部一致性快照,持久化存储不会自动清理。

Q2: 使用 RocksDBStateBackend 时为什么状态访问变慢了?
RocksDB 基于磁盘存储,访问需经过 LSM-Tree 查找,延迟高于纯内存。可通过启用 Block Cache、调优 state.backend.rocksdb.memory.managed 以及使用 SSD 来缓解。

Q3: Flink SQL 的流批一体如何切换执行模式?
设置 SET 'execution.runtime-mode' = 'batch';streaming;相同 SQL 逻辑在 Table API 中可通过 StreamTableEnvironmentBatchTableEnvironment 切换。

Q4: 为什么启用了 Exactly-Once 但下游仍有重复数据?
端到端 Exactly-Once 要求 Source 支持重放(如 Kafka offset)、Flink Checkpoint 保证内部状态一致,且 Sink 支持事务或幂等写入。若 Sink 非事务性且未做幂等处理,故障恢复时可能重复输出。

Q5: Flink on Kubernetes 中 TaskManager 重启后如何恢复状态?
Native K8s 模式支持通过 Checkpoint / Savepoint 路径自动恢复。在 FlinkDeployment 中配置 initialSavepointPath 或保留 checkpointDir,作业重启时指定 -s 参数即可从状态恢复。

总结

本文从 DataStream API 编程模型出发,系统性地覆盖了 Flink 的 Checkpoint 容错机制、Watermark 事件时间处理、窗口计算、Exactly-Once 语义实现与两阶段提交原理。随后深入 State Backend 选型与增量 Checkpoint 配置、CEP 复杂事件处理在金融风控中的实战、Flink SQL 的流批一体查询与 Window TVF 语法。还探讨了 Flink on Kubernetes 的 Native 部署模式、自动扩缩容、性能调优中的背压诊断与对象重用,以及 Pulsar 集成的 Key-Shared 订阅和 Schema 演进。最后通过生产故障案例与 FAQ 巩固常见问题的排查思路。掌握上述内容,即可在生产环境中稳定、高效地运行 Flink 实时计算作业。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「data-engineering」更多文章

  1. 数据工程深度指南:Modern Data Stack 全栈实践
  2. 数据平台工程:Data Mesh、FinOps 与 DataOps 生产实践
  3. Kafka Connect CDC 实战:Debezium 数据同步与变更捕获