概述
Flink DataStream API 不仅提供了基础的流处理能力,还包含了许多高级特性,使其能够应对复杂的生产环境挑战。本文将深入探讨这些高级用法,包括容错机制、性能优化、高级连接器等方面。
💡 提示: 本文假设读者已经掌握了 Flink DataStream API 的基础知识。如果您是初学者,建议先阅读 Flink DataStream API 中的算子。
本文重点
- 容错机制: Checkpoint、Savepoint 和状态一致性保证
- 性能优化: 并行度、状态后端、网络缓冲和算子链优化
- 高级集成: Table & SQL API、CDC、数据湖连接器
- 复杂处理: CEP 模式匹配、侧输出流、广播状态
- 生产实践: 监控、调试和最佳实践
容错机制
Checkpoint 机制
Checkpoint 是 Flink 容错的核心机制,通过定期保存应用状态来实现故障恢复。
⚠️ 注意: Checkpoint 的配置直接影响作业的性能和容错能力,需要根据实际场景权衡。
基础配置
import org.apache.flink.streaming.api.CheckpointingMode;
import org.apache.flink.streaming.api.environment.CheckpointConfig;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.contrib.streaming.state.RocksDBStateBackend;
public class CheckpointExample {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 启用 Checkpoint,间隔 10 秒
env.enableCheckpointing(10000);
// 高级配置
CheckpointConfig checkpointConfig = env.getCheckpointConfig();
// 设置 exactly-once 语义(默认)
checkpointConfig.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
// 设置最小间隔,避免 checkpoint 太频繁
checkpointConfig.setMinPauseBetweenCheckpoints(5000);
// 设置超时时间(1分钟)
checkpointConfig.setCheckpointTimeout(60000);
// 设置并发 checkpoint 数量(通常设为1)
checkpointConfig.setMaxConcurrentCheckpoints(1);
// 启用外部化 checkpoint,作业取消后保留
checkpointConfig.enableExternalizedCheckpoints(
CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION
);
// 设置容忍 checkpoint 失败次数
checkpointConfig.setTolerableCheckpointFailureNumber(3);
// 设置状态后端(使用 RocksDB 处理大状态)
env.setStateBackend(new RocksDBStateBackend("hdfs://namenode:9000/flink/checkpoints"));
}
}
对齐与非对齐 Checkpoint
💡 提示: 非对齐 Checkpoint 可以显著减少背压情况下的 Checkpoint 时间,但会增加状态大小。
import java.time.Duration;
// 启用非对齐 checkpoint(Flink 1.11+)
checkpointConfig.enableUnalignedCheckpoints();
// 设置对齐超时,超时后自动切换到非对齐
checkpointConfig.setAlignmentTimeout(Duration.ofSeconds(30));
// 非对齐 checkpoint 的缓冲区配置
checkpointConfig.setUnalignedCheckpointRetention(Duration.ofMinutes(5));
Savepoint 管理
Savepoint 是手动触发的检查点,主要用于升级、迁移和 A/B 测试。
🔧 配置: Savepoint 与 Checkpoint 的主要区别在于 Savepoint 是手动触发的,且会保存完整的状态信息。
触发 Savepoint
# 触发 savepoint
flink savepoint <jobId> hdfs://namenode:9000/savepoints
# 带状态停止作业
flink stop --savepointPath hdfs://namenode:9000/savepoints <jobId>
# 取消作业并创建 savepoint
flink cancel -s hdfs://namenode:9000/savepoints <jobId>
从 Savepoint 恢复
# 从 savepoint 启动
flink run -s hdfs://namenode:9000/savepoints/savepoint-xxx -c com.example.MyJob job.jar
# 允许跳过无法恢复的算子
flink run -s hdfs://namenode:9000/savepoints/savepoint-xxx \
--allowNonRestoredState \
-c com.example.MyJob job.jar
编程式 Savepoint 管理
import org.apache.flink.runtime.state.FunctionInitializationContext;
import org.apache.flink.runtime.state.FunctionSnapshotContext;
import org.apache.flink.streaming.api.checkpoint.CheckpointedFunction;
public class SavepointCompatibleFunction implements MapFunction<String, String>, CheckpointedFunction {
private transient ListState<String> checkpointedState;
private List<String> bufferedElements;
@Override
public void initializeState(FunctionInitializationContext context) throws Exception {
ListStateDescriptor<String> descriptor = new ListStateDescriptor<>(
"buffered-elements",
TypeInformation.of(String.class)
);
checkpointedState = context.getOperatorStateStore().getListState(descriptor);
// 从 checkpoint/savepoint 恢复
if (context.isRestored()) {
for (String element : checkpointedState.get()) {
bufferedElements.add(element);
}
}
}
@Override
public void snapshotState(FunctionSnapshotContext context) throws Exception {
checkpointedState.clear();
checkpointedState.addAll(bufferedElements);
}
@Override
public String map(String value) throws Exception {
bufferedElements.add(value);
return value.toUpperCase();
}
}
状态一致性保证
端到端 Exactly-Once
⚠️ 注意: 端到端 Exactly-Once 需要 Source 和 Sink 都支持事务语义。
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.connector.kafka.source.KafkaSource;
import org.apache.flink.connector.kafka.source.enumerator.initializer.OffsetsInitializer;
import org.apache.flink.connector.kafka.sink.KafkaSink;
import org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.CheckpointingMode;
import org.apache.flink.connector.base.DeliveryGuarantee;
import java.util.Properties;
// Kafka exactly-once 示例
public class ExactlyOnceKafkaExample {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(5000);
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
// Kafka Source 配置
Properties sourceProps = new Properties();
sourceProps.setProperty("bootstrap.servers", "localhost:9092");
sourceProps.setProperty("group.id", "flink-consumer");
sourceProps.setProperty("isolation.level", "read_committed");
KafkaSource<String> source = KafkaSource.<String>builder()
.setProperties(sourceProps)
.setTopics("input-topic")
.setStartingOffsets(OffsetsInitializer.latest())
.setValueOnlyDeserializer(new SimpleStringSchema())
.build();
// Kafka Sink 配置(exactly-once)
KafkaSink<String> sink = KafkaSink.<String>builder()
.setBootstrapServers("localhost:9092")
.setKafkaProducerConfig(getExactlyOnceProducerConfig())
.setRecordSerializer(
KafkaRecordSerializationSchema.builder()
.setTopic("output-topic")
.setValueSerializationSchema(new SimpleStringSchema())
.build()
)
.setDeliverGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
.setTransactionalIdPrefix("flink-sink")
.build();
DataStream<String> stream = env.fromSource(source,
WatermarkStrategy.noWatermarks(), "Kafka Source");
stream.map(value -> processValue(value))
.sinkTo(sink);
env.execute("Exactly-Once Kafka Example");
}
private static Properties getExactlyOnceProducerConfig() {
Properties props = new Properties();
props.setProperty("transaction.timeout.ms", "900000"); // 15 分钟
props.setProperty("enable.idempotence", "true");
return props;
}
}
Table & SQL API 集成
Flink 的 Table & SQL API 提供了声明式的流处理接口,可以与 DataStream API 无缝集成。
DataStream 与 Table 转换
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
import org.apache.flink.table.api.Table;
import org.apache.flink.types.Row;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.$;
import lombok.Data;
import lombok.AllArgsConstructor;
import java.time.Instant;
public class StreamTableConversionExample {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
// DataStream 转 Table
DataStream<Order> orderStream = env.fromElements(
new Order("order1", "user1", 100.0, Instant.now()),
new Order("order2", "user2", 200.0, Instant.now())
);
// 方式1:自动推断 schema
Table orderTable = tableEnv.fromDataStream(orderStream);
// 方式2:指定字段和时间属性
Table orderTableWithTime = tableEnv.fromDataStream(
orderStream,
$("orderId"),
$("userId"),
$("amount"),
$("orderTime").rowtime()
);
// 注册为临时视图
tableEnv.createTemporaryView("orders", orderTableWithTime);
// SQL 查询
Table result = tableEnv.sqlQuery(
"SELECT " +
" userId, " +
" TUMBLE_START(orderTime, INTERVAL '1' HOUR) as window_start, " +
" SUM(amount) as total_amount " +
"FROM orders " +
"GROUP BY userId, TUMBLE(orderTime, INTERVAL '1' HOUR)"
);
// Table 转回 DataStream
DataStream<Row> resultStream = tableEnv.toDataStream(result);
// 处理 changelog 流
Table updatingTable = tableEnv.sqlQuery(
"SELECT userId, COUNT(*) as order_count FROM orders GROUP BY userId"
);
DataStream<Row> changelogStream = tableEnv.toChangelogStream(updatingTable);
changelogStream.print();
env.execute();
}
}
@Data
@AllArgsConstructor
public class Order {
private String orderId;
private String userId;
private Double amount;
private Instant orderTime;
}
SQL 查询优化
public class SqlOptimizationExample {
public static void main(String[] args) {
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
// 配置优化参数
Configuration configuration = tableEnv.getConfig().getConfiguration();
// 启用对象重用
configuration.setBoolean("pipeline.object-reuse", true);
// 设置并行度
configuration.setInteger("parallelism.default", 4);
// 启用 minibatch 优化
configuration.setBoolean("table.exec.mini-batch.enabled", true);
configuration.setString("table.exec.mini-batch.allow-latency", "5 s");
configuration.setLong("table.exec.mini-batch.size", 5000L);
// 设置状态 TTL
configuration.setString("table.exec.state.ttl", "1 h");
// 使用 Blink planner 特性
tableEnv.executeSql(
"CREATE TABLE orders (" +
" order_id STRING, " +
" user_id STRING, " +
" amount DECIMAL(10, 2), " +
" order_time TIMESTAMP(3), " +
" WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND " +
") WITH (" +
" 'connector' = 'kafka', " +
" 'topic' = 'orders', " +
" 'properties.bootstrap.servers' = 'localhost:9092', " +
" 'format' = 'json' " +
")"
);
// 创建维表(用于 JOIN)
tableEnv.executeSql(
"CREATE TABLE users (" +
" user_id STRING PRIMARY KEY NOT ENFORCED, " +
" user_name STRING, " +
" user_level STRING " +
") WITH (" +
" 'connector' = 'jdbc', " +
" 'url' = 'jdbc:mysql://localhost:3306/db', " +
" 'table-name' = 'users', " +
" 'lookup.cache.max-rows' = '10000', " +
" 'lookup.cache.ttl' = '10 min' " +
")"
);
// 维表 JOIN 查询
Table result = tableEnv.sqlQuery(
"SELECT " +
" o.order_id, " +
" o.user_id, " +
" u.user_name, " +
" u.user_level, " +
" o.amount " +
"FROM orders AS o " +
"JOIN users FOR SYSTEM_TIME AS OF o.order_time AS u " +
"ON o.user_id = u.user_id"
);
}
}
动态表与流表对偶性
public class DynamicTableExample {
public static void main(String[] args) {
// 创建 CDC 源表
tableEnv.executeSql(
"CREATE TABLE products_cdc (" +
" id INT PRIMARY KEY NOT ENFORCED, " +
" name STRING, " +
" price DECIMAL(10, 2), " +
" update_time TIMESTAMP(3) METADATA FROM 'value.source.timestamp' VIRTUAL " +
") WITH (" +
" 'connector' = 'mysql-cdc', " +
" 'hostname' = 'localhost', " +
" 'port' = '3306', " +
" 'username' = 'root', " +
" 'password' = 'password', " +
" 'database-name' = 'inventory', " +
" 'table-name' = 'products' " +
")"
);
// 创建去重查询
Table deduplicatedTable = tableEnv.sqlQuery(
"SELECT * FROM (" +
" SELECT *, " +
" ROW_NUMBER() OVER (PARTITION BY id ORDER BY update_time DESC) as rn " +
" FROM products_cdc" +
") WHERE rn = 1"
);
// 创建物化视图
tableEnv.executeSql(
"CREATE TABLE product_summary (" +
" category STRING PRIMARY KEY NOT ENFORCED, " +
" total_products BIGINT, " +
" avg_price DECIMAL(10, 2) " +
") WITH (" +
" 'connector' = 'upsert-kafka', " +
" 'topic' = 'product-summary', " +
" 'properties.bootstrap.servers' = 'localhost:9092', " +
" 'key.format' = 'json', " +
" 'value.format' = 'json' " +
")"
);
// 持续查询并写入
tableEnv.executeSql(
"INSERT INTO product_summary " +
"SELECT " +
" SUBSTRING(name, 1, 3) as category, " +
" COUNT(*) as total_products, " +
" AVG(price) as avg_price " +
"FROM products_cdc " +
"GROUP BY SUBSTRING(name, 1, 3)"
);
}
}
性能优化
并行度优化
💡 提示: 并行度的设置需要考虑数据量、计算复杂度和可用资源。过高或过低的并行度都会影响性能。
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.api.common.functions.FlatMapFunction;
import org.apache.flink.api.common.functions.ReduceFunction;
import org.apache.flink.runtime.jobgraph.JobGraph;
public class ParallelismOptimization {
public static void main(String[] args) {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 全局并行度
env.setParallelism(8);
// 算子级并行度
DataStream<String> source = env
.readTextFile("hdfs://path/to/file")
.setParallelism(4); // Source 并行度
DataStream<String> processed = source
.flatMap(new Tokenizer())
.setParallelism(16) // FlatMap 高并行度
.keyBy(value -> value)
.reduce(new Counter())
.setParallelism(8); // Reduce 中等并行度
processed
.addSink(new CustomSink())
.setParallelism(2); // Sink 低并行度
// 动态并行度调整
int parallelism = calculateOptimalParallelism(
env.getStreamGraph().getJobGraph()
);
env.setParallelism(parallelism);
}
private static int calculateOptimalParallelism(JobGraph jobGraph) {
// 基于资源和数据量计算最优并行度
int availableCores = Runtime.getRuntime().availableProcessors();
long estimatedDataSize = getEstimatedDataSize();
int taskManagers = getClusterTaskManagers();
return Math.min(
availableCores * taskManagers,
(int) (estimatedDataSize / (1024 * 1024 * 100)) // 每 100MB 一个并行度
);
}
}
状态后端优化
🔧 配置: RocksDB 状态后端适合大状态场景,但需要合理配置以获得最佳性能。
import org.apache.flink.contrib.streaming.state.RocksDBStateBackend;
import org.apache.flink.contrib.streaming.state.RocksDBOptionsFactory;
import org.apache.flink.contrib.streaming.state.RocksDBMemoryConfiguration;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.api.common.state.MapState;
import org.apache.flink.api.common.state.MapStateDescriptor;
import org.apache.flink.api.common.typeinfo.Types;
import org.rocksdb.*;
import java.util.Collection;
public class StateBackendOptimization {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// RocksDB 状态后端配置
RocksDBStateBackend rocksDBStateBackend = new RocksDBStateBackend(
"hdfs://namenode:9000/flink/checkpoints"
);
// 增量检查点
rocksDBStateBackend.setEnableIncrementalCheckpointing(true);
// 配置 RocksDB 选项
rocksDBStateBackend.setOptions(new OptionsFactory() {
@Override
public DBOptions createDBOptions(DBOptions currentOptions, Collection<AutoCloseable> handlesToClose) {
return currentOptions
.setMaxBackgroundJobs(4)
.setMaxOpenFiles(-1);
}
@Override
public ColumnFamilyOptions createColumnOptions(
ColumnFamilyOptions currentOptions, Collection<AutoCloseable> handlesToClose) {
return currentOptions
.setTableFormatConfig(
new BlockBasedTableConfig()
.setBlockCacheSize(256 * 1024 * 1024) // 256 MB
.setBlockSize(128 * 1024) // 128 KB
.setCacheIndexAndFilterBlocks(true)
.setPinL0FilterAndIndexBlocksInCache(true)
)
.setMemTableConfig(
new HashLinkedListMemTableConfig()
.setBucketCount(131072)
)
.setCompactionStyle(CompactionStyle.LEVEL)
.setLevel0FileNumCompactionTrigger(10)
.setMaxBytesForLevelBase(256 * 1024 * 1024);
}
});
env.setStateBackend(rocksDBStateBackend);
// 使用定时器优化
DataStream<Event> events = env.addSource(new EventSource());
events.keyBy(Event::getKey)
.process(new OptimizedProcessFunction())
.print();
env.execute();
}
}
class OptimizedProcessFunction extends KeyedProcessFunction<String, Event, Result> {
private transient MapState<Long, List<Event>> timerState;
private final long TIMER_INTERVAL = 60000; // 1 分钟
@Override
public void open(Configuration parameters) throws Exception {
MapStateDescriptor<Long, List<Event>> descriptor = new MapStateDescriptor<>(
"timer-state",
Types.LONG,
Types.LIST(Types.POJO(Event.class))
);
timerState = getRuntimeContext().getMapState(descriptor);
}
@Override
public void processElement(Event event, Context ctx, Collector<Result> out) throws Exception {
long timerKey = ctx.timestamp() / TIMER_INTERVAL * TIMER_INTERVAL;
List<Event> events = timerState.get(timerKey);
if (events == null) {
events = new ArrayList<>();
ctx.timerService().registerEventTimeTimer(timerKey + TIMER_INTERVAL);
}
events.add(event);
timerState.put(timerKey, events);
}
@Override
public void onTimer(long timestamp, OnTimerContext ctx, Collector<Result> out) throws Exception {
List<Event> events = timerState.get(timestamp - TIMER_INTERVAL);
if (events != null) {
out.collect(processEvents(events));
timerState.remove(timestamp - TIMER_INTERVAL);
}
}
}
网络缓冲优化
public class NetworkBufferOptimization {
public static void main(String[] args) {
Configuration config = new Configuration();
// 网络缓冲配置
config.setString("taskmanager.network.memory.fraction", "0.15");
config.setString("taskmanager.network.memory.min", "128mb");
config.setString("taskmanager.network.memory.max", "1gb");
// 网络缓冲池大小
config.setInteger("taskmanager.network.numberOfBuffers", 4096);
// 配置网络超时
config.setString("akka.ask.timeout", "60 s");
config.setString("web.timeout", "60000");
// 启用网络缓冲反压
config.setBoolean("taskmanager.network.credit-model", true);
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(config);
// 使用批处理模式提高吞吐量
env.setBufferTimeout(100); // 100ms 的缓冲超时
DataStream<String> stream = env.socketTextStream("localhost", 9999);
// 使用 rebalance 避免数据倾斜
stream.rebalance()
.map(new HeavyComputation())
.keyBy(value -> value.hashCode() % 100)
.reduce((a, b) -> a + b)
.print();
}
}
算子链优化
public class OperatorChainOptimization {
public static void main(String[] args) {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 全局禁用算子链
// env.disableOperatorChaining();
DataStream<String> source = env.readTextFile("input.txt");
DataStream<String> result = source
.filter(line -> line.length() > 0)
.startNewChain() // 开始新的算子链
.map(String::toLowerCase)
.filter(line -> line.contains("error"))
.disableChaining() // 禁用链接
.map(line -> "ERROR: " + line)
.slotSharingGroup("group1") // 设置 slot 共享组
.print();
// 手动设置算子名称便于监控
source.filter(line -> line.length() > 0)
.name("FilterEmpty")
.uid("filter-empty-lines")
.map(String::toUpperCase)
.name("ToUpperCase")
.uid("to-upper-case");
}
}
高级连接器
Flink CDC
Flink CDC 提供了实时捕获数据库变更的能力,支持 MySQL、PostgreSQL、Oracle 等多种数据库。
💡 提示: CDC 连接器可以直接读取数据库的 binlog/redo log,实现低延迟的数据同步。
import com.ververica.cdc.connectors.mysql.source.MySqlSource;
import com.ververica.cdc.connectors.mysql.table.StartupOptions;
import com.ververica.cdc.debezium.JsonDebeziumDeserializationSchema;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.datastream.DataStream;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.databind.JsonNode;
import lombok.Data;
import lombok.AllArgsConstructor;
public class FlinkCDCExample {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(3000);
// MySQL CDC Source
MySqlSource<String> mySqlSource = MySqlSource.<String>builder()
.hostname("localhost")
.port(3306)
.databaseList("inventory")
.tableList("inventory.products", "inventory.orders")
.username("flink")
.password("flink")
.serverId("5400-5404")
.deserializer(new JsonDebeziumDeserializationSchema())
.includeSchemaChanges(true) // 包含 DDL 变更
.startupOptions(StartupOptions.initial())
.build();
DataStream<String> cdcStream = env.fromSource(
mySqlSource,
WatermarkStrategy.noWatermarks(),
"MySQL CDC Source"
);
// 解析 CDC 数据
DataStream<CDCEvent> events = cdcStream.map(json -> {
ObjectMapper mapper = new ObjectMapper();
JsonNode node = mapper.readTree(json);
String op = node.get("op").asText();
String table = node.get("source").get("table").asText();
JsonNode after = node.get("after");
JsonNode before = node.get("before");
return new CDCEvent(op, table, before, after);
});
// 处理不同的操作类型
events.process(new ProcessFunction<CDCEvent, String>() {
@Override
public void processElement(CDCEvent event, Context ctx, Collector<String> out) {
switch (event.getOp()) {
case "c": // CREATE
handleInsert(event, out);
break;
case "u": // UPDATE
handleUpdate(event, out);
break;
case "d": // DELETE
handleDelete(event, out);
break;
case "r": // READ (snapshot)
handleSnapshot(event, out);
break;
}
}
});
// 多表关联示例
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
tableEnv.executeSql(
"CREATE TABLE products (" +
" id INT PRIMARY KEY NOT ENFORCED, " +
" name STRING, " +
" description STRING, " +
" weight DECIMAL(10,2) " +
") WITH (" +
" 'connector' = 'mysql-cdc', " +
" 'hostname' = 'localhost', " +
" 'port' = '3306', " +
" 'username' = 'root', " +
" 'password' = 'password', " +
" 'database-name' = 'inventory', " +
" 'table-name' = 'products' " +
")"
);
tableEnv.executeSql(
"CREATE TABLE orders (" +
" order_id INT PRIMARY KEY NOT ENFORCED, " +
" product_id INT, " +
" quantity INT, " +
" order_date TIMESTAMP(3) " +
") WITH (" +
" 'connector' = 'mysql-cdc', " +
" 'hostname' = 'localhost', " +
" 'port' = '3306', " +
" 'username' = 'root', " +
" 'password' = 'password', " +
" 'database-name' = 'inventory', " +
" 'table-name' = 'orders' " +
")"
);
// 实时物化视图
Table result = tableEnv.sqlQuery(
"SELECT " +
" p.name as product_name, " +
" SUM(o.quantity) as total_quantity, " +
" COUNT(DISTINCT o.order_id) as order_count " +
"FROM products p " +
"JOIN orders o ON p.id = o.product_id " +
"GROUP BY p.name"
);
tableEnv.toChangelogStream(result).print();
env.execute("Flink CDC Example");
}
}
@Data
@AllArgsConstructor
class CDCEvent {
private String op;
private String table;
private JsonNode before;
private JsonNode after;
}
数据湖集成
public class DataLakeIntegration {
public static void main(String[] args) {
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
// Iceberg 集成
tableEnv.executeSql(
"CREATE CATALOG iceberg_catalog WITH (" +
" 'type' = 'iceberg', " +
" 'catalog-type' = 'hive', " +
" 'uri' = 'thrift://localhost:9083', " +
" 'warehouse' = 'hdfs://namenode:9000/user/hive/warehouse' " +
")"
);
tableEnv.useCatalog("iceberg_catalog");
// 创建 Iceberg 表
tableEnv.executeSql(
"CREATE TABLE IF NOT EXISTS events (" +
" event_id BIGINT, " +
" event_type STRING, " +
" event_time TIMESTAMP(3), " +
" properties MAP<STRING, STRING>, " +
" PRIMARY KEY (event_id) NOT ENFORCED " +
") PARTITIONED BY (event_type, DATE(event_time)) " +
"WITH (" +
" 'format-version' = '2', " +
" 'write.upsert.enabled' = 'true' " +
")"
);
// Hudi 集成
tableEnv.executeSql(
"CREATE TABLE hudi_events (" +
" uuid STRING PRIMARY KEY NOT ENFORCED, " +
" ts TIMESTAMP(3), " +
" event_type STRING, " +
" payload STRING " +
") WITH (" +
" 'connector' = 'hudi', " +
" 'path' = 'hdfs://namenode:9000/hudi/events', " +
" 'table.type' = 'MERGE_ON_READ', " +
" 'write.operation' = 'upsert', " +
" 'hoodie.datasource.write.recordkey.field' = 'uuid', " +
" 'hoodie.datasource.write.precombine.field' = 'ts' " +
")"
);
// Delta Lake 集成
DataStream<Event> eventStream = env.addSource(new EventSource());
// 写入 Delta Lake
eventStream.map(event -> Row.of(
event.getId(),
event.getType(),
event.getTimestamp(),
event.getPayload()
))
.addSink(
DeltaSink.forRowData(
new Path("hdfs://namenode:9000/delta/events"),
new Configuration(),
Types.ROW(
Types.LONG,
Types.STRING,
Types.SQL_TIMESTAMP,
Types.STRING
))
.build()
);
// 实时 ETL 到数据湖
tableEnv.executeSql(
"INSERT INTO iceberg_catalog.default.events " +
"SELECT " +
" CAST(uuid AS BIGINT) as event_id, " +
" event_type, " +
" ts as event_time, " +
" STR_TO_MAP(payload) as properties " +
"FROM hudi_events " +
"WHERE ts > CURRENT_TIMESTAMP - INTERVAL '1' DAY"
);
}
}
自定义 Source 和 Sink
// 自定义 Source
public class CustomKafkaSource implements SourceFunction<CustomEvent>, CheckpointedFunction {
private volatile boolean isRunning = true;
private transient ListState<Long> offsetState;
private long currentOffset = 0;
private final String topic;
private transient KafkaConsumer<String, String> consumer;
@Override
public void run(SourceContext<CustomEvent> ctx) throws Exception {
consumer = createConsumer();
// 从检查点恢复偏移量
if (currentOffset > 0) {
consumer.seek(new TopicPartition(topic, 0), currentOffset);
}
while (isRunning) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
synchronized (ctx.getCheckpointLock()) {
CustomEvent event = parseEvent(record.value());
ctx.collectWithTimestamp(event, event.getTimestamp());
currentOffset = record.offset();
}
}
// 发送水位线
ctx.emitWatermark(new Watermark(System.currentTimeMillis() - 5000));
}
}
@Override
public void cancel() {
isRunning = false;
if (consumer != null) {
consumer.close();
}
}
@Override
public void snapshotState(FunctionSnapshotContext context) throws Exception {
offsetState.clear();
offsetState.add(currentOffset);
}
@Override
public void initializeState(FunctionInitializationContext context) throws Exception {
ListStateDescriptor<Long> descriptor = new ListStateDescriptor<>(
"offset-state",
TypeInformation.of(Long.class)
);
offsetState = context.getOperatorStateStore().getListState(descriptor);
if (context.isRestored()) {
for (Long offset : offsetState.get()) {
currentOffset = offset;
}
}
}
}
// 自定义 Sink
public class CustomJDBCSink extends RichSinkFunction<Event> implements CheckpointedFunction {
private transient Connection connection;
private transient PreparedStatement preparedStatement;
private final List<Event> bufferedEvents = new ArrayList<>();
private final int batchSize = 1000;
@Override
public void open(Configuration parameters) throws Exception {
connection = DriverManager.getConnection(
"jdbc:mysql://localhost:3306/db",
"user",
"password"
);
connection.setAutoCommit(false);
preparedStatement = connection.prepareStatement(
"INSERT INTO events (id, type, timestamp, data) VALUES (?, ?, ?, ?) " +
"ON DUPLICATE KEY UPDATE type = VALUES(type), timestamp = VALUES(timestamp), data = VALUES(data)"
);
}
@Override
public void invoke(Event event, Context context) throws Exception {
bufferedEvents.add(event);
if (bufferedEvents.size() >= batchSize) {
flush();
}
}
private void flush() throws Exception {
for (Event event : bufferedEvents) {
preparedStatement.setLong(1, event.getId());
preparedStatement.setString(2, event.getType());
preparedStatement.setTimestamp(3, Timestamp.from(event.getTimestamp()));
preparedStatement.setString(4, event.getData());
preparedStatement.addBatch();
}
preparedStatement.executeBatch();
connection.commit();
bufferedEvents.clear();
}
@Override
public void close() throws Exception {
if (!bufferedEvents.isEmpty()) {
flush();
}
if (preparedStatement != null) {
preparedStatement.close();
}
if (connection != null) {
connection.close();
}
}
@Override
public void snapshotState(FunctionSnapshotContext context) throws Exception {
flush(); // 确保所有数据都被写入
}
@Override
public void initializeState(FunctionInitializationContext context) throws Exception {
// 恢复时不需要特殊处理,因为我们在快照时已经 flush 了
}
}
复杂事件处理
CEP 模式匹配
Complex Event Processing (CEP) 库允许在流中检测复杂的事件模式。
💡 提示: CEP 适用于实时监控、欺诈检测、异常检测等场景。
import org.apache.flink.cep.CEP;
import org.apache.flink.cep.PatternStream;
import org.apache.flink.cep.pattern.Pattern;
import org.apache.flink.cep.pattern.conditions.SimpleCondition;
import org.apache.flink.cep.PatternSelectFunction;
import org.apache.flink.cep.PatternTimeoutFunction;
import org.apache.flink.streaming.api.windowing.time.Time;
import org.apache.flink.util.OutputTag;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import java.util.List;
import java.util.Map;
public class CEPExample {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<Event> eventStream = env.addSource(new EventSource());
// 定义复杂事件模式
Pattern<Event, ?> pattern = Pattern.<Event>begin("start")
.where(new SimpleCondition<Event>() {
@Override
public boolean filter(Event event) {
return event.getType().equals("login");
}
})
.followedBy("middle")
.where(new SimpleCondition<Event>() {
@Override
public boolean filter(Event event) {
return event.getType().equals("add_to_cart");
}
})
.times(3, 5) // 3到5次加购
.within(Time.minutes(10)) // 10分钟内
.followedByAny("end")
.where(new SimpleCondition<Event>() {
@Override
public boolean filter(Event event) {
return event.getType().equals("purchase");
}
});
PatternStream<Event> patternStream = CEP.pattern(
eventStream.keyBy(Event::getUserId),
pattern
);
// 提取匹配的事件序列
DataStream<Alert> alerts = patternStream.select(
(Map<String, List<Event>> pattern) -> {
List<Event> start = pattern.get("start");
List<Event> middle = pattern.get("middle");
List<Event> end = pattern.get("end");
return new Alert(
start.get(0).getUserId(),
"购物行为",
middle.size(),
end.get(0).getTimestamp()
);
}
);
// 复杂的模式:检测异常登录
Pattern<LoginEvent, ?> loginPattern = Pattern.<LoginEvent>begin("first")
.where(new SimpleCondition<LoginEvent>() {
@Override
public boolean filter(LoginEvent event) {
return event.isSuccess();
}
})
.next("second")
.where(new SimpleCondition<LoginEvent>() {
@Override
public boolean filter(LoginEvent event) {
return !event.isSuccess();
}
})
.times(3).consecutive() // 连续3次失败
.within(Time.minutes(5));
PatternStream<LoginEvent> loginPatternStream = CEP.pattern(
loginStream.keyBy(LoginEvent::getUserId),
loginPattern
);
// 侧输出流处理超时事件
OutputTag<LoginEvent> timedOutTag = new OutputTag<LoginEvent>("timed-out"){};
SingleOutputStreamOperator<Alert> loginAlerts = loginPatternStream
.select(
timedOutTag,
// 超时事件处理
(Map<String, List<LoginEvent>> pattern, long timeoutTimestamp) -> {
LoginEvent first = pattern.get("first").get(0);
return new Alert(
first.getUserId(),
"登录超时",
pattern.get("second").size(),
timeoutTimestamp
);
},
// 正常匹配处理
(Map<String, List<LoginEvent>> pattern) -> {
LoginEvent first = pattern.get("first").get(0);
return new Alert(
first.getUserId(),
"异常登录尝试",
3,
System.currentTimeMillis()
);
}
);
// 获取超时的侧输出流
DataStream<Alert> timedOutAlerts = loginAlerts.getSideOutput(timedOutTag);
alerts.print("正常告警");
timedOutAlerts.print("超时告警");
env.execute("CEP Example");
}
}
侧输出流
public class SideOutputExample {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 定义侧输出标签
final OutputTag<String> errorTag = new OutputTag<String>("errors"){};
final OutputTag<Event> lateTag = new OutputTag<Event>("late-events"){};
final OutputTag<Event> anomalyTag = new OutputTag<Event>("anomalies"){};
SingleOutputStreamOperator<ProcessedEvent> mainDataStream = env
.addSource(new EventSource())
.assignTimestampsAndWatermarks(
WatermarkStrategy.<Event>forBoundedOutOfOrderness(Duration.ofSeconds(10))
.withTimestampAssigner((event, timestamp) -> event.getTimestamp())
)
.keyBy(Event::getKey)
.window(TumblingEventTimeWindows.of(Time.minutes(5)))
.allowedLateness(Time.minutes(1))
.sideOutputLateData(lateTag)
.process(new ProcessWindowFunction<Event, ProcessedEvent, String, TimeWindow>() {
@Override
public void process(String key, Context context,
Iterable<Event> elements,
Collector<ProcessedEvent> out) throws Exception {
List<Event> events = new ArrayList<>();
elements.forEach(events::add);
// 异常检测
for (Event event : events) {
if (isAnomaly(event)) {
context.output(anomalyTag, event);
}
}
try {
ProcessedEvent result = processEvents(events);
out.collect(result);
} catch (Exception e) {
// 错误数据输出到侧输出流
context.output(errorTag,
String.format("Error processing key %s: %s", key, e.getMessage()));
}
}
private boolean isAnomaly(Event event) {
// 异常检测逻辑
return event.getValue() > 1000 || event.getValue() < 0;
}
});
// 获取侧输出流
DataStream<Event> lateEvents = mainDataStream.getSideOutput(lateTag);
DataStream<String> errors = mainDataStream.getSideOutput(errorTag);
DataStream<Event> anomalies = mainDataStream.getSideOutput(anomalyTag);
// 处理迟到数据
lateEvents
.map(event -> "Late: " + event)
.addSink(new FileSink<>("late-events.txt"));
// 错误日志
errors.addSink(new LogSink<>());
// 异常告警
anomalies
.map(event -> createAlert(event))
.addSink(new AlertSink());
// 主流数据
mainDataStream
.addSink(new ElasticsearchSink<>());
env.execute("Side Output Example");
}
}
广播状态模式
public class BroadcastStateExample {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// 事件流
DataStream<Transaction> transactions = env
.addSource(new TransactionSource())
.assignTimestampsAndWatermarks(
WatermarkStrategy.<Transaction>forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withTimestampAssigner((tx, timestamp) -> tx.getTimestamp())
);
// 规则流(低吞吐量)
DataStream<Rule> rules = env
.addSource(new RuleSource())
.assignTimestampsAndWatermarks(
WatermarkStrategy.<Rule>forBoundedOutOfOrderness(Duration.ofSeconds(1))
.withTimestampAssigner((rule, timestamp) -> rule.getTimestamp())
);
// 定义广播状态描述符
MapStateDescriptor<String, Rule> ruleStateDescriptor = new MapStateDescriptor<>(
"rules",
BasicTypeInfo.STRING_TYPE_INFO,
TypeInformation.of(Rule.class)
);
// 广播规则流
BroadcastStream<Rule> ruleBroadcastStream = rules.broadcast(ruleStateDescriptor);
// 连接两个流并处理
DataStream<Alert> alerts = transactions
.keyBy(Transaction::getUserId)
.connect(ruleBroadcastStream)
.process(new DynamicAlertFunction());
alerts.print();
env.execute("Broadcast State Example");
}
}
class DynamicAlertFunction extends KeyedBroadcastProcessFunction<String, Transaction, Rule, Alert> {
private final MapStateDescriptor<String, Rule> ruleStateDescriptor =
new MapStateDescriptor<>(
"rules",
BasicTypeInfo.STRING_TYPE_INFO,
TypeInformation.of(Rule.class)
);
private transient ValueState<Double> transactionSumState;
@Override
public void open(Configuration parameters) throws Exception {
ValueStateDescriptor<Double> sumDescriptor = new ValueStateDescriptor<>(
"transaction-sum",
TypeInformation.of(Double.class)
);
transactionSumState = getRuntimeContext().getState(sumDescriptor);
}
@Override
public void processElement(Transaction transaction,
ReadOnlyContext ctx,
Collector<Alert> out) throws Exception {
// 获取当前的所有规则
ReadOnlyBroadcastState<String, Rule> broadcastState =
ctx.getBroadcastState(ruleStateDescriptor);
Double currentSum = transactionSumState.value();
if (currentSum == null) {
currentSum = 0.0;
}
currentSum += transaction.getAmount();
transactionSumState.update(currentSum);
// 检查所有规则
for (Map.Entry<String, Rule> entry : broadcastState.immutableEntries()) {
Rule rule = entry.getValue();
if (evaluateRule(rule, transaction, currentSum)) {
out.collect(new Alert(
rule.getId(),
transaction.getUserId(),
rule.getDescription(),
transaction.getTimestamp()
));
}
}
}
@Override
public void processBroadcastElement(Rule rule,
Context ctx,
Collector<Alert> out) throws Exception {
// 更新广播状态
BroadcastState<String, Rule> broadcastState =
ctx.getBroadcastState(ruleStateDescriptor);
if (rule.getAction() == Rule.Action.ADD || rule.getAction() == Rule.Action.UPDATE) {
broadcastState.put(rule.getId(), rule);
} else if (rule.getAction() == Rule.Action.DELETE) {
broadcastState.remove(rule.getId());
}
// 可以在这里输出规则变更日志
System.out.println("Rule updated: " + rule);
}
private boolean evaluateRule(Rule rule, Transaction transaction, double currentSum) {
switch (rule.getType()) {
case THRESHOLD:
return transaction.getAmount() > rule.getThreshold();
case ACCUMULATION:
return currentSum > rule.getThreshold();
case PATTERN:
return transaction.getType().matches(rule.getPattern());
default:
return false;
}
}
}
监控与调试
Metrics 系统
🔧 配置: Flink 的 Metrics 系统支持多种 Reporter,包括 Prometheus、Graphite、InfluxDB 等。
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.api.common.functions.RichMapFunction;
import org.apache.flink.metrics.Counter;
import org.apache.flink.metrics.Meter;
import org.apache.flink.metrics.Histogram;
import org.apache.flink.metrics.MeterView;
import org.apache.flink.runtime.metrics.DescriptiveStatisticsHistogram;
import org.apache.flink.metrics.Gauge;
import org.apache.flink.metrics.MetricGroup;
import org.apache.flink.streaming.api.functions.KeyedProcessFunction;
import org.apache.flink.util.Collector;
import java.util.HashSet;
import java.util.Set;
public class MetricsExample {
public static void main(String[] args) {
Configuration config = new Configuration();
// 配置 Prometheus reporter
config.setString("metrics.reporter.promgateway.class",
"org.apache.flink.metrics.prometheus.PrometheusPushGatewayReporter");
config.setString("metrics.reporter.promgateway.host", "localhost");
config.setString("metrics.reporter.promgateway.port", "9091");
config.setString("metrics.reporter.promgateway.jobName", "flink-job");
config.setString("metrics.reporter.promgateway.randomJobNameSuffix", "true");
config.setString("metrics.reporter.promgateway.deleteOnShutdown", "false");
config.setString("metrics.reporter.promgateway.interval", "10 SECONDS");
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(config);
env.addSource(new MetricSource())
.map(new MetricMapper())
.keyBy(value -> value.getKey())
.process(new MetricProcessFunction())
.print();
env.execute("Metrics Example");
}
}
class MetricMapper extends RichMapFunction<Event, Event> {
private transient Counter eventCounter;
private transient Meter eventMeter;
private transient Histogram valueHistogram;
@Override
public void open(Configuration parameters) throws Exception {
// 注册指标
eventCounter = getRuntimeContext()
.getMetricGroup()
.counter("events_processed");
eventMeter = getRuntimeContext()
.getMetricGroup()
.meter("events_per_second", new MeterView(60));
valueHistogram = getRuntimeContext()
.getMetricGroup()
.histogram("event_values", new DescriptiveStatisticsHistogram(1000));
}
@Override
public Event map(Event event) throws Exception {
// 更新指标
eventCounter.inc();
eventMeter.markEvent();
valueHistogram.update(event.getValue());
return event;
}
}
class MetricProcessFunction extends KeyedProcessFunction<String, Event, Result> {
private transient Gauge<Integer> activeKeysGauge;
private final Set<String> activeKeys = new HashSet<>();
@Override
public void open(Configuration parameters) throws Exception {
// 注册 Gauge
activeKeysGauge = getRuntimeContext()
.getMetricGroup()
.gauge("active_keys", () -> activeKeys.size());
// 自定义指标组
MetricGroup customGroup = getRuntimeContext()
.getMetricGroup()
.addGroup("custom")
.addGroup("business");
customGroup.counter("business_events");
customGroup.gauge("business_value", () -> calculateBusinessValue());
}
@Override
public void processElement(Event event, Context ctx, Collector<Result> out) {
activeKeys.add(event.getKey());
// 注册定时器清理不活跃的 key
ctx.timerService().registerProcessingTimeTimer(
ctx.timerService().currentProcessingTime() + 60000
);
out.collect(processEvent(event));
}
@Override
public void onTimer(long timestamp, OnTimerContext ctx, Collector<Result> out) {
// 清理逻辑
activeKeys.remove(ctx.getCurrentKey());
}
}
分布式追踪
public class DistributedTracingExample {
public static void main(String[] args) {
// 配置 Jaeger
Configuration config = new Configuration();
config.setString("metrics.latency.tracking", "true");
config.setString("metrics.latency.history-size", "100");
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(config);
DataStream<TracedEvent> stream = env
.addSource(new TracedSource())
.map(new TracedMapper())
.keyBy(TracedEvent::getKey)
.process(new TracedProcessFunction());
stream.print();
env.execute();
}
}
class TracedEvent {
private String traceId;
private String spanId;
private String parentSpanId;
private String key;
private String value;
private long timestamp;
// getters and setters
}
class TracedMapper extends RichMapFunction<TracedEvent, TracedEvent> {
private transient Tracer tracer;
@Override
public void open(Configuration parameters) throws Exception {
// 初始化 OpenTracing tracer
tracer = GlobalTracer.get();
}
@Override
public TracedEvent map(TracedEvent event) throws Exception {
// 创建或继续 span
SpanContext parentContext = extractContext(event);
Span span = tracer.buildSpan("map-operation")
.asChildOf(parentContext)
.withTag("event.key", event.getKey())
.start();
try (Scope scope = tracer.activateSpan(span)) {
// 业务逻辑
TracedEvent result = processEvent(event);
// 注入追踪信息
injectContext(span.context(), result);
return result;
} finally {
span.finish();
}
}
private SpanContext extractContext(TracedEvent event) {
Map<String, String> headers = new HashMap<>();
headers.put("trace-id", event.getTraceId());
headers.put("span-id", event.getSpanId());
headers.put("parent-span-id", event.getParentSpanId());
return tracer.extract(
Format.Builtin.TEXT_MAP,
new TextMapAdapter(headers)
);
}
private void injectContext(SpanContext context, TracedEvent event) {
Map<String, String> headers = new HashMap<>();
tracer.inject(
context,
Format.Builtin.TEXT_MAP,
new TextMapAdapter(headers)
);
event.setTraceId(headers.get("trace-id"));
event.setSpanId(headers.get("span-id"));
event.setParentSpanId(headers.get("parent-span-id"));
}
}
性能分析工具
public class PerformanceAnalysisExample {
public static void main(String[] args) throws Exception {
// 启用 Flame Graph
Configuration config = new Configuration();
config.setBoolean("rest.flamegraph.enabled", true);
config.setString("rest.flamegraph.directory", "/tmp/flamegraphs");
config.setInteger("rest.flamegraph.stack-depth", 100);
config.setInteger("rest.flamegraph.sample-interval", 50);
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(config);
// 启用延迟追踪
env.getConfig().setLatencyTrackingInterval(5000);
DataStream<Event> stream = env
.addSource(new PerformanceTestSource())
.name("source")
.uid("source")
.setParallelism(1)
.rebalance()
.map(new HeavyComputationMapper())
.name("heavy-computation")
.uid("heavy-computation")
.setParallelism(8)
.keyBy(Event::getKey)
.window(TumblingEventTimeWindows.of(Time.seconds(10)))
.aggregate(new PerformanceAggregator())
.name("aggregation")
.uid("aggregation")
.setParallelism(4);
// 添加性能监控 sink
stream.addSink(new PerformanceMonitoringSink())
.name("monitoring-sink")
.uid("monitoring-sink")
.setParallelism(1);
env.execute("Performance Analysis Job");
}
}
class PerformanceMonitoringSink extends RichSinkFunction<AggregatedResult> {
private transient DescriptiveStatistics latencyStats;
private transient DescriptiveStatistics throughputStats;
private transient long lastReportTime;
private transient long recordCount;
@Override
public void open(Configuration parameters) throws Exception {
latencyStats = new DescriptiveStatistics(1000);
throughputStats = new DescriptiveStatistics(60);
lastReportTime = System.currentTimeMillis();
recordCount = 0;
}
@Override
public void invoke(AggregatedResult value, Context context) throws Exception {
long currentTime = System.currentTimeMillis();
long latency = currentTime - value.getEventTime();
latencyStats.addValue(latency);
recordCount++;
// 每秒报告一次
if (currentTime - lastReportTime >= 1000) {
double throughput = recordCount * 1000.0 / (currentTime - lastReportTime);
throughputStats.addValue(throughput);
System.out.printf(
"Performance Report - Latency: p50=%.2fms, p95=%.2fms, p99=%.2fms | " +
"Throughput: %.2f records/sec | " +
"Total Records: %d%n",
latencyStats.getPercentile(50),
latencyStats.getPercentile(95),
latencyStats.getPercentile(99),
throughput,
recordCount
);
// 发送到监控系统
sendToMonitoringSystem(latencyStats, throughputStats);
lastReportTime = currentTime;
recordCount = 0;
}
}
private void sendToMonitoringSystem(
DescriptiveStatistics latency,
DescriptiveStatistics throughput) {
// 发送到 Prometheus/Grafana 等监控系统
}
}
最佳实践
本节总结了在生产环境中使用 Flink DataStream API 的最佳实践。
1. 状态管理最佳实践
public class StateManagementBestPractices {
// 使用状态 TTL 防止状态无限增长
public static class TTLProcessFunction extends KeyedProcessFunction<String, Event, Result> {
private transient ValueState<UserSession> sessionState;
@Override
public void open(Configuration parameters) throws Exception {
ValueStateDescriptor<UserSession> descriptor = new ValueStateDescriptor<>(
"user-session",
TypeInformation.of(UserSession.class)
);
// 配置 TTL
StateTtlConfig ttlConfig = StateTtlConfig
.newBuilder(Time.hours(24))
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
.cleanupInRocksdbCompactFilter(1000)
.build();
descriptor.enableTimeToLive(ttlConfig);
sessionState = getRuntimeContext().getState(descriptor);
}
}
// 使用 MapState 替代嵌套的 ValueState
public static class OptimizedStateFunction extends KeyedProcessFunction<String, Event, Result> {
// 不好的做法
// private transient ValueState<Map<String, List<Event>>> complexState;
// 好的做法
private transient MapState<String, List<Event>> optimizedState;
@Override
public void open(Configuration parameters) throws Exception {
MapStateDescriptor<String, List<Event>> descriptor = new MapStateDescriptor<>(
"events-by-type",
TypeInformation.of(String.class),
TypeInformation.of(new TypeHint<List<Event>>() {})
);
optimizedState = getRuntimeContext().getMapState(descriptor);
}
}
}
2. 性能优化最佳实践
public class PerformanceOptimizationBestPractices {
// 避免频繁的序列化/反序列化
public static class OptimizedMapper extends RichMapFunction<String, Event> {
private transient ObjectMapper objectMapper;
@Override
public void open(Configuration parameters) throws Exception {
// 重用 ObjectMapper 实例
objectMapper = new ObjectMapper();
objectMapper.configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false);
}
@Override
public Event map(String value) throws Exception {
return objectMapper.readValue(value, Event.class);
}
}
// 合理使用算子链
public static void configureOperatorChaining(DataStream<Event> stream) {
stream
.filter(event -> event.getType() != null) // 轻量级操作,保持链接
.map(event -> enrichEvent(event))
.startNewChain() // 重量级操作,开始新链
.keyBy(Event::getUserId)
.process(new HeavyProcessFunction())
.disableChaining() // 禁用链接,独立运行
.addSink(new CustomSink());
}
}
3. 容错性最佳实践
public class FaultToleranceBestPractices {
// 实现幂等性 Sink
public static class IdempotentSink extends RichSinkFunction<Event>
implements CheckpointedFunction, CheckpointListener {
private transient Set<String> pendingEvents;
private transient ListState<String> checkpointedState;
@Override
public void invoke(Event event, Context context) throws Exception {
String eventId = event.getId();
// 检查是否已处理
if (!isProcessed(eventId)) {
processEvent(event);
pendingEvents.add(eventId);
}
}
@Override
public void notifyCheckpointComplete(long checkpointId) throws Exception {
// 检查点完成后,确认待处理事件
for (String eventId : pendingEvents) {
markAsProcessed(eventId);
}
pendingEvents.clear();
}
@Override
public void snapshotState(FunctionSnapshotContext context) throws Exception {
checkpointedState.clear();
for (String eventId : pendingEvents) {
checkpointedState.add(eventId);
}
}
@Override
public void initializeState(FunctionInitializationContext context) throws Exception {
ListStateDescriptor<String> descriptor = new ListStateDescriptor<>(
"pending-events",
TypeInformation.of(String.class)
);
checkpointedState = context.getOperatorStateStore().getListState(descriptor);
if (context.isRestored()) {
pendingEvents = new HashSet<>();
for (String eventId : checkpointedState.get()) {
pendingEvents.add(eventId);
}
}
}
}
}
常见问题
本节整理了使用 Flink DataStream API 高级特性时的常见问题和解决方案。
Q1: 如何处理大状态?
解决方案:
// 1. 使用 RocksDB 状态后端
env.setStateBackend(new RocksDBStateBackend("hdfs://checkpoints"));
// 2. 启用增量检查点
backend.setEnableIncrementalCheckpointing(true);
// 3. 配置 RocksDB 参数
backend.setOptions(new OptionsFactory() {
@Override
public DBOptions createDBOptions(DBOptions currentOptions,
Collection<AutoCloseable> handlesToClose) {
return currentOptions
.setMaxBackgroundJobs(4)
.setMaxOpenFiles(-1);
}
});
// 4. 使用状态 TTL
StateTtlConfig ttlConfig = StateTtlConfig
.newBuilder(Time.hours(24))
.cleanupFullSnapshot()
.build();
Q2: 如何优化检查点性能?
解决方案:
// 1. 使用非对齐检查点
env.getCheckpointConfig().enableUnalignedCheckpoints();
// 2. 调整检查点间隔
env.enableCheckpointing(60000); // 1分钟
// 3. 设置合理的超时
env.getCheckpointConfig().setCheckpointTimeout(120000); // 2分钟
// 4. 限制并发检查点
env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);
Q3: 如何处理数据倾斜?
解决方案:
// 1. 使用自定义分区器
stream.partitionCustom(new Partitioner<String>() {
@Override
public int partition(String key, int numPartitions) {
// 实现均匀分布的逻辑
return (key.hashCode() & Integer.MAX_VALUE) % numPartitions;
}
}, Event::getKey);
// 2. 二次聚合
stream
.keyBy(event -> event.getKey() + "_" + random.nextInt(10))
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.aggregate(new LocalAggregator())
.keyBy(result -> result.getOriginalKey())
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.aggregate(new GlobalAggregator());
// 3. 使用 KeyGroup 优化
DataStream<Event> rebalanced = stream
.keyBy(new KeySelector<Event, Integer>() {
@Override
public Integer getKey(Event event) throws Exception {
return KeyGroupRangeAssignment.assignKeyToParallelOperator(
event.getKey(),
env.getParallelism(),
env.getParallelism()
);
}
});
总结
Flink DataStream API 的高级用法涵盖了容错机制、性能优化、复杂事件处理等多个方面。掌握这些高级特性能够帮助开发者构建更加健壮、高效的流处理应用。
🚫 错误: 不要在生产环境直接使用本文的示例代码,需要根据实际场景进行调整和测试。
核心要点
- 容错机制:合理配置 Checkpoint 和 Savepoint,确保数据一致性
- 性能优化:从并行度、状态后端、网络缓冲等多维度优化
- Table & SQL:利用 SQL 简化流处理开发
- 高级连接器:CDC、数据湖集成扩展应用场景
- 监控调试:完善的监控体系保障生产稳定性
相关文章
DataStream API 系列
- Flink DataStream API 中的算子 - 基础算子使用
- Flink DataStream API 窗口操作 - 窗口机制详解
核心概念系列
- Flink 状态管理详解 - 状态类型和使用方法
- Flink 时间和水位线 - 时间语义深入理解
- Flink 容错机制 - Checkpoint 和 Savepoint 详解
进阶学习
- Flink SQL 开发指南 - SQL API 使用
- Flink 性能调优实战 - 性能优化技巧
- Flink CDC 实战 - CDC 连接器详解