跳到正文
SL Blog 技术探索 · 工程实践 · AI 时代思考
返回
🌊 Flink 教程 · 4 / 7 查看系列简介 →

Flink 教程(四):Java 开发实践

文章目录
  1. 4.1 环境配置
  2. 4.2 DataStream API 基础
  3. 4.2.1 创建执行环境
  4. 4.2.2 Lambda 表达式支持
  5. 4.3 状态管理
  6. 4.3.1 Keyed State 使用
  7. 4.3.2 Operator State 使用
  8. 4.4 窗口操作
  9. 4.4.1 滚动窗口(Tumbling Window)
  10. 4.4.2 滑动窗口(Sliding Window)
  11. 4.4.3 会话窗口(Session Window)
  12. 4.4.4 窗口函数
  13. 4.5 时间语义和水位线
  14. 4.6 容错配置
  15. 4.7 数据源与接收器
  16. 4.8 Process Function
  17. 4.9 旁路输出(Side Output)
  18. 4.10 Kafka 连接器
  19. 4.10.1 Kafka 数据源(FlinkKafkaConsumer)
  20. 4.10.2 Kafka 接收器(FlinkKafkaProducer)
  21. 4.10.3 Kafka 连接器配置
  22. 4.11 JDBC 连接器
  23. 4.12 复杂事件处理(CEP)
  24. 4.12.1 添加 Maven 依赖
  25. 4.12.2 定义模式
  26. 4.12.3 循环模式
  27. 4.12.4 连续性策略
  28. 4.12.5 条件
  29. 4.12.6 检测模式并处理匹配
  30. 4.13 图计算库(Gelly)
  31. 4.13.1 添加 Maven 依赖
  32. 4.13.2 创建图
  33. 4.13.3 图转换
  34. 4.13.4 图算法
  35. 4.13.5 图生成器
  36. 4.14 大状态调优
  37. 4.14.1 Checkpoint 调优
  38. 4.14.2 RocksDB 调优
  39. 4.14.3 内存管理
  40. 4.14.4 压缩配置
  41. 4.14.5 任务本地恢复
  42. 4.15 Savepoint 操作

📌 本文是「Flink 教程」系列第 四 篇 · 系列目录

4.1 环境配置

Maven 依赖:

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-java</artifactId>
    <version>1.14.4</version>
</dependency>
<dependency>
     <groupId>org.apache.flink</groupId>
    <artifactId>flink-streaming-java_2.12</artifactId>
    <version>1.14.4</version>
</dependency>
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-cep_2.11</artifactId>
    <version>1.14.4</version>
</dependency>
<!-- Flink Gelly 已弃用,Flink 1.14+ 不推荐使用
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-gelly_2.11</artifactId>
    <version>1.14.4</version>
</dependency>
-->

4.2 DataStream API 基础

4.2.1 创建执行环境

import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;

public class WordCount {
    public static void main(String[] args) throws Exception {
        // 创建执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        // 设置并行度
        env.setParallelism(1);

        // 从集合创建数据源
        DataStream<String> textDataStream = env.fromElements(
            "Hello world",
            "Hello Flink",
            "Hello Java"
        );

        // 执行转换
        DataStream<Tuple2<String, Integer>> wordCounts = textDataStream
            .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(0)
            .sum(1);

        // 输出结果
        wordCounts.print();

        // 执行作业
        env.execute("WordCount Example");
    }
}

4.2.2 Lambda 表达式支持

简单 Lambda(类型可自动推断):

env.fromElements(1, 2, 3)
   .map(i -> i * i)
   .print();
// 输出: 1, 4, 9

带泛型的 Lambda(需要显式类型声明):

// 使用 flatMap 时需要指定 Collector 类型
input.flatMap((Integer number, Collector<String> out) -> {
    for (int i = 0; i < number; i++) {
        out.collect("a".repeat(i + 1));
    }
})
.returns(Types.STRING)
.print();

4.3 状态管理

4.3.1 Keyed State 使用

public class CountWindowAverage extends RichFlatMapFunction<Tuple2<Long, Long>, Tuple2<Long, Long>> {

    // 声明状态变量
    private transient ValueState<Tuple2<Long, Long>> sum;

    @Override
    public void flatMap(Tuple2<Long, Long> input, Collector<Tuple2<Long, Long>> out) throws Exception {
        // 获取当前状态
        Tuple2<Long, Long> currentSum = sum.value();
        if (currentSum == null) {
            currentSum = Tuple2.of(0L, 0L);
        }

        // 更新状态
        currentSum.f0 += 1;  // 计数
        currentSum.f1 += input.f1;  // 求和
        sum.update(currentSum);

        // 每 2 个元素输出一次
        if (currentSum.f0 >= 2) {
            out.collect(Tuple2.of(input.f0, currentSum.f1 / currentSum.f0));
            sum.clear();
        }
    }

    @Override
    public void open(Configuration config) {
        // 初始化状态描述符
        ValueStateDescriptor<Tuple2<Long, Long>> descriptor =
            new ValueStateDescriptor<>(
                "average",
                TypeInformation.of(new TypeHint<Tuple2<Long, Long>>() {}));

        // 获取运行时上下文中的状态
        sum = getRuntimeContext().getState(descriptor);
    }
}

4.3.2 Operator State 使用

public class BufferingSink implements SinkFunction<Tuple2<String, Integer>>,
                                      CheckpointedFunction {

    private final int threshold;
    private transient ListState<Tuple2<String, Integer>> checkpointedState;
    private List<Tuple2<String, Integer>> bufferedElements;

    public BufferingSink(int threshold) {
        this.threshold = threshold;
        this.bufferedElements = new ArrayList<>();
    }

    @Override
    public void invoke(Tuple2<String, Integer> value, Context context) throws Exception {
        bufferedElements.add(value);
        if (bufferedElements.size() >= threshold) {
            for (Tuple2<String, Integer> element : bufferedElements) {
                // 输出到外部系统
            }
            bufferedElements.clear();
        }
    }

    @Override
    public void snapshotState(FunctionSnapshotContext context) throws Exception {
        checkpointedState.clear();
        for (Tuple2<String, Integer> element : bufferedElements) {
            checkpointedState.add(element);
        }
    }

    @Override
    public void initializeState(FunctionInitializationContext context) throws Exception {
        ListStateDescriptor<Tuple2<String, Integer>> descriptor =
            new ListStateDescriptor<>(
                "buffered-elements",
                TypeInformation.of(new TypeHint<Tuple2<String, Integer>>() {}));

        checkpointedState = context.getOperatorStateStore().getListState(descriptor);

        // 从 Checkpoint 恢复
        if (context.isRestored()) {
            for (Tuple2<String, Integer> element : checkpointedState.get()) {
                bufferedElements.add(element);
            }
        }
    }
}

4.4 窗口操作

4.4.1 滚动窗口(Tumbling Window)

// 基于时间的滚动窗口
DataStream<Tuple2<String, Integer>> wordCounts = input
    .keyBy(t -> t.f0)
    .window(TumblingProcessingTimeWindows.of(Time.seconds(5)))
    .sum(1);

// 基于计数的滚动窗口
DataStream<Tuple2<String, Integer>> counts = input
    .keyBy(t -> t.f0)
    .countWindow(100)  // 每 100 个元素一个窗口
    .sum(1);

4.4.2 滑动窗口(Sliding Window)

// 窗口大小 10 秒,滑动步长 5 秒
DataStream<Tuple2<String, Integer>> wordCounts = input
    .keyBy(t -> t.f0)
    .window(SlidingProcessingTimeWindows.of(Time.seconds(10), Time.seconds(5)))
    .sum(1);

4.4.3 会话窗口(Session Window)

⚠️ 未实现:当前代码库中未包含会话窗口的示例代码

// 待实现
DataStream<Tuple2<String, Integer>> wordCounts = input
    .keyBy(t -> t.f0)
    .window(ProcessingTimeSessionWindows.withGap(Time.seconds(10)))
    .sum(1);

4.4.4 窗口函数

ProcessWindowFunction(全量窗口函数):

input
    .keyBy(x -> x.key)
    .window(TumblingEventTimeWindows.of(Time.minutes(1)))
    .process(new ProcessWindowFunction<...>() {
        @Override
        public void process(String key, Context context,
                          Iterable<...> values, Collector<...> out) {
            int count = 0;
            for (... value : values) {
                count += value;
            }
            out.collect(Tuple3.of(key, context.window().getEnd(), count));
        }
    });

增量聚合函数(ReduceFunction):

input
    .keyBy(...)
    .window(...)
    .reduce((value1, value2) -> Tuple2.of(value1.f0, value1.f1 + value2.f1));

4.5 时间语义和水位线

DataStream<Event> stream = ...;

WatermarkStrategy<Event> strategy = WatermarkStrategy
    .<Event>forBoundedOutOfOrderness(Duration.ofSeconds(20))
    .withTimestampAssigner((event, timestamp) -> event.timestamp);

DataStream<Event> withWatermarks = stream.assignTimestampsAndWatermarks(strategy);

💡 Tips: 窗口触发机制与延迟处理

窗口触发时机:

  • 事件时间窗口:当 Watermark > 窗口结束时间时触发
  • 处理时间窗口:当系统时间 > 窗口结束时间时触发
  • 使用 allowedLateness() 设置允许的迟到时间
  • 使用 sideOutputLateData() 收集迟到数据

窗口类型对比:

类型特点适用场景
滚动窗口窗口不重叠、固定大小周期性统计
滑动窗口窗口可重叠、步长可调移动平均、趋势分析
会话窗口根据活动间隙动态创建用户会话分析

迟到数据处理最佳实践:

stream.keyBy(...)
    .window(TumblingEventTimeWindows.of(Time.seconds(5)))
    .allowedLateness(Time.seconds(10))      // 允许迟到 10 秒
    .sideOutputLateData(lateOutputTag)       // 收集迟到数据
    .process(...)

4.6 容错配置

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

// 启用 Checkpoint
env.enableCheckpointing(60000); // 60 秒间隔

// Checkpoint 配置
CheckpointConfig config = env.getCheckpointConfig();
config.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
config.setCheckpointTimeout(60000);
config.setMinPauseBetweenCheckpoints(30);

// 任务取消时保留 Checkpoint
config.setExternalizedCheckpointCleanup(
    CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);

// 启用非对齐 Checkpoint
config.enableUnalignedCheckpoints();

// 配置状态后端
env.setStateBackend(new FsStateBackend("hdfs://namenode:40010/flink/checkpoints"));

4.7 数据源与接收器

常用数据源:

// 从集合创建
env.fromCollection(Collection)

// 从元素创建
env.fromElements(1, 2, 3)

// 读取文件
env.readTextFile("path/to/file")

// Socket 数据源
env.socketTextStream("localhost", 9999)

// Kafka 数据源
FlinkKafkaConsumer<String> consumer = new FlinkKafkaConsumer<>(
    "topic",
    new SimpleStringSchema(),
    properties);
env.addSource(consumer);

常用接收器:

// 打印到控制台
stream.print();

// 写入文件
stream.writeAsText("path/to/output");

// Kafka 接收器
stream.addSink(new FlinkKafkaProducer<>(
    "topic",
    new SimpleStringSchema(),
    properties));

💡 Tips: 数据源 Exactly-Once 限制

使用文件作为数据源时,需要注意以下限制:

模式Exactly-Once 支持说明
PROCESS_CONTINUOUSLY❌ 不支持文件内容变化时重新加载全部内容,无法保证精确一次
PROCESS_ONCE✅ 支持只读取变化的数据,可以保证精确一次

注意: 如果使用文件作为数据源,当某个节点异常停止时,Checkpoints 不会更新。如果数据持续生成,可能导致该节点数据积压,需要较长时间从最新的 Checkpoint 恢复。

Kafka 精确一次配置:

 // Kafka Source - 从 Checkpoint 恢复消费位置
 consumer.setCommitOffsetsOnCheckpoints(true);

 // Kafka Sink - 使用事务
 FlinkKafkaProducer<String> producer = new FlinkKafkaProducer<>(
     "output-topic",
     new SimpleStringSchema(),
     properties,
     FlinkKafkaProducer.Semantic.EXACTLY_ONCE,  // 关键:使用事务
     "producer-id"
 );

4.8 Process Function

// KeyedProcessFunction
input
    .keyBy(...)
    .process(new KeyedProcessFunction<K, IN, OUT>() {
        @Override
        public void processElement(IN value, Context ctx, Collector<OUT> out) {
            // 处理每个元素
            long currentTime = ctx.timerService().currentProcessingTime();
            ctx.timerService().registerProcessingTimeTimer(currentTime + 5000);
        }

        @Override
        public void onTimer(long timestamp, OnTimerContext ctx, Collector<OUT> out) {
            // 定时器触发时的逻辑
        }
    });

4.9 旁路输出(Side Output)

// 定义旁路输出标签
final OutputTag<String> outputTag = new OutputTag<String>("side-output") {};

// 使用 ProcessFunction 处理
SingleOutputStreamOperator<Integer> mainDataStream = input
    .process(new ProcessFunction<Integer, Integer>() {
        @Override
        public void processElement(Integer value, Context ctx, Collector<Integer> out) {
            out.collect(value);
            ctx.output(outputTag, "sideout-" + value);
        }
    });

// 获取旁路输出流
DataStream<String> sideOutputStream = mainDataStream.getSideOutput(outputTag);

4.10 Kafka 连接器

4.10.1 Kafka 数据源(FlinkKafkaConsumer)

import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer;
import org.apache.flink.api.common.serialization.SimpleStringSchema;

Properties properties = new Properties();
properties.setProperty("bootstrap.servers", "localhost:9092");
properties.setProperty("group.id", "my-group");

FlinkKafkaConsumer<String> consumer = new FlinkKafkaConsumer<>(
    "my-topic",
    new SimpleStringSchema(),
    properties);

// 从最早的记录开始消费
consumer.setStartFromEarliest();

// 从最新的记录开始消费
// consumer.setStartFromLatest();

// 从指定的 offset 开始消费
// consumer.setStartFromOffset(100L);

// 从指定的时间戳开始消费
// consumer.setStartFromTimestamp(1609459200000L);

DataStream<String> stream = env.addSource(consumer);

4.10.2 Kafka 接收器(FlinkKafkaProducer)

import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer;
import org.apache.flink.api.common.serialization.SimpleStringSchema;

// 精确一次语义(需要事务支持)
FlinkKafkaProducer<String> producer = new FlinkKafkaProducer<>(
    "my-output-topic",
    new SimpleStringSchema(),
    properties,
    FlinkKafkaProducer.Semantic.EXACTLY_ONCE,
    "my-producer-id"
);

// 至少一次语义
// FlinkKafkaProducer<String> producer = new FlinkKafkaProducer<>(
//     "my-output-topic",
//     new SimpleStringSchema(),
//     properties,
//     FlinkKafkaProducer.Semantic.AT_LEAST_ONCE,
//     "my-producer-id"
// );

stream.addSink(producer);

4.10.3 Kafka 连接器配置

// 完整的 Kafka 连接器配置
FlinkKafkaConsumer<String> consumer = new FlinkKafkaConsumer<>(
    "topic",
    new SimpleStringSchema(),
    properties);

// 启用消费者自动分区发现
properties.setProperty("flink.partition-discovery.interval", "5000");

// 配置超时时间
properties.setProperty("request.timeout.ms", "30000");
properties.setProperty("session.timeout.ms", "30000");

// 启用事务 Producer 的幂等性(精确一次)
properties.setProperty("enable.idempotence", "true");

4.11 JDBC 连接器

import org.apache.flink.streaming.connectors.jdbc.JdbcConnectionOptions;
import org.apache.flink.streaming.connectors.jdbc.JdbcExecutionOptions;
import org.apache.flink.streaming.connectors.jdbc.JdbcSink;

// 基本的 JDBC Sink
env.fromElements(...)
    .addSink(JdbcSink.sink(
        "INSERT INTO users (id, name, age) VALUES (?, ?, ?)",
        (statement, user) -> {
            statement.setInt(1, user.id);
            statement.setString(2, user.name);
            statement.setInt(3, user.age);
        },
        new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
            .withUrl("jdbc:mysql://localhost:3306/mydb")
            .withDriverName("com.mysql.cj.jdbc.Driver")
            .withUsername("root")
            .withPassword("password")
            .build()
    ));

// 带有执行选项的 JDBC Sink
env.fromElements(...)
    .addSink(JdbcSink.sink(
        "INSERT INTO orders (id, product, quantity) VALUES (?, ?, ?)",
        (statement, order) -> {
            statement.setInt(1, order.id);
            statement.setString(2, order.product);
            statement.setInt(3, order.quantity);
        },
        JdbcExecutionOptions.builder()
            .withBatchSize(1000)           // 每 1000 条批量写入
            .withBatchIntervalMs(200)      // 批量间隔 200ms
            .withMaxRetries(5)             // 最大重试次数
            .build(),
        new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
            .withUrl("jdbc:mysql://localhost:3306/mydb")
            .withDriverName("com.mysql.cj.jdbc.Driver")
            .build()
    ));

4.12 复杂事件处理(CEP)

FlinkCEP 是 Flink 上层的复杂事件处理库,用于在无限事件流中检测出特定的事件模式。

CEP 处理流程:

flowchart LR
    subgraph "事件流"
        E1[事件 A] --> E2[事件 B] --> E3[事件 C] --> E4[事件 D] --> E5[事件 E]
    end

    subgraph "模式匹配"
        P["Pattern: A → B → C<br/>within 10s"]
    end

    subgraph "匹配结果"
        R[匹配到完整模式<br/>触发处理]
    end

    E1 --> P
    E2 --> P
    E3 --> P
    E4 -.->|"不匹配"| P
    E5 -.->|"不匹配"| P

    P -->|"匹配成功"| R

    style P fill:#e3f2fd
    style R fill:#c8e6c9

4.12.1 添加 Maven 依赖

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-cep_2.11</artifactId>
    <version>1.14.4</version>
</dependency>

4.12.2 定义模式

import org.apache.flink.cep.pattern.Pattern;
import org.apache.flink.cep.pattern.conditions.SimpleCondition;

// 定义模式:start -> middle -> end
Pattern<Event, ?> pattern = Pattern.<Event>begin("start")
    .where(new SimpleCondition<Event>() {
        @Override
        public boolean filter(Event event) {
            return event.getId() == 42;  // 匹配 id 为 42 的事件
        }
    })
    .next("middle")  // 严格连续
    .subtype(SubEvent.class)  // 限制事件类型
    .where(new SimpleCondition<SubEvent>() {
        @Override
        public boolean filter(SubEvent subEvent) {
            return subEvent.getVolume() >= 10.0;
        }
    })
    .followedBy("end")  // 松散连续
    .where(new SimpleCondition<Event>() {
        @Override
        public boolean filter(Event event) {
            return event.getName().equals("end");
        }
    })
    .within(Time.seconds(10));  // 10 秒内完成

4.12.3 循环模式

// 期望出现 4 次
Pattern.<Event>begin("start").times(4);

// 期望出现 2-4 次
Pattern.<Event>begin("start").times(2, 4);

// 期望出现 1 到多次
Pattern.<Event>begin("start").oneOrMore();

// 期望出现 0 到多次(可选)
Pattern.<Event>begin("start").oneOrMore().optional();

// 期望出现至少 2 次
Pattern.<Event>begin("start").timesOrMore(2);

// 贪婪模式(尽可能多的匹配)
Pattern.<Event>begin("start").times(2, 4).greedy();

// 严格连续(循环中)
Pattern.<Event>begin("start").followedBy("middle").oneOrMore().consecutive();

4.12.4 连续性策略

// 严格连续:事件必须一个接一个出现,中间不能有其他事件
Pattern.<Event>begin("start").next("middle");

// 松散连续:忽略匹配事件之间的不匹配事件
Pattern.<Event>begin("start").followedBy("middle");

// 不确定松散连续:允许忽略匹配事件的附加匹配
Pattern.<Event>begin("start").followedByAny("middle");

// 在循环模式中使用严格连续
Pattern.<Event>begin("start").followedBy("middle").oneOrMore().consecutive();

4.12.5 条件

// 简单条件:只取决于事件自身属性
Pattern.<Event>begin("start").where(event -> event.getAmount() > 100);

// 迭代条件:可以访问之前匹配的事件
Pattern.<Event>begin("middle").where(
    new IterativeCondition<SubEvent>() {
        @Override
        public boolean filter(SubEvent value, Context<SubEvent> ctx) throws Exception {
            if (!value.getName().startsWith("foo")) {
                return false;
            }
            double sum = value.getPrice();
            for (Event event : ctx.getEventsForPattern("middle")) {
                sum += event.getPrice();
            }
            return sum < 5.0;
        }
    }
);

// 停止条件(用于循环模式)
Pattern.<Event>begin("start")
    .oneOrMore()
    .until(new SimpleCondition<Event>() {
        @Override
        public boolean filter(Event event) {
            return event.getStatus().equals("stop");
        }
    });

4.12.6 检测模式并处理匹配

import org.apache.flink.cep.PatternStream;
import org.apache.flink.cep.functions.PatternProcessFunction;
import org.apache.flink.util.Collector;

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

DataStream<Alert> result = patternStream.process(
    new PatternProcessFunction<Event, Alert>() {
        @Override
        public void processMatch(
                Map<String, List<Event>> pattern,
                Context ctx,
                Collector<Alert> out) throws Exception {
            Event start = pattern.get("start").get(0);
            List<Event> middle = pattern.get("middle");
            Event end = pattern.get("end").get(0);

            out.collect(createAlert(start, middle, end));
        }
    });

// 处理超时部分匹配
DataStream<Alert> result = patternStream.process(
    new PatternProcessFunction<Event, Alert>() {
        @Override
        public void processMatch(
                Map<String, List<Event>> pattern,
                Context ctx,
                Collector<Alert> out) throws Exception {
            // 正常匹配处理
        }

        @Override
        public void processTimedOutMatch(
                Map<String, List<Event>> pattern,
                Context ctx,
                Collector<Alert> out) throws Exception {
            // 超时部分匹配处理
        }
    });

4.13 图计算库(Gelly)

⚠️ 未实现:Flink Gelly 在 Flink 1.14+ 版本已弃用,当前代码库未包含相关示例

Gelly 是 Flink 的图计算 API,提供了一套简化和加速 Flink 图分析应用开发的方法和工具。

4.13.1 添加 Maven 依赖

<!-- Java 版 -->
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-gelly_2.11</artifactId>
    <version>1.14.4</version>
</dependency>

<!-- Scala 版 -->
<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-gelly-scala_2.11</artifactId>
    <version>1.14.4</version>
</dependency>

4.13.2 创建图

import org.apache.flink.graph.Graph;
import org.apache.flink.graph.vertexData.LongVertexValue;
import org.apache.flink.graph.edgeData.LongEdgeValue;

// 从边列表创建图
Graph<Long, Long, Long> graph = Graph.fromEdges(
    edges,
    new LongVertexValueInitializer(0L)
);

// 从顶点和边创建图
Graph<Long, Long, Long> graph = Graph.fromDataSet(
    vertices,
    edges,
    new LongVertexValueInitializer(0L)
);

// 从 CSV 文件创建图
Graph<Long, Long, Long> graph = Graph.fromCsvReader(
    "/path/to/vertices.csv",
    "/path/to/edges.csv",
    env
);

4.13.3 图转换

// 映射顶点值
Graph<Long, Double, Long> mappedGraph = graph.mapVertices(
    new MapFunction<Vertex<Long, Long>, Double>() {
        @Override
        public Vertex<Long, Double> map(Vertex<Long, Long> value) {
            return new Vertex<>(value.getId(), value.getValue() * 2.0);
        }
    }
);

// 映射边值
Graph<Long, Long, Double> mappedGraph = graph.mapEdges(
    new MapFunction<Edge<Long, Long>, Double>() {
        @Override
        public Edge<Long, Double> map(Edge<Long, Long> value) {
            return new Edge<>(value.getSource(), value.getTarget(), value.getValue() * 1.5);
        }
    }
);

// 过滤顶点
Graph<Long, Long, Long> filteredGraph = graph.filterVertices(
    new FilterFunction<Vertex<Long, Long>>() {
        @Override
        public boolean filter(Vertex<Long, Long> vertex) {
            return vertex.getValue() > 0;
        }
    }
);

// 过滤边
Graph<Long, Long, Long> filteredGraph = graph.filterEdges(
    new FilterFunction<Edge<Long, Long>>() {
        @Override
        public boolean filter(Edge<Long, Long> edge) {
            return edge.getValue() > 0;
        }
    }
);

4.13.4 图算法

import org.apache.flink.graph.library.LabelPropagation;
import org.apache.flink.graph.library.PageRank;
import org.apache.flink.graph.library.ConnectedComponents;

// 标签传播(社区发现)
DataSet<Vertex<Long, Long>> labels = graph.run(
    new LabelPropagation<Long, Long, Long>(maxIterations)
);

// PageRank
DataSet<Vertex<Long, Double>> ranks = graph.run(
    new PageRank<>(dampingFactor, maxIterations)
);

// 连通分量
DataSet<Vertex<Long, Long>> components = graph.run(
    new ConnectedComponents<>(maxIterations)
);

// 三角计数
DataSet<Triplet<Long, Long, Long>> triangles = graph.run(
    new Triangles<>()
);

4.13.5 图生成器

import org.apache.flink.graph.generator.RMatGraph;
import org.apache.flink.graph.generator.GraphGenerators;

// RMat 图(用于测试)
Graph<Long, Long, Long> rmatGraph = GraphGenerators
    .rmatGraph(env, numVertices, numEdges)
    .setConf(new RMatGraphConfiguration(20, 64))
    .setSeed(42);

// 网格图
Graph<Long, Long, Long> gridGraph = GraphGenerators
    .gridGraph(env, rows, cols);

// 循环图
Graph<Long, Long, Long> cycleGraph = GraphGenerators
    .cycleGraph(env, numVertices);

4.14 大状态调优

4.14.1 Checkpoint 调优

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

// 设置 Checkpoint 间隔
env.enableCheckpointing(60000);  // 60 秒

CheckpointConfig config = env.getCheckpointConfig();

// 设置 Checkpoint 之间的最小间隔
config.setMinPauseBetweenCheckpoints(30);  // 30 秒

// 设置 Checkpoint 超时时间
config.setCheckpointTimeout(120000);  // 2 分钟

// 允许并发 Checkpoint 数量
config.setMaxConcurrentCheckpoints(1);

// 启用非对齐 Checkpoint(加速)
config.enableUnalignedCheckpoints();

4.14.2 RocksDB 调优

⚠️ 未实现:当前代码库中未包含 RocksDB 调优的示例代码

// 待实现
// 使用 RocksDB 状态后端
EmbeddedRocksDBStateBackend rocksDB = new EmbeddedRocksDBStateBackend(true);
env.setStateBackend(rocksDB);

// 启用增量 Checkpoint
env.getCheckpointConfig().enableIncrementalCheckpointing();

💡 Tips: RocksDB 调优与状态 TTL 配置

状态 TTL 配置(防止状态无限增长):

StateTtlConfig ttlConfig = StateTtlConfig
    .newBuilder(Time.hours(1))  // 1 小时后过期
    .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
    .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
    .build();

ValueStateDescriptor<String> descriptor = new ValueStateDescriptor<>("state", String.class);
descriptor.enableTimeToLive(ttlConfig);

4.14.3 内存管理

⚠️ 未实现:当前代码库中未包含内存管理示例

// 待实现 - 配置通过 flink-conf.yaml 完成
// state.backend.rocksdb.memory.managed: true
// state.backend.rocksdb.memory.fixed-per-slot: 256mb

4.14.4 压缩配置

⚠️ 未实现:当前代码库中未包含压缩配置示例

// 待实现
ExecutionConfig executionConfig = new ExecutionConfig();
executionConfig.setUseSnapshotCompression(true);  // 启用 Checkpoints 压缩

4.14.5 任务本地恢复

⚠️ 未实现:当前代码库中未包含本地恢复示例

// 待实现 - 配置通过 flink-conf.yaml 完成
// state.backend.local-recovery: true

4.15 Savepoint 操作

# 触发 Savepoint
bin/flink savepoint <jobId> [targetDirectory]

# 使用 YARN 触发 Savepoint
bin/flink savepoint <jobId> [targetDirectory] -yid <yarnAppId>

# 带 Savepoint 取消作业
bin/flink cancel -s [targetDirectory] <jobId>

# 从 Savepoint 恢复
bin/flink run -s <savepointPath> [runArgs]

# 跳过无法映射的状态恢复
bin/flink run -s <savepointPath> -n [runArgs]

# 删除 Savepoint
bin/flink savepoint -d <savepointPath>
// 为算子分配唯一 ID(必须)
DataStream<String> stream = env
    .addSource(new StatefulSource()).uid("source-id")
    .shuffle()
    .map(new StatefulMapper()).uid("mapper-id")
    .print();  // 无状态算子可不分配 ID

参见: 01-应用场景 | 02-架构与原理 | 03-部署与运维 | 附录


标签: #flink #java #开发 #实践 #编程


RAG 智能问答

针对本文继续提问:《Flink 教程(四):Java 开发实践》