10分钟搞定Flink CDC:MySQL到HBase实时同步流程详解
·
Flink CDC:MySQL到HBase实时同步(10分钟速通)
1. 核心原理
- Flink CDC 捕获MySQL的binlog变更事件(增删改),通过流式计算引擎实时处理。
- 数据流向:MySQL → Flink CDC Connector → Flink SQL/Table API → HBase Connector → HBase表。
- 关键公式:同步延迟 $T_{\text{delay}} = T_{\text{extract}} + T_{\text{process}} + T_{\text{load}}$,通常控制在毫秒级。
2. 环境准备
# 依赖组件
Flink 1.13+ | MySQL 5.7+ | HBase 2.0+ | Flink CDC Connector 2.3+
3. 四步实现同步
步骤1:MySQL配置
-- 启用binlog(my.cnf)
[mysqld]
server_id=1
log_bin=mysql-bin
binlog_format=ROW
步骤2:Flink SQL作业(核心代码)
-- 创建MySQL CDC源表
CREATE TABLE mysql_users (
id BIGINT PRIMARY KEY,
name STRING,
email STRING
) WITH (
'connector' = 'mysql-cdc',
'hostname' = 'localhost',
'port' = '3306',
'username' = 'user',
'password' = 'pass',
'database-name' = 'test_db',
'table-name' = 'users'
);
-- 创建HBase目标表
CREATE TABLE hbase_users (
rowkey BIGINT,
cf ROW<name STRING, email STRING>,
PRIMARY KEY (rowkey) NOT ENFORCED
) WITH (
'connector' = 'hbase-2.2',
'table-name' = 'users',
'zookeeper.quorum' = 'localhost:2181',
'zookeeper.znode.parent' = '/hbase'
);
-- 启动同步管道
INSERT INTO hbase_users
SELECT id AS rowkey, ROW(name, email) AS cf
FROM mysql_users;
步骤3:HBase表结构设计
- RowKey:直接使用MySQL主键(如
id) - 列族:
cf包含所有字段(name,email)
步骤4:提交Flink作业
./bin/flink run -d -c org.apache.flink.table.api.SqlRunner /path/to/sync_job.sql
4. 验证同步
# 在MySQL插入数据
INSERT INTO users VALUES (1001, '张三', 'zhangsan@example.com');
# 在HBase查询
hbase> scan 'users', {LIMIT => 1}
ROW COLUMN+CELL
1001 column=cf:name, timestamp=... value=张三
5. 常见问题解决
| 问题现象 | 解决方案 |
|---|---|
Can't connect to MySQL | 检查CDC用户权限:GRANT SELECT, REPLICATION SLAVE ON *.* TO 'user' |
| HBase写入超时 | 增大Flink checkpoint间隔:execution.checkpointing.interval: 30s |
| 数据乱码 | 在DDL中指定编码:'sink.properties.hbase.client.encoding' = 'UTF-8' |
耗时统计:环境配置(3分钟)→ SQL编写(2分钟)→ 部署测试(5分钟)
6. 优化建议
- 增量快照:使用
scan.incremental.snapshot.enabled=true避免全表锁 - 并行度:根据QPS调整
parallelism.default - 容错:开启Flink checkpoint+WAL保障Exactly-Once
通过Flink Web UI监控实时同步状态:
注:图中Records Sent和Records Received差值应趋近于0
更多推荐

所有评论(0)