全部笔记All notes

Flink DataStream API 高级用法详解

阅读 8m 26s8m 26s read

概述

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 提供了实时捕获数据库变更的能力,支持 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 的高级用法涵盖了容错机制、性能优化、复杂事件处理等多个方面。掌握这些高级特性能够帮助开发者构建更加健壮、高效的流处理应用。

🚫 错误: 不要在生产环境直接使用本文的示例代码,需要根据实际场景进行调整和测试。

核心要点

  1. 容错机制:合理配置 Checkpoint 和 Savepoint,确保数据一致性
  2. 性能优化:从并行度、状态后端、网络缓冲等多维度优化
  3. Table & SQL:利用 SQL 简化流处理开发
  4. 高级连接器:CDC、数据湖集成扩展应用场景
  5. 监控调试:完善的监控体系保障生产稳定性

相关文章

DataStream API 系列

核心概念系列

  • Flink 状态管理详解 - 状态类型和使用方法
  • Flink 时间和水位线 - 时间语义深入理解
  • Flink 容错机制 - Checkpoint 和 Savepoint 详解

进阶学习

  • Flink SQL 开发指南 - SQL API 使用
  • Flink 性能调优实战 - 性能优化技巧
  • Flink CDC 实战 - CDC 连接器详解