📌 本文是「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 #开发 #实践 #编程