在业务系统迭代过程中,经常需要将MySQL数据库的变更数据实时同步到其他存储系统或者业务模块中,比如同步到Elasticsearch做搜索索引、同步到Redis做缓存、同步到另一套数据库做数据备份。传统的同步方式多采用定时轮询数据库的方式,这种方式不仅会增加数据库的压力,还会存在秒级甚至分钟级的同步延迟,无法满足对实时性要求高的场景。Canal作为一款基于MySQL Binlog增量日志解析的工具,能够实时捕获数据库的增删改操作,结合PHP的异步处理能力,可以高效实现数据同步的实时性要求。

Canal监听Binlog的基本原理
MySQL的Binlog是数据库记录所有数据变更操作的二进制日志,包含了数据的增删改详细记录。Canal模拟MySQL Slave的交互协议,伪装成MySQL的从库,向主库发送dump协议请求,主库收到请求后会推送Binlog日志给Canal,Canal解析这些日志后,将变更数据以结构化形式输出,供下游消费者处理。
整个同步流程的核心步骤如下:
- MySQL开启Binlog并设置为ROW模式,确保能记录每一行数据的变更细节
- Canal服务端部署并配置MySQL连接信息,启动后连接MySQL主库
- Canal解析接收到的Binlog日志,封装成变更事件
- PHP客户端连接Canal服务端,订阅变更事件并处理同步逻辑
环境准备与配置
1. MySQL配置
首先需要确保MySQL开启了Binlog,并且格式为ROW,修改MySQL配置文件my.cnf,添加以下配置:
[mysqld] # 开启Binlog log-bin=mysql-bin # Binlog格式设置为ROW binlog_format=ROW # 设置Binlog的过期时间,避免日志占用过多磁盘 expire_logs_days=7 # 给Canal使用的账号开启复制权限 server-id=1
配置完成后重启MySQL,然后创建Canal使用的账号并授权:
-- 创建canal用户 CREATE USER 'canal'@'%' IDENTIFIED BY 'canal_password'; -- 授予复制权限 GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'canal'@'%'; -- 刷新权限 FLUSH PRIVILEGES;
2. Canal服务端部署
下载Canal最新稳定版本,解压后修改conf/example/instance.properties配置文件,配置MySQL连接信息:
# MySQL地址 canal.instance.master.address=127.0.0.1:3306 # MySQL账号 canal.instance.dbUsername=canal # MySQL密码 canal.instance.dbPassword=canal_password # 需要监听的数据库,不填则监听所有 canal.instance.filter.regex=.*\..*
启动Canal服务端,默认会监听11111端口,等待客户端连接。
PHP客户端实现数据同步
PHP需要连接Canal服务端,订阅Binlog变更事件,解析事件内容后执行对应的同步操作。这里我们使用PHP的socket相关函数实现Canal客户端的连接和通信。
1. Canal客户端连接与订阅
Canal使用自定义的协议进行通信,PHP需要先发送握手请求,然后发送订阅请求,之后循环接收变更数据。以下是基础的连接和订阅代码:
<?php
class CanalClient {
private $host;
private $port;
private $socket;
private $clientId = 1001; // 客户端ID,自定义唯一值
public function __construct($host = '127.0.0.1', $port = 11111) {
$this->host = $host;
$this->port = $port;
}
// 连接Canal服务端
public function connect() {
$this->socket = socket_create(AF_INET, SOCK_STREAM, SOL_TCP);
if (!$this->socket) {
throw new Exception('创建socket失败:' . socket_strerror(socket_last_error()));
}
$result = socket_connect($this->socket, $this->host, $this->port);
if (!$result) {
throw new Exception('连接Canal服务端失败:' . socket_strerror(socket_last_error()));
}
// 发送握手请求
$this->sendHandshake();
// 发送客户端认证
$this->sendAuth();
// 订阅所有数据库变更
$this->subscribe();
}
// 发送握手请求
private function sendHandshake() {
// Canal握手协议包,具体内容参考Canal官方协议文档
$packet = pack('C', 0x01); // 协议版本
socket_write($this->socket, $packet, strlen($packet));
// 读取服务端响应
$response = socket_read($this->socket, 1024);
if (empty($response)) {
throw new Exception('握手失败');
}
}
// 发送认证信息
private function sendAuth() {
// 认证包格式:clientId + 认证信息,这里简化为发送clientId
$packet = pack('N', $this->clientId);
socket_write($this->socket, $packet, strlen($packet));
$response = socket_read($this->socket, 1024);
if (empty($response)) {
throw new Exception('认证失败');
}
}
// 发送订阅请求
private function subscribe() {
// 订阅请求包,指定订阅的filter规则,这里订阅所有
$filter = '';
$packet = pack('N', strlen($filter)) . $filter;
socket_write($this->socket, $packet, strlen($packet));
$response = socket_read($this->socket, 1024);
if (empty($response)) {
throw new Exception('订阅失败');
}
}
// 获取变更数据
public function getChanges() {
// 发送获取数据的请求
$packet = pack('N', 1); // 请求一次获取一批数据
socket_write($this->socket, $packet, strlen($packet));
// 读取返回的数据包
$lengthData = socket_read($this->socket, 4);
if (strlen($lengthData) != 4) {
return null;
}
$length = unpack('N', $lengthData)[1];
$data = socket_read($this->socket, $length);
return $this->parseChangeData($data);
}
// 解析变更数据,这里简化为解析基础结构,实际需要根据Canal协议完整解析
private function parseChangeData($data) {
// 实际解析逻辑需要处理Canal返回的RowChange对象,包含数据库、表、操作类型、变更前后数据
// 这里返回模拟的解析结果
return [
'database' => 'test_db',
'table' => 'user',
'type' => 'INSERT',
'data' => [
'id' => 1,
'name' => '张三',
'age' => 20
]
];
}
// 关闭连接
public function close() {
if ($this->socket) {
socket_close($this->socket);
}
}
}
// 使用示例
$client = new CanalClient();
try {
$client->connect();
while (true) {
$changes = $client->getChanges();
if (!empty($changes)) {
// 处理同步逻辑
handleSync($changes);
}
sleep(1); // 适当休眠,避免频繁请求
}
} catch (Exception $e) {
echo '同步异常:' . $e->getMessage() . PHP_EOL;
} finally {
$client->close();
}
// 处理同步逻辑的函数
function handleSync($change) {
$db = $change['database'];
$table = $change['table'];
$type = $change['type'];
$data = $change['data'];
// 根据操作类型执行不同的同步操作
switch ($type) {
case 'INSERT':
// 执行插入同步逻辑,比如同步到Redis
echo "新增数据:{$db}.{$table}," . json_encode($data) . PHP_EOL;
break;
case 'UPDATE':
// 执行更新同步逻辑
echo "更新数据:{$db}.{$table}," . json_encode($data) . PHP_EOL;
break;
case 'DELETE':
// 执行删除同步逻辑
echo "删除数据:{$db}.{$table}," . json_encode($data) . PHP_EOL;
break;
}
}
?>
2. 优化建议
上述代码是基础的演示实现,实际生产环境中需要做以下优化:
- 使用PHP的常驻进程运行客户端,比如用Swoole扩展实现异步IO,避免频繁创建销毁连接
- 增加错误重试机制,当连接断开时自动重连,保证同步服务的稳定性
- 对解析后的变更数据做批量处理,减少下游系统的请求次数,提升同步效率
- 增加同步日志和监控,方便排查同步过程中出现的问题
常见问题与解决方案
在实际落地过程中,可能会遇到以下问题:
- Canal无法连接MySQL:检查MySQL的Binlog配置是否正确,Canal使用的账号权限是否足够,网络是否通畅
- 同步延迟过高:检查Canal服务端的性能,是否有大量堆积的Binlog未处理,PHP客户端的消费速度是否跟得上
- 数据解析错误:确认MySQL的Binlog格式是否为ROW,Canal的版本是否和MySQL版本兼容
通过Canal监听Binlog配合PHP处理,能够很好地实现数据同步的实时性要求,相比传统的定时同步方式,延迟可以控制在毫秒级,同时不会对业务数据库造成额外的查询压力,适合大多数需要实时数据同步的业务场景。