Flink CDC实战:5分钟搞定MySQL数据实时同步到Elasticsearch(Flink 1.18 + Docker Compose一键部署)
Flink CDC极速体验:MySQL到Elasticsearch的实时数据管道搭建
在数据驱动的时代,企业对于实时数据同步的需求日益增长。想象一下,当MySQL数据库中的客户信息发生变更时,搜索服务能立即反映这些变化,这将极大提升用户体验。Flink CDC正是为此而生的利器,它让数据实时同步变得前所未有的简单。
本文将带你用最短的时间搭建一个完整的实时数据同步系统。我们采用Docker Compose一键部署所有组件,通过Flink SQL实现零代码开发,整个过程只需5-10分钟。这种开箱即用的体验特别适合快速验证原型、学习测试或小型项目部署。
1. 环境准备与一键部署
1.1 Docker Compose编排
我们使用Docker Compose来管理所有服务依赖,包括MySQL、Elasticsearch、Kibana和Flink集群。这种容器化部署方式保证了环境隔离,也避免了复杂的配置过程。
创建一个docker-compose.yml文件,内容如下:
version: '2.1'
services:
mysql:
image: debezium/example-mysql:1.1
ports:
- "3306:3306"
environment:
- MYSQL_ROOT_PASSWORD=123456
- MYSQL_USER=root
- MYSQL_PASSWORD=123456
elasticsearch:
image: elastic/elasticsearch:7.6.0
environment:
- cluster.name=docker-cluster
- bootstrap.memory_lock=true
- "ES_JAVA_OPTS=-Xms512m -Xmx512m"
- discovery.type=single-node
ports:
- "9200:9200"
- "9300:9300"
kibana:
image: elastic/kibana:7.6.0
ports:
- "5601:5601"
depends_on:
- elasticsearch
启动所有服务只需一条命令:
docker-compose up -d
验证服务是否正常运行:
docker ps
你应该看到三个容器正在运行:MySQL、Elasticsearch和Kibana。
1.2 Flink集群准备
虽然我们可以将Flink也加入Docker Compose,但为了更贴近生产环境,这里选择独立部署Flink 1.18版本。下载并解压Flink后,需要获取两个关键依赖包:
flink-cdc-pipeline-connector-mysql-3.0.0.jarflink-sql-connector-elasticsearch7-3.0.1-1.17.jar
将这两个JAR文件放入Flink的lib目录后,启动集群:
./bin/start-cluster.sh
2. 数据准备与管道配置
2.1 MySQL测试数据
在MySQL中创建一个测试数据库和表,并插入一些初始数据:
CREATE DATABASE cdctest;
USE cdctest;
CREATE TABLE user (
id INT PRIMARY KEY,
name VARCHAR(255),
age INT
);
INSERT INTO user VALUES
(1, '张三', 28),
(2, '李四', 32),
(3, '王五', 25);
2.2 Flink SQL配置
启动Flink SQL客户端:
./bin/sql-client.sh
设置必要的参数:
SET sql-client.execution.result-mode = tableau;
SET execution.checkpointing.interval = 3s;
3. 实时同步管道搭建
3.1 创建CDC源表
定义MySQL的CDC源表,Flink将从这个表捕获变更:
CREATE TABLE mysql_user (
id INT,
name STRING,
age INT,
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'mysql-cdc',
'hostname' = 'localhost',
'port' = '3306',
'username' = 'root',
'password' = '123456',
'database-name' = 'cdctest',
'table-name' = 'user'
);
3.2 创建Elasticsearch目标表
定义Elasticsearch目标表结构:
CREATE TABLE es_user (
id INT,
name STRING,
age INT,
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'elasticsearch-7',
'hosts' = 'http://elasticsearch:9200',
'index' = 'user_index'
);
3.3 启动同步任务
通过一个简单的INSERT语句启动同步:
INSERT INTO es_user SELECT * FROM mysql_user;
这个任务会持续运行,自动将MySQL中的变更同步到Elasticsearch。
4. 验证与监控
4.1 Kibana数据查看
访问Kibana界面(http://localhost:5601),通过Dev Tools查询数据:
GET /user_index/_search
{
"query": {
"match_all": {}
}
}
你应该能看到从MySQL同步过来的用户数据。
4.2 实时变更测试
在MySQL中执行以下操作,观察Elasticsearch中的数据变化:
-
插入新记录:
INSERT INTO user VALUES (4, '赵六', 40); -
更新现有记录:
UPDATE user SET age = 30 WHERE id = 1; -
删除记录:
DELETE FROM user WHERE id = 2;
每次操作后,在Kibana中刷新查询,都能立即看到变更后的结果。
4.3 Flink作业监控
访问Flink Web UI(默认http://localhost:8081),可以看到正在运行的CDC作业。这里可以监控作业状态、吞吐量和延迟等关键指标。
5. 进阶配置与优化
5.1 初始快照控制
对于大表,初始快照可能耗时较长。可以通过以下参数控制:
WITH (
...
'scan.incremental.snapshot.enabled' = 'true',
'scan.incremental.snapshot.chunk.size' = '8096',
'scan.snapshot.fetch.size' = '1024'
)
5.2 故障恢复配置
确保作业具备容错能力:
SET execution.checkpointing.interval = 10s;
SET execution.checkpointing.tolerable-failed-checkpoints = 3;
SET restart-strategy.fixed-delay.attempts = 5;
5.3 字段映射与转换
可以在同步过程中进行简单的数据转换:
INSERT INTO es_user
SELECT
id,
UPPER(name) AS name,
age + 1 AS age
FROM mysql_user;
5.4 多表同步策略
对于需要同步整个数据库的场景,可以使用正则表达式匹配多个表:
CREATE TABLE mysql_source (
...
) WITH (
...
'database-name' = 'cdctest',
'table-name' = 'cdctest\.user, cdctest\.product.*'
);
6. 生产环境考量
6.1 性能调优参数
# MySQL CDC连接器
'server-id' = '5400-5404' # 确保唯一
'connect.timeout' = '30s'
'connect.max-retries' = '3'
# Elasticsearch连接器
'sink.bulk-flush.max-actions' = '100'
'sink.bulk-flush.interval' = '1s'
6.2 安全配置
# MySQL SSL连接
'use-ssl' = 'true'
'ssl-mode' = 'REQUIRED'
# Elasticsearch认证
'username' = 'elastic'
'password' = 'yourpassword'
6.3 监控与告警
建议配置以下监控指标:
- 源数据库的binlog延迟
- Flink检查点持续时间
- Elasticsearch批量写入延迟
- 网络吞吐量
7. 常见问题排查
7.1 连接问题检查清单
| 问题现象 | 可能原因 | 解决方案 |
|---|---|---|
| 无法连接MySQL | 网络不通/权限不足 | 检查防火墙、验证凭证 |
| 读取不到binlog | server-id冲突 | 更换server-id范围 |
| ES写入失败 | 索引不存在/字段类型冲突 | 预先创建索引或启用自动创建 |
7.2 性能问题优化
-- 增加并行度
SET parallelism.default = 4;
-- 调整批处理参数
'sink.bulk-flush.max-actions' = '500'
'sink.bulk-flush.interval' = '2s'
7.3 数据一致性验证
定期运行校验查询,比较源和目标的数据差异:
-- MySQL计数
SELECT COUNT(*) FROM user;
-- ES计数
GET /user_index/_count
8. 环境清理
完成测试后,可以优雅地停止所有服务:
# 停止Flink作业
./bin/flink list
./bin/flink cancel <job-id>
# 停止Docker容器
docker-compose down
这套环境特别适合作为学习平台或原型验证。在实际生产部署时,需要考虑高可用配置、更完善的安全策略和监控体系。
更多推荐
所有评论(0)