一、Flink 基础概念

1.1 什么是 Flink

Apache Flink 是一个开源的分布式流处理框架,专门用于高效处理大规模数据流。它支持流处理和批处理两种模式,但核心是流处理。

主要特点:

  • 高吞吐、低延迟的流处理引擎
  • 支持事件时间(Event Time)处理,解决数据乱序问题
  • 提供精确一次(Exactly-Once)的语义保证
  • 支持有状态的计算
  • 灵活的窗口操作
  • 完善的容错机制

1.2 Flink 与其他框架的对比

  • Flink vs Spark Streaming:Flink 是真正的流处理(微批处理),Spark Streaming 是微批处理框架
  • Flink vs Kafka Streams:Flink 更强大灵活,支持更复杂的计算逻辑和更多状态管理能力
  • Flink vs Storm:Flink 提供更高的吞吐量和更强大的 API

1.3 Flink 架构中的核心概念

Source:数据源,负责读取外部数据(Kafka、文件、Socket 等)

Operator:算子,执行具体的计算逻辑(map、filter、reduce 等)

Sink:数据汇,负责输出结果到外部系统(Kafka、数据库、文件等)

Stream:数据流,无限的数据序列

Transformation:转换操作,将一个或多个 DataStream 转换为新的 DataStream

二、Flink 部署架构

2.1 系统架构

Flink 采用典型的 Master-Slave 架构:

JobManager(主节点)

  • 接收提交的作业
  • 负责作业的调度和协调
  • 管理 TaskManager 的生命周期
  • 存储和维护作业状态

TaskManager(从节点)

  • 执行任务的实际工作节点
  • 管理自己的资源(内存、CPU、网络)
  • 与 JobManager 定期心跳通信

Client

  • 用户提交作业的客户端
  • 将作业转换成 JobGraph 并提交给 JobManager

2.2 部署模式

Standalone 模式:独立部署,适合开发测试

Yarn 模式:运行在 Hadoop Yarn 之上,适合与其他大数据系统集成

Kubernetes 模式:容器化部署,适合云原生场景

Flink on K8s:现在推荐的部署方式

三、DataStream API

3.1 基本编程模型

// 1. 获取执行环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

// 2. 创建数据源(Source)
DataStream<String> stream = env.addSource(new FlinkKafkaConsumer<>(...));

// 3. 数据转换(Transformation)
stream.map(x -> x * 2)
      .filter(x -> x > 10)
      .print();

// 4. 数据输出(Sink)
stream.addSink(new FlinkKafkaProducer<>(...));

// 5. 提交执行
env.execute("Job Name");

3.2 常用转换算子

map:一对一转换,逐个元素映射

filter:过滤元素,保留满足条件的元素

flatMap:一对多转换,每个元素可以产生 0 个或多个结果

keyBy:按键分组,返回 KeyedStream,后续操作会分组执行

reduce:聚合操作,同一分组内的元素进行规约

aggregate:更灵活的聚合方式,可以使用自定义聚合逻辑

union:合并多个 DataStream

connect:连接两个 DataStream,可以进行联合处理

split / select:分流操作(旧版本,新版本用 side output)

3.3 Source 常见实现

Kafka Source

Properties props = new Properties();
props.setProperty("bootstrap.servers", "localhost:9092");
props.setProperty("group.id", "flink-group");
DataStream<String> stream = env.addSource(
    new FlinkKafkaConsumer<>("topic", new SimpleStringSchema(), props)
);

Socket Source(用于测试)

DataStream<String> stream = env.socketTextStream("localhost", 9999);

集合 Source

DataStream<Integer> stream = env.fromCollection(Arrays.asList(1, 2, 3, 4, 5));

自定义 Source:实现 SourceFunction 接口

3.4 Sink 常见实现

Print Sink:直接打印到控制台

stream.print();

Kafka Sink

stream.addSink(new FlinkKafkaProducer<>(
    "topic",
    new SimpleStringSchema(),
    props
));

自定义 Sink:实现 SinkFunction 接口

四、时间概念(Time Semantics)

4.1 三种时间

Processing Time(处理时间):事件被处理的时间,即系统时钟时间。简单但无法处理乱序数据。

Event Time(事件时间):事件发生的真实时间,即数据本身携带的时间戳。能准确处理乱序数据,是生产环境推荐方式。

Ingestion Time(摄入时间):事件进入 Flink 系统的时间。介于两者之间,使用较少。

4.2 配置时间语义

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 配置使用事件时间
env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime);
// 配置最大允许延迟
env.getConfig().setAutoWatermarkInterval(1000);

4.3 Watermark(水位线)

Watermark 用于衡量事件时间的进度,标记"在这个时间点之前的数据都已经到达"。

作用

  • 触发窗口计算
  • 处理迟到数据
  • 确保计算的完整性

生成 Watermark

DataStream<Event> stream = env.addSource(...)
    .assignTimestampsAndWatermarks(
        WatermarkStrategy
            .forBoundedOutOfOrderness(Duration.ofSeconds(20))
            .withTimestampAssigner((event, recordTimestamp) -> event.getTimestamp())
    );

allowedLateness:指定允许迟到数据的最大时间

stream.keyBy(...)
      .window(TumblingEventTimeWindows.of(Time.seconds(10)))
      .allowedLateness(Time.minutes(1))
      .process(...)

五、窗口(Window)

5.1 窗口的概念

窗口是处理无限数据流的关键机制,将无限数据流分割成有限的数据集合。

5.2 窗口类型

时间窗口

  • 滚动窗口(Tumbling Window):固定大小,无重叠

    .window(TumblingEventTimeWindows.of(Time.seconds(10)))
    
  • 滑动窗口(Sliding Window):固定大小,有重叠

    .window(SlidingEventTimeWindows.of(Time.seconds(10), Time.seconds(5)))
    
  • 会话窗口(Session Window):根据数据活跃度动态生成

    .window(EventTimeSessionWindows.withGap(Time.seconds(5)))
    

计数窗口

  • 滚动计数窗口

    .countWindow(100)
    
  • 滑动计数窗口

    .countWindow(100, 50)
    

5.3 窗口计算

stream.keyBy(x -> x.userId)
      .window(TumblingEventTimeWindows.of(Time.seconds(10)))
      .reduce((x, y) -> new Event(x.userId, x.count + y.count))
      .print();

或使用 process() 获得更多控制权:

stream.keyBy(...)
      .window(...)
      .process(new ProcessWindowFunction<Event, Result, String, TimeWindow>() {
          @Override
          public void process(String key, Context context, Iterable<Event> elements, Collector<Result> out) {
              // 自定义窗口计算逻辑
          }
      })

六、状态管理(State)

6.1 状态概念

状态是指在流处理过程中维护的信息,用于存储计算的中间结果。

6.2 状态类型

Operator State(算子状态)

  • 绑定到单个算子实例
  • 分布式场景下不推荐使用
  • 应用场景有限

Keyed State(键值状态)

  • 按键分组存储
  • 每个键对应一份状态
  • 生产环境推荐使用

6.3 Keyed State 类型

ValueState:存储单个值

ValueState<Integer> countState = getRuntimeContext()
    .getState(new ValueStateDescriptor<>("count", Integer.class));

ListState:存储列表

ListState<String> listState = getRuntimeContext()
    .getListState(new ListStateDescriptor<>("list", String.class));

MapState:存储键值对

MapState<String, Integer> mapState = getRuntimeContext()
    .getMapState(new MapStateDescriptor<>("map", String.class, Integer.class));

ReducingStateAggregatingState:专门用于聚合

6.4 有状态计算示例

public class StatefulCounter extends KeyedProcessFunction<String, Event, Result> {
    private transient ValueState<Integer> countState;
    
    @Override
    public void open(Configuration parameters) throws Exception {
        ValueStateDescriptor<Integer> descriptor = 
            new ValueStateDescriptor<>("count", Integer.class);
        countState = getRuntimeContext().getState(descriptor);
    }
    
    @Override
    public void processElement(Event event, Context ctx, Collector<Result> out) 
            throws Exception {
        Integer current = countState.value();
        current = current == null ? 1 : current + 1;
        countState.update(current);
        out.collect(new Result(event.getKey(), current));
    }
}

七、容错机制

7.1 Checkpoint

Checkpoint 是 Flink 故障恢复的核心机制,定期保存整个应用的状态快照。

启用 Checkpoint

env.enableCheckpointing(60000); // 每 60 秒进行一次 checkpoint
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000);
env.getCheckpointConfig().setCheckpointTimeout(600000);
env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);

Checkpoint 模式

  • EXACTLY_ONCE:精确一次语义,性能略低但最安全
  • AT_LEAST_ONCE:至少一次语义,性能更好但可能重复处理

7.2 State Backend

状态后端用于存储状态数据。

MemoryStateBackend:存储在内存中,适合开发测试

env.setStateBackend(new MemoryStateBackend());

FsStateBackend:存储在文件系统中,支持大状态

env.setStateBackend(new FsStateBackend("hdfs:///flink/checkpoints"));

RocksDBStateBackend:基于 RocksDB,支持最大的状态量

env.setStateBackend(new RocksDBStateBackend("hdfs:///flink/checkpoints"));

7.3 Savepoint

Savepoint 是人工触发的状态快照,用于应用版本升级、迁移或维护。

# 触发 savepoint
flink savepoint <job-id> [target-directory]

# 从 savepoint 恢复
flink run -s <savepoint-path> <jar-path>

八、CEP(复杂事件处理)

8.1 什么是 CEP

CEP 用于在事件流中检测和提取复杂的事件模式。

8.2 基本用法

DataStream<Event> input = env.addSource(...);

Pattern<Event, ?> pattern = Pattern
    .<Event>begin("start").where(e -> e.getId() == 1)
    .next("middle").where(e -> e.getId() == 2)
    .followedBy("end").where(e -> e.getId() == 3)
    .within(Time.seconds(10));

PatternStream<Event> patternStream = CEP.pattern(input, pattern);

DataStream<String> result = patternStream.process(
    new PatternProcessFunction<Event, String>() {
        @Override
        public void processMatch(Map<String, List<Event>> pattern, 
                               Context ctx, 
                               Collector<String> out) {
            out.collect("Pattern matched!");
        }
    }
);

九、高级特性

9.1 Side Output(旁路输出)

用于将数据分为多个输出流,例如将异常数据和正常数据分开处理。

OutputTag<String> outputTag = new OutputTag<String>("side-output"){};

SingleOutputStreamOperator<Integer> mainStream = stream
    .process(new ProcessFunction<String, Integer>() {
        @Override
        public void processElement(String value, Context ctx, Collector<Integer> out) {
            if (value.startsWith("ERROR")) {
                ctx.output(outputTag, value);
            } else {
                out.collect(Integer.parseInt(value));
            }
        }
    });

DataStream<String> sideStream = mainStream.getSideOutput(outputTag);

9.2 Broadcast State

广播状态用于将某个流的数据广播到所有并行任务。

MapStateDescriptor<String, String> descriptor = 
    new MapStateDescriptor<>("broadcast", String.class, String.class);
BroadcastStream<String> broadcastStream = configStream.broadcast(descriptor);

stream.connect(broadcastStream)
      .process(new BroadcastProcessFunction<Event, String, Result>() {
          @Override
          public void processElement(Event event, ReadOnlyContext ctx, Collector<Result> out) {
              ReadOnlyBroadcastState<String, String> state = ctx.getBroadcastState(descriptor);
              // 使用广播状态
          }
          
          @Override
          public void processBroadcastElement(String element, Context ctx, Collector<Result> out) {
              BroadcastState<String, String> state = ctx.getBroadcastState(descriptor);
              state.put("key", element);
          }
      });

9.3 异步函数(Async I/O)

用于异步调用外部服务,提高吞吐量。

AsyncDataStream.unorderedWait(
    stream,
    new AsyncFunction<Event, Result>() {
        @Override
        public void asyncInvoke(Event input, ResultFuture<Result> resultFuture) {
            // 异步调用外部服务
            CompletableFuture.supplyAsync(() -> {
                return callExternalService(input);
            }).thenAccept(resultFuture::complete);
        }
    },
    1000,
    TimeUnit.MILLISECONDS
);

十、调优和最佳实践

10.1 性能调优

并行度设置:与数据源分区数相匹配

env.setParallelism(8);

减少序列化开销:使用 Kryo 序列化器

env.getConfig().enableForceKryo();

状态大小控制:定期清理过期状态,使用 TTL

StateTtlConfig ttlConfig = StateTtlConfig
    .newBuilder(Time.days(1))
    .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
    .build();
descriptor.enableTimeToLive(ttlConfig);

10.2 最佳实践

  • 使用事件时间而非处理时间,保证计算的准确性
  • 正确配置 Watermark,平衡延迟和数据完整性
  • 启用 Checkpoint,确保故障恢复能力
  • 监控作业状态,及时发现和解决问题
  • 合理使用状态,避免状态过大导致性能下降
  • 充分测试,特别是复杂的业务逻辑

十一、常见问题排查

数据处理延迟高

  • 检查并行度设置是否合理
  • 检查状态大小是否过大
  • 优化 Checkpoint 配置

数据丢失

  • 确保启用了 Checkpoint
  • 检查 Source 和 Sink 的容错能力
  • 验证数据源的可靠性

内存溢出

  • 减少并行度或增加 TaskManager 内存
  • 检查是否存在状态泄漏
  • 优化数据结构

更多推荐