Flink CDC案例

wen java案例 3

Flink CDC 实战案例详解

基础概念

Flink CDC(Change Data Capture)是Apache Flink提供的一个强大的数据捕获框架,能够实时捕获数据库中的数据变更(插入、更新、删除),并将其转换为流式事件进行处理。

Flink CDC案例

经典案例:MySQL实时同步到Kafka

场景描述

将MySQL中订单表的变更实时同步到Kafka,供下游系统消费。

环境准备

<!-- pom.xml 依赖配置 -->
<dependencies>
    <!-- Flink CDC 依赖 -->
    <dependency>
        <groupId>com.ververica</groupId>
        <artifactId>flink-connector-mysql-cdc</artifactId>
        <version>2.4.0</version>
    </dependency>
    <!-- Flink Kafka 连接器 -->
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-connector-kafka</artifactId>
        <version>1.17.0</version>
    </dependency>
    <!-- Flink Stream API -->
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-streaming-java</artifactId>
        <version>1.17.0</version>
    </dependency>
</dependencies>

MySQL数据准备

-- 创建数据库
CREATE DATABASE flink_cdc_demo;
USE flink_cdc_demo;
-- 创建订单表
CREATE TABLE orders (
    id INT PRIMARY KEY AUTO_INCREMENT,
    order_no VARCHAR(50) NOT NULL,
    user_id INT NOT NULL,
    product_name VARCHAR(100),
    amount DECIMAL(10,2),
    status VARCHAR(20),
    create_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
    update_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP
);
-- 插入测试数据
INSERT INTO orders (order_no, user_id, product_name, amount, status) 
VALUES 
('ORD-001', 1001, 'iPhone 15', 6999.00, 'CREATED'),
('ORD-002', 1002, 'MacBook Pro', 15999.00, 'PAID'),
('ORD-003', 1003, 'AirPods', 1299.00, 'SHIPPED');
-- 创建用户表
CREATE TABLE users (
    id INT PRIMARY KEY AUTO_INCREMENT,
    name VARCHAR(50),
    email VARCHAR(100)
);
INSERT INTO users (name, email) VALUES ('张三', 'zhangsan@example.com');

核心代码实现

import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.ProcessFunction;
import org.apache.flink.util.Collector;
import com.ververica.cdc.connectors.mysql.source.MySqlSource;
import com.ververica.cdc.debezium.JsonDebeziumDeserializationSchema;
import org.apache.flink.connector.kafka.sink.KafkaSink;
import org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema;
import org.apache.flink.connector.kafka.sink.KafkaSinkBuilder;
public class MySqlCDCToKafka {
    public static void main(String[] args) throws Exception {
        // 1. 创建执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);
        // 启用Checkpoint
        env.enableCheckpointing(5000);
        // 2. 创建MySQL CDC Source
        MySqlSource<String> mySqlSource = MySqlSource.<String>builder()
            .hostname("localhost")
            .port(3306)
            .databaseList("flink_cdc_demo")  // 监控的数据库
            .tableList("flink_cdc_demo.orders")  // 监控的表
            .username("root")
            .password("password")
            .deserializer(new JsonDebeziumDeserializationSchema()) // 转换为JSON
            .startupOptions(StartupOptions.initial())  // 启动模式:初始化
            .build();
        // 3. 创建Kafka Sink
        KafkaSink<String> kafkaSink = KafkaSink.<String>builder()
            .setBootstrapServers("localhost:9092")
            .setRecordSerializer(KafkaRecordSerializationSchema.builder()
                .setTopic("cdc-orders")
                .setValueSerializationSchema(new SimpleStringSchema())
                .build()
            )
            .setDeliveryGuarantee(DeliveryGuarantee.AT_LEAST_ONCE)
            .build();
        // 4. 构建数据流
        DataStream<String> stream = env
            .fromSource(mySqlSource, WatermarkStrategy.noWatermarks(), "MySQL CDC Source")
            .name("MySQL CDC Source");
        // 5. 数据处理(可选:解析和转换)
        DataStream<String> processedStream = stream
            .process(new ProcessFunction<String, String>() {
                @Override
                public void processElement(String value, Context ctx, Collector<String> out) throws Exception {
                    // 这里可以添加自定义处理逻辑
                    // 数据清洗、格式转换、业务逻辑处理等
                    out.collect(value);
                }
            })
            .name("Data Processing");
        // 6. 写入Kafka
        processedStream.sinkTo(kafkaSink).name("Kafka Sink");
        // 7. 执行任务
        env.execute("MySQL CDC to Kafka");
    }
}

配置其他启动模式

// 多种启动模式配置
// 1. 初始化模式(默认):先读取现有数据,然后继续读取变更
StartupOptions.initial()
// 2. 最早模式:从最早的binlog开始读取
StartupOptions.earliest()
// 3. 最新模式:只从当前时间开始读取变更
StartupOptions.latest()
// 4. 指定时间戳
StartupOptions.timestamp(1700000000000L)
// 5. 指定偏移量
StartupOptions.specificOffset("mysql-bin.000001", 4L, 12345L)

进阶案例:多表关联与状态管理

import org.apache.flink.api.common.state.MapState;
import org.apache.flink.api.common.state.MapStateDescriptor;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.co.KeyedBroadcastProcessFunction;
import org.apache.flink.util.Collector;
public class MultiTableJoinCDC {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(1);
        // 1. 创建订单CDC Source
        MySqlSource<String> orderSource = MySqlSource.<String>builder()
            .hostname("localhost")
            .port(3306)
            .databaseList("flink_cdc_demo")
            .tableList("flink_cdc_demo.orders")
            .username("root")
            .password("password")
            .deserializer(new JsonDebeziumDeserializationSchema())
            .build();
        // 2. 创建用户CDC Source
        MySqlSource<String> userSource = MySqlSource.<String>builder()
            .hostname("localhost")
            .port(3306)
            .databaseList("flink_cdc_demo")
            .tableList("flink_cdc_demo.users")
            .username("root")
            .password("password")
            .deserializer(new JsonDebeziumDeserializationSchema())
            .build();
        // 3. 获取数据流
        DataStream<String> orderStream = env.fromSource(
            orderSource, 
            WatermarkStrategy.noWatermarks(), 
            "Order CDC Source"
        );
        DataStream<String> userStream = env.fromSource(
            userSource, 
            WatermarkStrategy.noWatermarks(), 
            "User CDC Source"
        );
        // 4. 使用Connect合并两个流并关联
        DataStream<String> result = orderStream
            .connect(userStream)
            .process(new OrderUserJoinFunction())
            .name("Order-User Join");
        // 5. 输出结果
        result.print();
        env.execute("Multi-table Join CDC");
    }
    // 自定义连接函数
    public static class OrderUserJoinFunction 
        extends KeyedBroadcastProcessFunction<String, String, String, String> {
        private MapState<String, String> userState;
        @Override
        public void open(Configuration parameters) {
            MapStateDescriptor<String, String> userDescriptor = 
                new MapStateDescriptor<>("users", String.class, String.class);
            userState = getRuntimeContext().getMapState(userDescriptor);
        }
        @Override
        public void processElement(String orderJson, ReadOnlyContext ctx, 
                                  Collector<String> out) throws Exception {
            // 解析订单JSON,获取user_id
            String userId = extractUserId(orderJson);
            String userInfo = userState.get(userId);
            if (userInfo != null) {
                // 关联用户信息
                String enrichedOrder = enrichOrderWithUser(orderJson, userInfo);
                out.collect(enrichedOrder);
            }
        }
        @Override
        public void processBroadcastElement(String userJson, Context ctx, 
                                           Collector<String> out) throws Exception {
            // 更新用户状态
            String userId = extractUserId(userJson);
            userState.put(userId, userJson);
        }
    }
}

实际生产案例:实时数仓同步

// 完整的生产级配置示例
public class ProductionCDCApplication {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        // 生产环境配置
        Configuration config = new Configuration();
        config.setInteger("taskmanager.numberOfTaskSlots", 4);
        env.configure(config);
        // 启用Checkpoint
        CheckpointConfig checkpointConfig = env.getCheckpointConfig();
        checkpointConfig.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
        checkpointConfig.setCheckpointInterval(60000); // 1分钟
        checkpointConfig.setCheckpointTimeout(60000);
        checkpointConfig.setMaxConcurrentCheckpoints(1);
        checkpointConfig.setMinPauseBetweenCheckpoints(5000);
        // 设置状态后端
        env.setStateBackend(new HashMapStateBackend());
        env.getCheckpointConfig().setCheckpointStorage("hdfs://namenode:8020/flink/checkpoints");
        // 创建多表监控的CDC Source
        MySqlSource<String> source = MySqlSource.<String>builder()
            .hostname("mysql-master")
            .port(3306)
            .databaseList("business_db")  // 业务数据库
            .tableList("business_db.orders,business_db.payments,business_db.shipments")  // 多表
            .username("cdc_user")
            .password("cdc_password")
            .serverId("5400-5404")  // 分配多个server ID用于并行读取
            .debeziumProperties(new HashMap<String, String>() {{
                put("snapshot.locking.mode", "none");  // 无锁快照
                put("database.server.name", "business_db");
                put("include.schema.changes", "true");
            }})
            .deserializer(new JsonDebeziumDeserializationSchema())
            .startupOptions(StartupOptions.latest())
            .build();
        DataStream<String> stream = env.fromSource(
            source, 
            WatermarkStrategy.noWatermarks(), 
            "Production CDC Source"
        )
        .setParallelism(4)  // 设置并行度
        .rebalance();
        // 写入多个目标
        // 1. 写入Kafka
        stream.sinkTo(createKafkaSink("cdc-topic"));
        // 2. 写入StarRocks/ClickHouse
        stream.addSink(createStarRocksSink());
        // 3. 本地调试输出
        if (isDebugMode()) {
            stream.print();
        }
        env.execute("Production CDC Sync Job");
    }
}

常见问题与解决方案

数据一致性

// 确保至少一次语义
CheckpointConfig config = env.getCheckpointConfig();
config.setCheckpointingMode(CheckpointingMode.AT_LEAST_ONCE);
// 或精确一次性
config.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);

性能优化

# flink-conf.yaml 配置
# 增加Socket接收缓冲区
taskmanager.network.memory.buffer-debloat.enabled: true
# 优化网络缓冲区
taskmanager.memory.network.min: 64mb
taskmanager.memory.network.max: 128mb
# 并行度配置
parallelism.default: 4

监控告警

// 添加Metrics监控
stream
    .map(new RichMapFunction<String, String>() {
        private transient Counter counter;
        @Override
        public void open(Configuration parameters) {
            counter = getRuntimeContext()
                .getMetricGroup()
                .counter("cdc_records_count");
        }
        @Override
        public String map(String value) throws Exception {
            counter.inc();
            return value;
        }
    });

CDC数据格式示例

// MySQL CDC 输出的JSON格式
{
    "before": {
        "id": 1,
        "order_no": "ORD-001",
        "status": "CREATED"
    },
    "after": {
        "id": 1,
        "order_no": "ORD-001", 
        "status": "PAID"
    },
    "source": {
        "db": "flink_cdc_demo",
        "table": "orders",
        "server_id": 123,
        "ts_sec": 1700000000
    },
    "op": "u",  // c:创建, u:更新, d:删除, r:快照读取
    "ts_ms": 1700000000000
}

最佳实践总结

  1. 启动策略选择

    • 首次使用:initial模式全量+增量
    • 已有offset:latest模式只读增量
  2. 性能调优

    • 增大server.id分配范围
    • 调整debezium.snapshot.fetch.size
    • 启用并行读取
  3. 数据质量

    • 实现数据校验机制
    • 添加重试和死信队列
    • 监控延迟和吞吐量
  4. 运维建议

    • 定期检查binlog保留时间
    • 监控MySQL性能指标
    • 做好Flink状态备份

这个案例涵盖了从基础到生产的Flink CDC应用,可以根据实际需求选择合适的实现方案。

抱歉,评论功能暂时关闭!