全部笔记All notes

Flink CDC 实时数据同步完全指南

阅读 10m 59s10m 59s read

概述

Flink CDC (Change Data Capture) 是基于 Apache Flink 的数据变更捕获框架,能够实时捕获数据库的增删改操作,并将变更事件流式传输到下游系统。它提供了一站式的实时数据同步解决方案。

核心特性

  • 全量 + 增量同步: 支持历史数据全量同步和实时增量同步
  • Exactly-Once 语义: 保证数据不丢失、不重复
  • 无锁读取: 不影响在线业务,对源库压力小
  • 分布式架构: 支持水平扩展,处理大规模数据
  • Schema 演进: 自动感知并同步表结构变更
  • 断点续传: 支持从中断位置继续同步

应用场景

场景描述优势
实时数据同步数据库之间的实时同步低延迟、高可靠
实时数仓构建构建实时数据仓库实时 ETL、数据清洗
缓存更新数据库到缓存的同步自动失效、实时更新
搜索索引构建数据库到搜索引擎同步近实时搜索
微服务数据同步跨服务数据共享解耦服务、数据一致性

💡 提示: Flink CDC 相比传统的批量同步工具,具有更低的延迟和更小的资源开销。

CDC 基础概念

什么是 CDC

Change Data Capture(变更数据捕获)是一种用于识别和捕获数据库中数据变更的技术:

传统批量同步 vs CDC 实时同步

批量同步:
┌─────────┐  定时全量  ┌─────────┐
│ Source  │ ────────> │ Target  │
│   DB    │   (高延迟) │   DB    │
└─────────┘           └─────────┘

CDC 同步:
┌─────────┐  实时增量  ┌─────────┐
│ Source  │ ────────> │ Target  │
│   DB    │ (低延迟)  │   DB    │
└─────────┘           └─────────┘
     ↓
  Binlog/WAL

CDC 实现方式对比

方式原理优点缺点
基于查询定期查询变更实现简单延迟高、压力大
基于触发器数据库触发器实时性好侵入性强、性能影响
基于日志解析数据库日志低侵入、低延迟实现复杂

Flink CDC 采用基于日志的方式,通过读取数据库的事务日志(如 MySQL 的 binlog、PostgreSQL 的 WAL)来捕获数据变更。

整体架构

┌─────────────────────────────────────────────────────────────┐
│                      Flink CDC 架构                          │
├─────────────────────────────────────────────────────────────┤
│                                                             │
│  数据源层                                                    │
│  ┌────────┐ ┌────────┐ ┌────────┐ ┌────────┐              │
│  │ MySQL  │ │  PG    │ │ Oracle │ │MongoDB │              │
│  └───┬────┘ └───┬────┘ └───┬────┘ └───┬────┘              │
│      │          │          │          │                     │
│      ▼          ▼          ▼          ▼                     │
│  ┌─────────────────────────────────────────┐               │
│  │        CDC Source Connector              │               │
│  │  ┌─────────┐ ┌─────────┐ ┌─────────┐   │               │
│  │  │Debezium │ │ Reader  │ │Snapshot │   │               │
│  │  │ Engine  │ │ Thread  │ │ Thread  │   │               │
│  │  └─────────┘ └─────────┘ └─────────┘   │               │
│  └─────────────────────────────────────────┘               │
│                        │                                     │
│                        ▼                                     │
│  ┌─────────────────────────────────────────┐               │
│  │           Flink Runtime                  │               │
│  │  ┌─────────┐ ┌─────────┐ ┌─────────┐   │               │
│  │  │  State  │ │Checkpoint│ │Watermark│   │               │
│  │  │ Backend │ │         │ │         │   │               │
│  │  └─────────┘ └─────────┘ └─────────┘   │               │
│  └─────────────────────────────────────────┘               │
│                        │                                     │
│                        ▼                                     │
│  数据处理层                                                  │
│  ┌─────────┐ ┌─────────┐ ┌─────────┐                      │
│  │Transform│ │ Filter  │ │ Route   │                      │
│  └─────────┘ └─────────┘ └─────────┘                      │
│                        │                                     │
│                        ▼                                     │
│  目标系统层                                                  │
│  ┌────────┐ ┌────────┐ ┌────────┐ ┌────────┐              │
│  │ Kafka  │ │  ES    │ │ HBase  │ │Iceberg │              │
│  └────────┘ └────────┘ └────────┘ └────────┘              │
│                                                             │
└─────────────────────────────────────────────────────────────┘

核心组件

1. Debezium Engine

负责从数据库日志中读取变更事件:

  • 管理数据库连接
  • 解析二进制日志
  • 转换为统一的事件格式

2. Source Connector

Flink 数据源连接器:

  • 管理读取位置(offset)
  • 协调全量和增量读取
  • 处理故障恢复

3. Snapshot Reader

全量数据读取器:

  • 并行读取历史数据
  • 生成一致性快照
  • 无锁读取实现

4. Stream Reader

增量数据读取器:

  • 实时读取变更日志
  • 维护读取位置
  • 处理 Schema 变更

快速开始

环境准备

<!-- Maven 依赖 -->
<dependencies>
    <!-- Flink CDC -->
    <dependency>
        <groupId>com.ververica</groupId>
        <artifactId>flink-sql-connector-mysql-cdc</artifactId>
        <version>2.4.0</version>
    </dependency>
    
    <!-- Flink 依赖 -->
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-streaming-java</artifactId>
        <version>1.17.0</version>
    </dependency>
    
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-table-api-java-bridge</artifactId>
        <version>1.17.0</version>
    </dependency>
</dependencies>

基础示例

1. DataStream API 方式

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;

public class MySqlCDCExample {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        
        // 设置并行度
        env.setParallelism(4);
        
        // 创建 MySQL CDC Source
        MySqlSource<String> mySqlSource = MySqlSource.<String>builder()
                .hostname("localhost")
                .port(3306)
                .databaseList("test_db") // 监控的数据库
                .tableList("test_db.users", "test_db.orders") // 监控的表
                .username("root")
                .password("password")
                .startupOptions(StartupOptions.initial()) // 从最早位置开始读取
                .deserializer(new JsonDebeziumDeserializationSchema()) // JSON 格式
                .build();
        
        // 使用 CDC Source
        env.fromSource(mySqlSource, WatermarkStrategy.noWatermarks(), "MySQL Source")
           .print();
        
        env.execute("MySQL CDC Example");
    }
}

2. Table API 方式

import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;

public class MySqlCDCTableExample {
    public static void main(String[] args) {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
        
        // 创建 CDC 源表
        tableEnv.executeSql(
            "CREATE TABLE users_source (" +
            "  id INT," +
            "  name STRING," +
            "  age INT," +
            "  email STRING," +
            "  update_time TIMESTAMP(3)," +
            "  PRIMARY KEY (id) NOT ENFORCED" +
            ") WITH (" +
            "  'connector' = 'mysql-cdc'," +
            "  'hostname' = 'localhost'," +
            "  'port' = '3306'," +
            "  'username' = 'root'," +
            "  'password' = 'password'," +
            "  'database-name' = 'test_db'," +
            "  'table-name' = 'users'" +
            ")"
        );
        
        // 创建目标表
        tableEnv.executeSql(
            "CREATE TABLE users_sink (" +
            "  id INT," +
            "  name STRING," +
            "  age INT," +
            "  email STRING," +
            "  update_time TIMESTAMP(3)," +
            "  PRIMARY KEY (id) NOT ENFORCED" +
            ") WITH (" +
            "  'connector' = 'elasticsearch-7'," +
            "  'hosts' = 'http://localhost:9200'," +
            "  'index' = 'users'" +
            ")"
        );
        
        // 同步数据
        tableEnv.executeSql("INSERT INTO users_sink SELECT * FROM users_source");
    }
}

支持的数据源

MySQL CDC

1. 前置条件

-- 检查 binlog 是否开启
SHOW VARIABLES LIKE 'log_bin';

-- 查看 binlog 格式(需要是 ROW)
SHOW VARIABLES LIKE 'binlog_format';

-- 配置 my.cnf
[mysqld]
log-bin=mysql-bin
binlog-format=ROW
binlog-row-image=FULL
expire_logs_days=7

2. 高级配置

public class MySqlCDCAdvanced {
    public static MySqlSource<String> createAdvancedMySqlSource() {
        Properties jdbcProperties = new Properties();
        jdbcProperties.put("useSSL", "false");
        jdbcProperties.put("allowPublicKeyRetrieval", "true");
        
        Properties debeziumProperties = new Properties();
        // 设置时区
        debeziumProperties.put("database.serverTimezone", "Asia/Shanghai");
        // 处理 decimal 类型
        debeziumProperties.put("decimal.handling.mode", "string");
        
        return MySqlSource.<String>builder()
                .hostname("localhost")
                .port(3306)
                .databaseList("db1", "db2") // 多数据库
                .tableList("db1.table1", "db2.*") // 支持正则
                .username("root")
                .password("password")
                .serverId("5400-5404") // 服务器 ID 范围
                .splitSize(8096) // 全量读取时的分片大小
                .fetchSize(1024) // JDBC fetch size
                .connectTimeout(Duration.ofSeconds(30))
                .connectionPoolSize(20) // 连接池大小
                .jdbcProperties(jdbcProperties)
                .debeziumProperties(debeziumProperties)
                .startupOptions(StartupOptions.specificOffset("mysql-bin.000003", 4))
                .deserializer(new JsonDebeziumDeserializationSchema())
                .build();
    }
}

3. 自定义反序列化

public class CustomDeserializationSchema implements DebeziumDeserializationSchema<RowData> {
    
    @Override
    public void deserialize(SourceRecord record, Collector<RowData> out) throws Exception {
        // 获取操作类型
        Envelope.Operation op = Envelope.operationFor(record);
        
        // 获取数据
        Struct value = (Struct) record.value();
        Struct source = value.getStruct("source");
        String database = source.getString("db");
        String table = source.getString("table");
        
        // 根据操作类型处理
        switch (op) {
            case CREATE:
            case READ:
                Struct after = value.getStruct("after");
                RowData insertRow = extractRowData(after, RowKind.INSERT);
                out.collect(insertRow);
                break;
                
            case UPDATE:
                Struct before = value.getStruct("before");
                Struct afterUpdate = value.getStruct("after");
                
                // 发送删除和插入事件(CDC 语义)
                if (before != null) {
                    RowData deleteRow = extractRowData(before, RowKind.DELETE);
                    out.collect(deleteRow);
                }
                RowData insertRowUpdate = extractRowData(afterUpdate, RowKind.INSERT);
                out.collect(insertRowUpdate);
                break;
                
            case DELETE:
                Struct beforeDelete = value.getStruct("before");
                RowData deleteRow = extractRowData(beforeDelete, RowKind.DELETE);
                out.collect(deleteRow);
                break;
        }
    }
    
    private RowData extractRowData(Struct struct, RowKind rowKind) {
        // 实现数据提取逻辑
        // ...
    }
    
    @Override
    public TypeInformation<RowData> getProducedType() {
        return TypeInformation.of(RowData.class);
    }
}

PostgreSQL CDC

1. 配置要求

-- postgresql.conf
wal_level = logical
max_replication_slots = 4
max_wal_senders = 4

-- 创建发布
CREATE PUBLICATION flink_publication FOR ALL TABLES;

-- 创建复制槽
SELECT * FROM pg_create_logical_replication_slot('flink_slot', 'pgoutput');

2. 使用示例

public class PostgreSQLCDCExample {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        
        PostgreSQLSource<String> postgresSource = PostgreSQLSource.<String>builder()
                .hostname("localhost")
                .port(5432)
                .database("postgres")
                .schemaList("public")
                .tableList("public.users", "public.orders")
                .username("postgres")
                .password("password")
                .slotName("flink_slot")
                .decodingPluginName("pgoutput") // 或 wal2json
                .deserializer(new JsonDebeziumDeserializationSchema())
                .build();
        
        env.fromSource(postgresSource, WatermarkStrategy.noWatermarks(), "PostgreSQL Source")
           .print();
        
        env.execute("PostgreSQL CDC Example");
    }
}

Oracle CDC

1. 配置 Oracle

-- 启用归档日志
ALTER DATABASE ARCHIVELOG;

-- 启用补充日志
ALTER DATABASE ADD SUPPLEMENTAL LOG DATA (ALL) COLUMNS;

-- 创建用户并授权
CREATE USER flinkuser IDENTIFIED BY flinkpw;
GRANT CREATE SESSION TO flinkuser;
GRANT SELECT ANY TABLE TO flinkuser;
GRANT SELECT_CATALOG_ROLE TO flinkuser;
GRANT EXECUTE_CATALOG_ROLE TO flinkuser;
GRANT SELECT ANY TRANSACTION TO flinkuser;
GRANT LOGMINING TO flinkuser;

2. 使用示例

public class OracleCDCExample {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        
        OracleSource<String> oracleSource = OracleSource.<String>builder()
                .hostname("localhost")
                .port(1521)
                .database("ORCL") // SID
                .schemaList("FLINKUSER")
                .tableList("FLINKUSER.USERS")
                .username("flinkuser")
                .password("flinkpw")
                .deserializer(new JsonDebeziumDeserializationSchema())
                .build();
        
        env.fromSource(oracleSource, WatermarkStrategy.noWatermarks(), "Oracle Source")
           .print();
        
        env.execute("Oracle CDC Example");
    }
}

MongoDB CDC

1. 前置条件

// MongoDB 必须是副本集模式
// 检查副本集状态
rs.status()

// 如果是单机,转换为副本集
rs.initiate()

2. 使用示例

public class MongoDBCDCExample {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        
        MongoDBSource<String> mongoSource = MongoDBSource.<String>builder()
                .hosts("localhost:27017")
                .username("flinkuser")
                .password("flinkpw")
                .databaseList("test") // 支持正则
                .collectionList("test.users", "test.orders")
                .deserializer(new JsonDebeziumDeserializationSchema())
                .build();
        
        env.fromSource(mongoSource, WatermarkStrategy.noWatermarks(), "MongoDB Source")
           .print();
        
        env.execute("MongoDB CDC Example");
    }
}

SQL Server CDC

1. 启用 CDC

-- 启用数据库级别 CDC
EXEC sys.sp_cdc_enable_db;

-- 启用表级别 CDC
EXEC sys.sp_cdc_enable_table
    @source_schema = N'dbo',
    @source_name   = N'users',
    @role_name     = NULL;

2. 使用示例

public class SqlServerCDCExample {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        
        SqlServerSource<String> sqlServerSource = SqlServerSource.<String>builder()
                .hostname("localhost")
                .port(1433)
                .database("testdb")
                .tableList("dbo.users")
                .username("sa")
                .password("Password!")
                .deserializer(new JsonDebeziumDeserializationSchema())
                .build();
        
        env.fromSource(sqlServerSource, WatermarkStrategy.noWatermarks(), "SQL Server Source")
           .print();
        
        env.execute("SQL Server CDC Example");
    }
}

高级特性

全量与增量同步

1. 启动模式配置

public class StartupOptionsExample {
    public static void demonstrateStartupOptions() {
        // 1. 从最早位置开始(包括全量)
        StartupOptions initial = StartupOptions.initial();
        
        // 2. 从最新位置开始(仅增量)
        StartupOptions latest = StartupOptions.latest();
        
        // 3. 从指定位置开始
        StartupOptions specificOffset = StartupOptions.specificOffset("mysql-bin.000003", 4);
        
        // 4. 从指定时间戳开始
        StartupOptions timestamp = StartupOptions.timestamp(1000L);
        
        // 5. 跳过全量阶段,从变更流开始
        StartupOptions snapshot = StartupOptions.snapshot();
    }
}

2. 分片并行读取

public class ParallelSnapshotReading {
    public static MySqlSource<String> createParallelSource() {
        return MySqlSource.<String>builder()
                .hostname("localhost")
                .port(3306)
                .databaseList("large_db")
                .tableList("large_db.huge_table")
                .username("root")
                .password("password")
                .splitSize(8096) // 每个分片的大小
                .splitMetaGroupSize(1000) // 分片元数据组大小
                .fetchSize(1024) // JDBC 批量大小
                .connectTimeout(Duration.ofSeconds(30))
                .deserializer(new JsonDebeziumDeserializationSchema())
                .build();
    }
}

Schema 演进

1. 自动 Schema 变更感知

public class SchemaEvolutionExample {
    public static void handleSchemaEvolution() {
        // 配置 Debezium 属性以处理 Schema 变更
        Properties debeziumProperties = new Properties();
        // 包含 Schema 变更事件
        debeziumProperties.put("include.schema.changes", "true");
        // 存储 Schema 历史
        debeziumProperties.put("database.history.store.only.monitored.tables.ddl", "true");
        
        MySqlSource<String> source = MySqlSource.<String>builder()
                .hostname("localhost")
                .port(3306)
                .databaseList("test")
                .tableList("test.*")
                .username("root")
                .password("password")
                .debeziumProperties(debeziumProperties)
                .deserializer(new SchemaAwareDeserializer())
                .build();
    }
}

// 自定义 Schema 感知反序列化器
public class SchemaAwareDeserializer implements DebeziumDeserializationSchema<String> {
    @Override
    public void deserialize(SourceRecord record, Collector<String> out) throws Exception {
        // 检查是否是 Schema 变更事件
        if (record.topic().endsWith(".DDL")) {
            handleSchemaChange(record);
            return;
        }
        
        // 处理数据变更事件
        String result = convertToJson(record);
        out.collect(result);
    }
    
    private void handleSchemaChange(SourceRecord record) {
        // 处理 DDL 变更
        Struct value = (Struct) record.value();
        String ddl = value.getString("ddl");
        System.out.println("Schema change detected: " + ddl);
        
        // 可以将 Schema 变更信息发送到元数据管理系统
    }
}

分布式快照

1. 一致性保证

public class ConsistentSnapshotExample {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        
        // 启用 Checkpoint 保证一致性
        env.enableCheckpointing(10000);
        env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
        
        MySqlSource<String> source = MySqlSource.<String>builder()
                .hostname("localhost")
                .port(3306)
                .databaseList("test")
                .tableList("test.*")
                .username("root")
                .password("password")
                // 使用一致性快照读取
                .startupOptions(StartupOptions.initial())
                .includeSchemaChanges(true) // 包含 schema 变更
                .deserializer(new JsonDebeziumDeserializationSchema())
                .build();
        
        DataStream<String> stream = env.fromSource(
                source,
                WatermarkStrategy.noWatermarks(),
                "MySQL Source"
        );
        
        // 处理数据
        stream.process(new ProcessFunction<String, String>() {
            private transient ValueState<Long> offsetState;
            
            @Override
            public void open(Configuration parameters) throws Exception {
                ValueStateDescriptor<Long> descriptor = 
                    new ValueStateDescriptor<>("offset", Long.class);
                offsetState = getRuntimeContext().getState(descriptor);
            }
            
            @Override
            public void processElement(String value, Context ctx, Collector<String> out) 
                    throws Exception {
                // 保存 offset 到状态
                JSONObject json = JSON.parseObject(value);
                Long offset = json.getLong("pos");
                if (offset != null) {
                    offsetState.update(offset);
                }
                out.collect(value);
            }
        });
        
        env.execute("Consistent Snapshot Example");
    }
}

Exactly-Once 语义

1. 端到端 Exactly-Once

public class ExactlyOnceExample {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        
        // 配置 Exactly-Once
        env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
        env.enableCheckpointing(5000);
        
        // MySQL CDC Source
        MySqlSource<RowData> mySqlSource = createMySqlSource();
        
        DataStream<RowData> mysqlStream = env.fromSource(
                mySqlSource,
                WatermarkStrategy.noWatermarks(),
                "MySQL Source"
        );
        
        // Kafka Sink with Exactly-Once
        KafkaSink<RowData> kafkaSink = KafkaSink.<RowData>builder()
                .setBootstrapServers("localhost:9092")
                .setRecordSerializer(new RowDataSerializationSchema())
                .setDeliverGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
                .setTransactionalIdPrefix("flink-cdc-")
                .build();
        
        mysqlStream.sinkTo(kafkaSink);
        
        env.execute("Exactly-Once CDC Pipeline");
    }
}

性能优化

并行度优化

1. Source 并行度设置

public class ParallelismOptimization {
    public static void optimizeParallelism(StreamExecutionEnvironment env) {
        // 全局并行度
        env.setParallelism(4);
        
        MySqlSource<String> source = MySqlSource.<String>builder()
                .hostname("localhost")
                .port(3306)
                .databaseList("test")
                .tableList("test.large_table")
                .username("root")
                .password("password")
                // 服务器 ID 范围,支持并行读取
                .serverId("5400-5403") // 4个并行度
                .splitSize(8096) // 合理的分片大小
                .deserializer(new JsonDebeziumDeserializationSchema())
                .build();
        
        // Source 并行度
        DataStream<String> stream = env.fromSource(
                source,
                WatermarkStrategy.noWatermarks(),
                "MySQL Source"
        ).setParallelism(4);
        
        // 下游处理并行度可以不同
        stream.map(new MyMapper()).setParallelism(8)
              .keyBy(new MyKeySelector()).setParallelism(16)
              .sink(new MySink()).setParallelism(4);
    }
}

2. 表分片策略

public class TableSplitStrategy {
    public static class CustomSplitStrategy {
        // 根据主键范围分片
        public List<MySqlSplit> splitTable(TableId tableId, Object[] primaryKeys) {
            List<MySqlSplit> splits = new ArrayList<>();
            
            // 获取主键范围
            Object minKey = getMinPrimaryKey(tableId);
            Object maxKey = getMaxPrimaryKey(tableId);
            
            // 计算分片
            long range = (Long)maxKey - (Long)minKey;
            long splitSize = range / 4; // 4个分片
            
            for (int i = 0; i < 4; i++) {
                long start = (Long)minKey + i * splitSize;
                long end = (i == 3) ? (Long)maxKey : start + splitSize;
                
                MySqlSplit split = new MySqlSplit(
                    tableId,
                    "split-" + i,
                    start,
                    end
                );
                splits.add(split);
            }
            
            return splits;
        }
    }
}

内存优化

1. 缓冲区配置

public class MemoryOptimization {
    public static void configureMemory() {
        Properties debeziumProperties = new Properties();
        
        // Kafka 缓冲区设置
        debeziumProperties.put("max.queue.size", "16384");
        debeziumProperties.put("max.batch.size", "8192");
        
        // 数据库连接池
        Properties jdbcProperties = new Properties();
        jdbcProperties.put("connectionPoolSize", "5");
        jdbcProperties.put("connectTimeout", "30");
        
        MySqlSource<String> source = MySqlSource.<String>builder()
                .hostname("localhost")
                .port(3306)
                .databaseList("test")
                .tableList("test.*")
                .username("root")
                .password("password")
                .fetchSize(1000) // JDBC fetch size
                .connectionPoolSize(5) // 连接池大小
                .jdbcProperties(jdbcProperties)
                .debeziumProperties(debeziumProperties)
                .deserializer(new JsonDebeziumDeserializationSchema())
                .build();
    }
}

2. 状态管理优化

public class StateOptimization {
    public static void optimizeState(StreamExecutionEnvironment env) {
        // 使用 RocksDB 状态后端
        env.setStateBackend(new RocksDBStateBackend("hdfs://checkpoints"));
        
        // 增量 Checkpoint
        env.getCheckpointConfig().enableIncrementalCheckpoints();
        
        // 状态 TTL
        StateTtlConfig ttlConfig = StateTtlConfig
                .newBuilder(Time.hours(24))
                .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
                .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
                .build();
        
        // 应用到状态描述符
        ValueStateDescriptor<String> descriptor = 
            new ValueStateDescriptor<>("offset-state", String.class);
        descriptor.enableTimeToLive(ttlConfig);
    }
}

网络优化

1. 批量处理

public class BatchProcessingOptimization {
    public static class BatchingProcessFunction 
            extends ProcessFunction<String, List<String>> {
        
        private transient ListState<String> bufferedElements;
        private final int batchSize = 1000;
        private final long batchInterval = 5000L; // 5秒
        
        @Override
        public void open(Configuration parameters) {
            ListStateDescriptor<String> descriptor = 
                new ListStateDescriptor<>("buffer", String.class);
            bufferedElements = getRuntimeContext().getListState(descriptor);
        }
        
        @Override
        public void processElement(String value, Context ctx, 
                                 Collector<List<String>> out) throws Exception {
            bufferedElements.add(value);
            
            // 注册定时器
            ctx.timerService().registerProcessingTimeTimer(
                ctx.timerService().currentProcessingTime() + batchInterval
            );
            
            // 检查批次大小
            List<String> batch = new ArrayList<>();
            for (String element : bufferedElements.get()) {
                batch.add(element);
            }
            
            if (batch.size() >= batchSize) {
                out.collect(batch);
                bufferedElements.clear();
            }
        }
        
        @Override
        public void onTimer(long timestamp, OnTimerContext ctx, 
                          Collector<List<String>> out) throws Exception {
            List<String> batch = new ArrayList<>();
            for (String element : bufferedElements.get()) {
                batch.add(element);
            }
            
            if (!batch.isEmpty()) {
                out.collect(batch);
                bufferedElements.clear();
            }
        }
    }
}

实战案例

MySQL 到 Elasticsearch

1. 完整示例

public class MySQL2ElasticsearchExample {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
        
        // 设置检查点
        env.enableCheckpointing(5000);
        
        // 创建 MySQL CDC 源表
        tableEnv.executeSql(
            "CREATE TABLE mysql_users (" +
            "  id BIGINT," +
            "  name STRING," +
            "  age INT," +
            "  email STRING," +
            "  address STRING," +
            "  update_time TIMESTAMP(3)," +
            "  PRIMARY KEY (id) NOT ENFORCED" +
            ") WITH (" +
            "  'connector' = 'mysql-cdc'," +
            "  'hostname' = 'localhost'," +
            "  'port' = '3306'," +
            "  'username' = 'root'," +
            "  'password' = 'password'," +
            "  'database-name' = 'test'," +
            "  'table-name' = 'users'," +
            "  'server-id' = '5400-5404'," +
            "  'scan.startup.mode' = 'initial'" +
            ")"
        );
        
        // 创建 Elasticsearch 目标表
        tableEnv.executeSql(
            "CREATE TABLE es_users (" +
            "  id BIGINT," +
            "  name STRING," +
            "  age INT," +
            "  email STRING," +
            "  address STRING," +
            "  update_time TIMESTAMP(3)," +
            "  PRIMARY KEY (id) NOT ENFORCED" +
            ") WITH (" +
            "  'connector' = 'elasticsearch-7'," +
            "  'hosts' = 'http://localhost:9200'," +
            "  'index' = 'users'," +
            "  'document-id.key-delimiter' = '_'," +
            "  'sink.bulk-flush.max-size' = '1000'," +
            "  'sink.bulk-flush.max-actions' = '1000'," +
            "  'sink.bulk-flush.interval' = '1s'," +
            "  'format' = 'json'" +
            ")"
        );
        
        // 数据同步与转换
        tableEnv.executeSql(
            "INSERT INTO es_users " +
            "SELECT " +
            "  id," +
            "  UPPER(name) as name," +  // 转换为大写
            "  age," +
            "  email," +
            "  address," +
            "  update_time " +
            "FROM mysql_users " +
            "WHERE age >= 18"  // 只同步成年用户
        );
    }
}

2. 自定义处理逻辑

public class CustomProcessingExample {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        
        MySqlSource<String> mySqlSource = MySqlSource.<String>builder()
                .hostname("localhost")
                .port(3306)
                .databaseList("test")
                .tableList("test.orders")
                .username("root")
                .password("password")
                .deserializer(new JsonDebeziumDeserializationSchema())
                .build();
        
        DataStream<String> mysqlStream = env.fromSource(
                mySqlSource,
                WatermarkStrategy.noWatermarks(),
                "MySQL Source"
        );
        
        // 解析并处理数据
        DataStream<Order> orderStream = mysqlStream
                .map(json -> JSON.parseObject(json, Order.class))
                .keyBy(Order::getUserId)
                .window(TumblingProcessingTimeWindows.of(Time.minutes(5)))
                .aggregate(new OrderAggregator());
        
        // 写入 Elasticsearch
        orderStream.addSink(
                new ElasticsearchSink.Builder<>(
                        httpHosts,
                        new ElasticsearchSinkFunction<Order>() {
                            @Override
                            public void process(Order order, RuntimeContext ctx, 
                                              RequestIndexer indexer) {
                                IndexRequest request = Requests.indexRequest()
                                        .index("order_stats")
                                        .id(order.getUserId())
                                        .source(JSON.toJSONString(order));
                                indexer.add(request);
                            }
                        }
                ).build()
        );
        
        env.execute("Custom MySQL to Elasticsearch Pipeline");
    }
}

异构数据库同步

1. MySQL 到 PostgreSQL

public class MySQL2PostgreSQLSync {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
        
        // MySQL 源
        tableEnv.executeSql(
            "CREATE TABLE mysql_source (" +
            "  id INT," +
            "  name STRING," +
            "  create_time TIMESTAMP(3)," +
            "  PRIMARY KEY (id) NOT ENFORCED" +
            ") WITH (" +
            "  'connector' = 'mysql-cdc'," +
            "  'hostname' = 'mysql-host'," +
            "  'port' = '3306'," +
            "  'username' = 'root'," +
            "  'password' = 'password'," +
            "  'database-name' = 'source_db'," +
            "  'table-name' = 'source_table'" +
            ")"
        );
        
        // PostgreSQL 目标
        tableEnv.executeSql(
            "CREATE TABLE pg_sink (" +
            "  id INT," +
            "  name STRING," +
            "  create_time TIMESTAMP(3)," +
            "  sync_time TIMESTAMP(3)," +  // 添加同步时间
            "  PRIMARY KEY (id) NOT ENFORCED" +
            ") WITH (" +
            "  'connector' = 'jdbc'," +
            "  'url' = 'jdbc:postgresql://pg-host:5432/target_db'," +
            "  'table-name' = 'target_table'," +
            "  'username' = 'postgres'," +
            "  'password' = 'password'," +
            "  'sink.buffer-flush.max-rows' = '1000'," +
            "  'sink.buffer-flush.interval' = '1s'" +
            ")"
        );
        
        // 同步数据
        tableEnv.executeSql(
            "INSERT INTO pg_sink " +
            "SELECT " +
            "  id," +
            "  name," +
            "  create_time," +
            "  CURRENT_TIMESTAMP as sync_time " +
            "FROM mysql_source"
        );
    }
}

2. 多源合并同步

public class MultiSourceMergeExample {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        
        // MySQL 源
        DataStream<UserEvent> mysqlStream = env.fromSource(
                createMySqlSource(),
                WatermarkStrategy.noWatermarks(),
                "MySQL Source"
        ).map(new MySqlEventParser());
        
        // PostgreSQL 源
        DataStream<UserEvent> pgStream = env.fromSource(
                createPostgreSQLSource(),
                WatermarkStrategy.noWatermarks(),
                "PostgreSQL Source"
        ).map(new PostgreSQLEventParser());
        
        // MongoDB 源
        DataStream<UserEvent> mongoStream = env.fromSource(
                createMongoDBSource(),
                WatermarkStrategy.noWatermarks(),
                "MongoDB Source"
        ).map(new MongoEventParser());
        
        // 合并多个数据源
        DataStream<UserEvent> mergedStream = mysqlStream
                .union(pgStream)
                .union(mongoStream)
                .keyBy(UserEvent::getUserId)
                .process(new DeduplicationFunction()); // 去重
        
        // 统一写入目标系统
        mergedStream.addSink(createUnifiedSink());
        
        env.execute("Multi-Source Merge Sync");
    }
}

实时数据湖构建

1. CDC 到 Iceberg

public class CDC2IcebergExample {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
        
        // 配置 Iceberg Catalog
        tableEnv.executeSql(
            "CREATE CATALOG iceberg_catalog WITH (" +
            "  'type' = 'iceberg'," +
            "  'catalog-type' = 'hive'," +
            "  'uri' = 'thrift://localhost:9083'," +
            "  'warehouse' = 'hdfs://localhost:9000/user/hive/warehouse'" +
            ")"
        );
        
        tableEnv.useCatalog("iceberg_catalog");
        
        // 创建 MySQL CDC 源
        tableEnv.executeSql(
            "CREATE TABLE mysql_orders (" +
            "  order_id BIGINT," +
            "  user_id BIGINT," +
            "  product_id BIGINT," +
            "  amount DECIMAL(10, 2)," +
            "  order_time TIMESTAMP(3)," +
            "  status STRING," +
            "  PRIMARY KEY (order_id) NOT ENFORCED" +
            ") WITH (" +
            "  'connector' = 'mysql-cdc'," +
            "  'hostname' = 'localhost'," +
            "  'port' = '3306'," +
            "  'username' = 'root'," +
            "  'password' = 'password'," +
            "  'database-name' = 'shop'," +
            "  'table-name' = 'orders'" +
            ")"
        );
        
        // 创建 Iceberg 表
        tableEnv.executeSql(
            "CREATE TABLE IF NOT EXISTS iceberg_orders (" +
            "  order_id BIGINT," +
            "  user_id BIGINT," +
            "  product_id BIGINT," +
            "  amount DECIMAL(10, 2)," +
            "  order_time TIMESTAMP(3)," +
            "  status STRING," +
            "  dt STRING," +  // 分区字段
            "  PRIMARY KEY (order_id) NOT ENFORCED" +
            ") PARTITIONED BY (dt) WITH (" +
            "  'format-version' = '2'," +
            "  'write.metadata.delete-after-commit.enabled' = 'true'," +
            "  'write.metadata.previous-versions-max' = '3'" +
            ")"
        );
        
        // 实时同步到 Iceberg
        tableEnv.executeSql(
            "INSERT INTO iceberg_orders " +
            "SELECT " +
            "  order_id," +
            "  user_id," +
            "  product_id," +
            "  amount," +
            "  order_time," +
            "  status," +
            "  DATE_FORMAT(order_time, 'yyyy-MM-dd') as dt " +
            "FROM mysql_orders"
        );
    }
}

2. 实时数仓分层

public class RealtimeDataWarehouse {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);
        
        // ODS 层 - 原始数据
        createODSLayer(tableEnv);
        
        // DWD 层 - 明细数据
        createDWDLayer(tableEnv);
        
        // DWS 层 - 汇总数据
        createDWSLayer(tableEnv);
        
        // ADS 层 - 应用数据
        createADSLayer(tableEnv);
        
        env.execute("Realtime Data Warehouse");
    }
    
    private static void createODSLayer(StreamTableEnvironment tableEnv) {
        // 订单表 CDC
        tableEnv.executeSql(
            "CREATE TABLE ods_orders (" +
            "  order_id BIGINT," +
            "  user_id BIGINT," +
            "  amount DECIMAL(10, 2)," +
            "  order_time TIMESTAMP(3)," +
            "  PRIMARY KEY (order_id) NOT ENFORCED" +
            ") WITH (" +
            "  'connector' = 'mysql-cdc'," +
            "  'hostname' = 'localhost'," +
            "  'port' = '3306'," +
            "  'username' = 'root'," +
            "  'password' = 'password'," +
            "  'database-name' = 'shop'," +
            "  'table-name' = 'orders'" +
            ")"
        );
    }
    
    private static void createDWDLayer(StreamTableEnvironment tableEnv) {
        // 订单明细宽表
        tableEnv.executeSql(
            "CREATE VIEW dwd_order_detail AS " +
            "SELECT " +
            "  o.order_id," +
            "  o.user_id," +
            "  u.user_name," +
            "  u.user_level," +
            "  p.product_name," +
            "  p.category," +
            "  o.amount," +
            "  o.order_time " +
            "FROM ods_orders o " +
            "JOIN ods_users u ON o.user_id = u.user_id " +
            "JOIN ods_products p ON o.product_id = p.product_id"
        );
    }
}

生产部署

集群部署配置

# flink-conf.yaml
jobmanager:
  memory:
    process.size: 2g
    
taskmanager:
  memory:
    process.size: 4g
  numberOfTaskSlots: 4
  
# CDC 相关配置
execution.checkpointing.interval: 10s
execution.checkpointing.mode: EXACTLY_ONCE
state.backend: rocksdb
state.checkpoints.dir: hdfs://namenode:9000/flink-checkpoints

2. Docker Compose 部署

version: '3.8'
services:
  jobmanager:
    image: flink:1.17-scala_2.12
    ports:
      - "8081:8081"
    command: jobmanager
    environment:
      - JOB_MANAGER_RPC_ADDRESS=jobmanager
    volumes:
      - ./conf:/opt/flink/conf
      - ./jars:/opt/flink/lib
      
  taskmanager:
    image: flink:1.17-scala_2.12
    depends_on:
      - jobmanager
    command: taskmanager
    scale: 3
    environment:
      - JOB_MANAGER_RPC_ADDRESS=jobmanager
    volumes:
      - ./conf:/opt/flink/conf
      - ./jars:/opt/flink/lib
      
  mysql:
    image: mysql:8.0
    environment:
      MYSQL_ROOT_PASSWORD: password
      MYSQL_DATABASE: test
    ports:
      - "3306:3306"
    command: --log-bin=mysql-bin --binlog-format=ROW

3. Kubernetes 部署

apiVersion: flink.apache.org/v1beta1
kind: FlinkDeployment
metadata:
  name: flink-cdc-cluster
spec:
  image: flink:1.17-scala_2.12
  flinkVersion: v1_17
  flinkConfiguration:
    taskmanager.numberOfTaskSlots: "4"
    state.backend: rocksdb
    state.checkpoints.dir: s3://flink-checkpoints
    execution.checkpointing.interval: 10s
    
  serviceAccount: flink
  
  jobManager:
    resource:
      memory: "2048m"
      cpu: 1
      
  taskManager:
    replicas: 3
    resource:
      memory: "4096m"
      cpu: 2
      
  job:
    jarURI: local:///opt/flink/lib/flink-cdc-job.jar
    parallelism: 4
    upgradeMode: savepoint
    savepointTriggerNonce: 0

监控和运维

监控指标

1. 自定义监控指标

public class CDCMetricsExample {
    public static class MetricsProcessFunction 
            extends ProcessFunction<String, String> {
        
        // 延迟监控
        private transient Histogram eventLatency;
        // 吞吐量监控
        private transient Meter throughput;
        // 错误计数
        private transient Counter errorCounter;
        // 表级别监控
        private transient Map<String, Counter> tableCounters;
        
        @Override
        public void open(Configuration parameters) {
            MetricGroup metricGroup = getRuntimeContext().getMetricGroup();
            
            eventLatency = metricGroup.histogram("event_latency", 
                new DescriptiveStatisticsHistogram(1000));
            throughput = metricGroup.meter("throughput", new MeterView(60));
            errorCounter = metricGroup.counter("errors");
            
            tableCounters = new HashMap<>();
        }
        
        @Override
        public void processElement(String value, Context ctx, 
                                 Collector<String> out) throws Exception {
            try {
                JSONObject event = JSON.parseObject(value);
                
                // 计算延迟
                long eventTime = event.getLong("ts_ms");
                long currentTime = System.currentTimeMillis();
                long latency = currentTime - eventTime;
                eventLatency.update(latency);
                
                // 记录吞吐量
                throughput.markEvent();
                
                // 表级别统计
                String table = event.getString("table");
                tableCounters.computeIfAbsent(table, k -> 
                    getRuntimeContext().getMetricGroup()
                        .counter("table_" + k + "_count")
                ).inc();
                
                out.collect(value);
                
            } catch (Exception e) {
                errorCounter.inc();
                LOG.error("Process error", e);
            }
        }
    }
}

2. Prometheus 集成

# flink-conf.yaml
metrics.reporters: prom
metrics.reporter.prom.class: org.apache.flink.metrics.prometheus.PrometheusReporter
metrics.reporter.prom.port: 9249

# Prometheus 配置
scrape_configs:
  - job_name: 'flink'
    static_configs:
      - targets: ['jobmanager:9249', 'taskmanager1:9249', 'taskmanager2:9249']

告警配置

public class AlertingExample {
    public static class LagMonitor extends ProcessFunction<String, Alert> {
        private final long lagThreshold = 60000; // 1分钟
        
        @Override
        public void processElement(String value, Context ctx, 
                                 Collector<Alert> out) throws Exception {
            JSONObject event = JSON.parseObject(value);
            
            long eventTime = event.getLong("ts_ms");
            long currentTime = System.currentTimeMillis();
            long lag = currentTime - eventTime;
            
            if (lag > lagThreshold) {
                Alert alert = new Alert(
                    "HIGH_LAG",
                    "CDC lag exceeds threshold: " + lag + "ms",
                    Alert.Severity.WARNING
                );
                out.collect(alert);
                
                // 发送告警
                sendAlert(alert);
            }
        }
        
        private void sendAlert(Alert alert) {
            // 发送到告警系统
            // 例如:钉钉、邮件、短信等
        }
    }
}

常见问题

Q1: 如何处理大表的初始同步?

解答:

  1. 增加并行度:使用多个 server-id 增加并行读取
  2. 调整分片大小:根据表大小设置合适的 splitSize
  3. 限流读取:控制 fetchSize 避免 OOM
  4. 分批同步:按时间范围或 ID 范围分批
// 大表同步优化
MySqlSource<String> source = MySqlSource.<String>builder()
    .hostname("localhost")
    .port(3306)
    .databaseList("large_db")
    .tableList("large_db.huge_table")
    .username("root")
    .password("password")
    .serverId("5400-5415") // 16个并行度
    .splitSize(4096) // 减小分片大小
    .fetchSize(500) // 限制每次拉取量
    .connectTimeout(Duration.ofMinutes(5))
    .deserializer(new JsonDebeziumDeserializationSchema())
    .build();

Q2: 如何保证数据一致性?

解答:

  1. 启用 Checkpoint:确保故障恢复
  2. 使用 Exactly-Once 语义:端到端一致性
  3. 处理重复数据:使用主键去重
  4. 监控延迟:确保及时同步

Q3: Schema 变更如何处理?

解答:

  1. 自动感知:启用 include.schema.changes
  2. 兼容性处理:使用 JSON 等灵活格式
  3. 版本管理:维护 Schema 版本历史
  4. 平滑升级:使用蓝绿部署

Q4: 性能优化建议?

解答:

  1. 合理设置并行度:根据源库负载调整
  2. 批量处理:减少网络开销
  3. 内存优化:调整缓冲区大小
  4. 监控指标:及时发现瓶颈

Q5: 如何处理删除操作?

解答:

// 处理删除事件
public class DeleteHandler implements DebeziumDeserializationSchema<String> {
    @Override
    public void deserialize(SourceRecord record, Collector<String> out) {
        Envelope.Operation op = Envelope.operationFor(record);
        
        if (op == Envelope.Operation.DELETE) {
            // 处理删除
            Struct before = value.getStruct("before");
            String id = before.getString("id");
            
            // 方案1:标记删除
            JSONObject deleteEvent = new JSONObject();
            deleteEvent.put("id", id);
            deleteEvent.put("_deleted", true);
            deleteEvent.put("_delete_time", System.currentTimeMillis());
            out.collect(deleteEvent.toJSONString());
            
            // 方案2:发送到删除主题
            // sendToDeleteTopic(id);
        }
    }
}

相关文章

高级特性

集成方案

相关技术

  • MySQL - MySQL 数据库
  • PostgreSQL - PostgreSQL 数据库
  • Kafka - 消息中间件
  • Elasticsearch - 搜索引擎

总结

Flink CDC 提供了强大的实时数据同步能力,是构建实时数据管道的理想选择:

🎯 核心优势

  1. 实时性高:毫秒级延迟的数据同步
  2. 一致性强:Exactly-Once 语义保证
  3. 扩展性好:支持大规模分布式部署
  4. 使用简单:SQL 和 DataStream API 双支持

✅ 最佳实践

  • 合理配置数据库参数(binlog、wal 等)
  • 根据数据量选择合适的并行度
  • 使用 Checkpoint 保证容错性
  • 建立完善的监控告警机制
  • 定期进行性能调优

🚀 适用场景

  • 实时数据同步和备份
  • 实时数仓和数据湖构建
  • 缓存同步和搜索索引更新
  • 微服务间数据共享
  • 异构数据库迁移

🚫 注意事项:

  • 确保源数据库开启必要的日志功能
  • 监控源库压力,避免影响业务
  • 合理设置资源,防止 OOM
  • 做好 Schema 变更的兼容处理

通过 Flink CDC,可以轻松构建企业级的实时数据同步平台,为数据驱动的业务决策提供有力支撑。

上一章 / 下一章