news 2026/8/14 1:31:23

基于Canal实现mysql数据同步到消息队列(RabbitMQ/Kafka)详细操作教程

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
基于Canal实现mysql数据同步到消息队列(RabbitMQ/Kafka)详细操作教程

由于网络网闸问题,目前采用canal进行解析mysql,然后推送到mq,再次进行消费写入数据库。

以下基于arm64 linux ubuntu24、mysql 8.0、 docker部署、canal-server v1.1.8、kafka v3.9.0、rabbitmq v4.3.2

默认已经安装好数据库(mysql 8.0)、mysql已开启binlog、主从同步、创建好对应权限的用户和密码

一、mysql 修改

1. 修改mysql.ini(类似的配置文件)

#MySQL8.0 兼容密码插件(canal客户端识别) #自行搜索是否有,也许可能不需要 default_authentication_plugin=mysql_native_password

2.创建同步使用账号

CREATE USER canal@'%' IDENTIFIED BY 'canal@113.'; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO canal@'%'; -- (可能不需要)MySQL8.0 必须重置认证方式为 native(否则canal连接报错) ALTER USER canal@'%' IDENTIFIED WITH mysql_native_password BY 'canal@113.'; FLUSH PRIVILEGES; --验证是否有权限 SHOW GRANTS FOR canal@'%'; ---结果如下(说明成功了) GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO `canal`@`%`

3.验证是否开启binlog,并且是ROW

---验证 binlog 是否开启 看到是 log_bin= On binlog_format=ROW SHOW VARIABLES LIKE 'log_bin'; SHOW VARIABLES LIKE 'binlog_format';

二、安装所需镜像

1. canal-server

docker pull canal/canal-server:v1.1.8

2. kafka

docker pull apache/kafka:3.9.0

3. rabbitmq

docker pull rabbitmq:4.3.2-management

三、创建所需要的文件夹

本文演示路径是/run/media/nvme0n1/rlzk

1.canal-server

mkdir canal chmod 777 canal cd canal/ mkdir logs mkdir data mkdir conf chmod 777 data chmod 777 logs chmod 777 conf

2.rabbitmq

mkdir rabbitmq chmod 777 rabbitmq

3.kafka

mkdir kafka chmod 777 kafka

四、容器创建及使用

1. rabbitmq

1.1 创建容器(注意其中的账号密码

docker run \ --name rabbitMQ \ --network xx-network \ --add-host=host.docker.internal:host-gateway \ -e TZ=Asia/Shanghai \ -p 7008:5672 \ -p 15672:15672 \ -v /run/media/nvme0n1/rlzk/rabbitmq:/var/lib/rabbitmq \ -e RABBITMQ_DEFAULT_USER=rabbitmq \ -e RABBITMQ_DEFAULT_PASS=rabbitmq@123654 \ --restart always \ -d rabbitmq:4.3.2-management

1.2登录其后台,进行一些配置(或者客户端直接连上去创建),并且可能导致canal-server推送数据 失败

登录地址 http://ip:15672 账号密码上上方代码里

1.2.1找到Exchanges需要添加canal.change topic类型,如下图

1.2.2 绑定queue, 需创建canal_queue,如下图

1.2.3 绑定队列(演示 canal.change1),如下图(效果入图三)

2.canal-server

2.1 canal-server conf文件

2.1.1 创建临时容器
docker run -d --name canal-tmp canal/canal-server:v1.1.8
2.1.2 复制容器内文件到本地目录(注意需要先定位到/run/media/nvme0n1/rlzk/canal/conf
docker cp canal-tmp:/home/admin/canal-server/conf .
2.1.3 移除临时容器
docker rm -f canal-tmp

2.2 修改配置文件(canal.properties

cd /run/media/nvme0n1/rlzk/canal/conf vi canal.properties --找到 添加 admin 和密码 canal.admin.user = admin canal.admin.passwd =admin@123654 --找到 serverMode= 改为rabbitMQ canal.serverMode = rabbitMQ --找到 rabbitMQ配置项,改为如下 ---由于采用docker配置 所以这个host是 docker 创建的名称 ----如果这个 rabbitmq 默认端口不是5672 那么 rabbitmq.host= xxx:端口 rabbitmq.host = rabbitMQ rabbitmq.virtual.host = / rabbitmq.exchange = canal.exchange rabbitmq.username = rabbitmq rabbitmq.password = rabbitmq@123654 rabbitmq.queue = rabbitmq.routingKey = ${database}.${table} rabbitmq.deliveryMode = 2 ---改完记得保存

2.3 修改配置文件(instance.properties)如下图

cd example vi instance.properties ----修改数据库所在地址(采用docker部署,数据库在宿主机所以使用host.docker.internal) canal.instance.master.address=host.docker.internal:3306 --找到设置连接数据库的账号和密码设置,这个账号在一中创建;如下 canal.instance.dbUsername=canal canal.instance.dbPassword=canal@113. --找到过滤(只要所需的数据库) canal.instance.filter.regex=xxx_db\\..* ---如果要从指定位置开始同步 找到其中的binlog 和 pos 进行修改,以下 两个值是假的 # 指定binlog文件名 canal.instance.master.journal.name=binlog.000326 # 指定偏移pos canal.instance.master.position=89100413 # 时间戳留空,二选一用文件+pos更精准 canal.instance.master.timestamp= 修改后进行保存

2.4 创建容器(canal-server)

docker run \ --name canal-server \ --network xx-network \ --add-host=host.docker.internal:host-gateway \ -e TZ=Asia/Shanghai \ -p 11111:11111 \ -p 11110:11110 \ -v /run/media/nvme0n1/rlzk/canal/conf:/home/admin/canal-server/conf \ -v /run/media/nvme0n1/rlzk/canal/logs:/home/admin/canal-server/logs \ --restart always \ -d canal/canal-server:v1.1.8

2.5 验证rabbitmq有没有收到数据了

3.kafka

3.1 创建容器(注意其中的172.0.10.140,这个到时客户端连接会使用到)
docker run \ --name kafka \ -p 9092:9092 \ --network xx-network \ --add-host=host.docker.internal:host-gateway \ -e TZ=Asia/Shanghai \ --restart always \ -e KAFKA_NODE_ID=1 \ -e KAFKA_PROCESS_ROLES=broker,controller \ -e KAFKA_CONTROLLER_QUORUM_VOTERS=1@kafka:9093 \ -e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092,CONTROLLER://kafka:9093 \ -e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://172.0.10.140:9092 \ -e KAFKA_CONTROLLER_LISTENER_NAMES=CONTROLLER \ -e KAFKA_LOG_DIRS=/opt/kraft-data \ -e KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1 \ -e KAFKA_DEFAULT_REPLICATION_FACTOR=1 \ -v /run/media/nvme0n1/rlzk/kafka:/opt/kraft-data \ -d apache/kafka:3.9.0
3.2 更改 canal-server 配置
3.2.1 只要改动其中的 canal.serverMode, 如下图

五、附带C# 客户端代码

5.1 引用常用类库(Newtonsoft.json v13.0.4、RabbitMQ.Client v7.2.1 、Confluent.Kafka v2.15.0

5.2 通用类

using System; using System.Collections.Generic; using System.Linq; using System.Text; using System.Threading.Tasks; namespace CanalSync { public class CanalMsg { public List<Dictionary<string, object>> data { get; set; } /// <summary> /// /// </summary> public string database { get; set; } public long es { get; set; } public string gtid { get; set; } public int id { get; set; } /// <summary> /// 是否ddl /// </summary> public bool? isDdl { get; set; } /// <summary> /// isDdl true 这个才有值 /// </summary> public string ddl { get; set; } public Dictionary<string, string> mysqlType { get; set; } public List<Dictionary<string, object>> old { get; set; } public List<string> pkNames { get; set; } public string sql { get; set; } //public Dictionary<string,string> sqlType { get; set; } /// <summary> /// /// </summary> public string table { get; set; } public long ts { get; set; } public string type { get; set; } } }

5.3 rabbitmq

using Newtonsoft.Json; using Newtonsoft.Json.Linq; using RabbitMQ.Client; using RabbitMQ.Client.Events; using System.Data.Common; using System.Text; using System.Threading.Channels; using CanalSync; partial class Program { private static IConnection _connection; private static IChannel _channel; static async Task Main(string[] args) { var factory = new ConnectionFactory() { HostName = "172.0.10.140", Port = 5672, UserName = "rabbitmq", Password = "rabbitmq@123654", AutomaticRecoveryEnabled = true, NetworkRecoveryInterval = TimeSpan.FromSeconds(3), RequestedHeartbeat = TimeSpan.FromSeconds(30) }; _connection = await factory.CreateConnectionAsync(); _channel = await _connection.CreateChannelAsync(); await _channel.ExchangeDeclareAsync("canal.exchange", "topic", true); await _channel.QueueDeclareAsync("canal_queue", true, exclusive: false, autoDelete: false); await _channel.QueueBindAsync("canal_queue", "canal.exchange", "#"); await _channel.BasicQosAsync(0, 1, false); var consume = new AsyncEventingBasicConsumer(_channel); consume.ReceivedAsync += OnMessageReceived; await _channel.BasicConsumeAsync("canal_queue", autoAck: false, consume); Console.WriteLine("Canal消费者启动成功,等待数据..."); // 常驻阻塞,程序不退出,连接不会被释放 await Task.Delay(-1); Console.WriteLine("程序退出了"); // 程序退出时手动关闭 await _channel.CloseAsync(); await _connection.CloseAsync(); await Task.Delay(TimeSpan.FromSeconds(3)); Console.WriteLine("程序完全退出"); } /// <summary> /// 消息处理回调(修复ack时通道关闭问题) /// </summary> private static async Task OnMessageReceived(object sender, BasicDeliverEventArgs ea) { try { // 先判断通道是否可用,避免关闭后ack抛异常 if (_channel.IsClosed) { Console.WriteLine("通道已关闭,跳过ACK,消息等待重连后重新投递"); return; } var body = ea.Body.ToArray(); string json = Encoding.UTF8.GetString(body); Console.WriteLine($"收到Canal数据:{json}"); var model = JsonConvert.DeserializeObject<CanalMsg>(json); var sql = BuildSql(model); System.IO.File.AppendAllText("1.txt", sql + "\r\n", Encoding.UTF8); Console.WriteLine($"解析后sql:【{sql}】ts->{model.ts}"); await _channel.BasicAckAsync(ea.DeliveryTag, multiple: false); Console.WriteLine($"消息消费确认完成{DateTime.Now}"); } catch (Exception ex) { Console.WriteLine($"消费异常:{ex.Message}"); if (!_channel.IsClosed) { await _channel.BasicNackAsync(ea.DeliveryTag, multiple: false, requeue: true); } } } #region 解析 private static string BuildSql(CanalMsg dto) { if (dto == null) { return string.Empty; } if (dto.isDdl.HasValue && dto.isDdl.Value) { return dto.sql; } var sb = new StringBuilder(); string db = dto.database; string tb = dto.table; var pks = dto.pkNames ?? new List<string>(); switch (dto.type) { case "INSERT": foreach (var item in dto.data) { sb.AppendLine(BuildInsert(db, tb, item, dto.mysqlType)); } break; case "UPDATE": for (var i = 0; i < dto.data.Count; i++) { sb.AppendLine(BuildUpdate(db, tb, dto.data[i], dto.old[i], pks, dto.mysqlType)); } break; case "DELETE": foreach (var item in dto.data) { sb.AppendLine(BuildDelete(db, tb, item, pks, dto.mysqlType)); } break; case "QUERY": if (dto.data == null && !string.IsNullOrEmpty(dto.sql) && !dto.sql.StartsWith("select", StringComparison.OrdinalIgnoreCase)) { sb.AppendLine(dto.sql); } break; } return sb.ToString(); } private static string BuildInsert(string db, string table, Dictionary<string, object> row, Dictionary<string, string> fieldType) { var keys = row.Keys.ToList(); var cols = string.Join(",", keys.Select(c => $"`{c}`")); var vals = string.Join(",", keys.Select(c => WrapValue(row[c], fieldType[c]))); return $"INSERT INTO `{db}`.`{table}`({cols}) VALUES({vals});"; } private static string BuildUpdate(string db, string table, Dictionary<string, object> row, Dictionary<string, object> old, List<string> pkNames, Dictionary<string, string> fieldType) { var updateFields = old.Keys .Where(field => !pkNames.Contains(field)) .Select(field => $"`{field}`={WrapValue(row[field], fieldType[field])}"); string setSql = string.Join(",", updateFields); var whereConditions = pkNames .Select(pk => $"`{pk}`={WrapValue(row[pk], fieldType[pk])}"); string whereSql = string.Join(" AND ", whereConditions); return $"UPDATE `{db}`.`{table}` SET {setSql} WHERE {whereSql};"; } private static string BuildDelete(string db, string table, Dictionary<string, object> row, List<string> pkNames, Dictionary<string, string> fieldType) { var whereSql = string.Join( " AND ", pkNames.Select(pk => $"`{pk}`={WrapValue(row[pk], fieldType[pk])}") ); return $"DELETE FROM `{db}`.`{table}` WHERE {whereSql};"; } private static string WrapValue(object val, string fullMysqlType) { if (val == null || val == DBNull.Value) return "NULL"; string str = val.ToString().Replace("'", "''"); string baseType = fullMysqlType.Split('(')[0].ToUpper(); var numberTypes = new HashSet<string> { "TINYINT", "SMALLINT", "INT", "BIGINT", "FLOAT", "DOUBLE", "DECIMAL", "NUMERIC", "BIT" }; if (numberTypes.Contains(baseType)) { return str; } return $"'{str}'"; } #endregion }

5.4 kafka

using Confluent.Kafka; using Newtonsoft.Json; using System.Text; using System.Threading.Channels; namespace kafka.client { internal class Program { static void Main(string[] args) { var config = new ConsumerConfig { // kafka服务地址 BootstrapServers = "172.0.10.140:9092", // 消费者组ID GroupId = "g4", AutoOffsetReset = AutoOffsetReset.Earliest, EnableAutoCommit = true, AutoCommitIntervalMs = 1000, SecurityProtocol = SecurityProtocol.Plaintext, }; using (var consumer = new ConsumerBuilder<Ignore, string>(config).Build()) { // 订阅topic名称 consumer.Subscribe("example"); Console.WriteLine("Kafka消费者启动成功,等待接收消息..."); while (true) { try { var consumeResult = consumer.Consume(CancellationToken.None); string msg = consumeResult.Message.Value; long offset = consumeResult.Offset.Value; int partition = consumeResult.Partition.Value; Console.WriteLine($"【收到消息】partition:{partition}, offset:{offset}"); Console.WriteLine($"内容:{msg}"); var model = JsonConvert.DeserializeObject<CanalMsg>(msg); var sql = BuildSql(model); Console.WriteLine($"解析后sql:【{sql}】ts->{model.ts}"); Console.WriteLine($"消息消费确认完成{DateTime.Now}"); File.AppendAllText(DateTime.Now.ToString("yyyy-MM-dd") + ".txt", $"---{partition}---{offset}--{DateTime.Now}--------\r\n{msg}\r\n------------\r\n", System.Text.Encoding.UTF8); File.AppendAllText(DateTime.Now.ToString("yyyy-MM-dd") + "1.txt", $"---{partition}---{offset}--{DateTime.Now}--------\r\n{sql}\r\n------------\r\n", System.Text.Encoding.UTF8); } catch (Exception ex) { Console.WriteLine($"消费异常:{ex.Message}"); } } consumer.Close(); } } #region 解析 private static string BuildSql(CanalMsg dto) { if (dto == null) { return string.Empty; } if (dto.isDdl.HasValue && dto.isDdl.Value) { return dto.sql; } var sb = new StringBuilder(); string db = dto.database; string tb = dto.table; var pks = dto.pkNames ?? new List<string>(); switch (dto.type) { case "INSERT": foreach (var item in dto.data) { sb.AppendLine(BuildInsert(db, tb, item, dto.mysqlType)); } break; case "UPDATE": for (var i = 0; i < dto.data.Count; i++) { sb.AppendLine(BuildUpdate(db, tb, dto.data[i], dto.old[i], pks, dto.mysqlType)); } break; case "DELETE": foreach (var item in dto.data) { sb.AppendLine(BuildDelete(db, tb, item, pks, dto.mysqlType)); } break; case "QUERY": if (dto.data == null && !string.IsNullOrEmpty(dto.sql) && !dto.sql.StartsWith("select", StringComparison.OrdinalIgnoreCase)) { sb.AppendLine(dto.sql); } break; } return sb.ToString(); } private static string BuildInsert(string db, string table, Dictionary<string, object> row, Dictionary<string, string> fieldType) { var keys = row.Keys.ToList(); var cols = string.Join(",", keys.Select(c => $"`{c}`")); var vals = string.Join(",", keys.Select(c => WrapValue(row[c], fieldType[c]))); return $"INSERT INTO `{db}`.`{table}`({cols}) VALUES({vals});"; } private static string BuildUpdate(string db, string table, Dictionary<string, object> row, Dictionary<string, object> old, List<string> pkNames, Dictionary<string, string> fieldType) { var updateFields = old.Keys .Where(field => !pkNames.Contains(field)) .Select(field => $"`{field}`={WrapValue(row[field], fieldType[field])}"); string setSql = string.Join(",", updateFields); var whereConditions = pkNames .Select(pk => $"`{pk}`={WrapValue(row[pk], fieldType[pk])}"); string whereSql = string.Join(" AND ", whereConditions); return $"UPDATE `{db}`.`{table}` SET {setSql} WHERE {whereSql};"; } private static string BuildDelete(string db, string table, Dictionary<string, object> row, List<string> pkNames, Dictionary<string, string> fieldType) { var whereSql = string.Join( " AND ", pkNames.Select(pk => $"`{pk}`={WrapValue(row[pk], fieldType[pk])}") ); return $"DELETE FROM `{db}`.`{table}` WHERE {whereSql};"; } private static string WrapValue(object val, string fullMysqlType) { if (val == null || val == DBNull.Value) return "NULL"; string str = val.ToString().Replace("'", "''"); string baseType = fullMysqlType.Split('(')[0].ToUpper(); var numberTypes = new HashSet<string> { "TINYINT", "SMALLINT", "INT", "BIGINT", "FLOAT", "DOUBLE", "DECIMAL", "NUMERIC", "BIT" }; if (numberTypes.Contains(baseType)) { return str; } return $"'{str}'"; } #endregion } }
版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/14 1:29:05

指令微调模型为何更易复用人类句法?机制、影响与应对策略

最近在分析大语言模型&#xff08;LLM&#xff09;的生成文本时&#xff0c;一个有趣的现象引起了我的注意&#xff1a;经过指令微调&#xff08;Instruction-Tuning&#xff09;的模型&#xff0c;在生成回复时&#xff0c;对人类句法结构的“复用”程度&#xff0c;有时甚至超…

作者头像 李华
网站建设 2026/8/14 1:28:06

门店质检报告审核Agent方案:AI如何用分级审核机制替代80%人工重复判断

一、业务背景&#xff1a;规模扩张带来的审核压力门店从50家扩张至500家&#xff0c;总部日均巡检报告从数十份增长至数百份。巡店照片、自检记录、整改凭证持续汇集。人工逐份核验、逐图比对、逐项判定的传统审核模式&#xff0c;面临的挑战不止是人力成本的线性增长&#xff…

作者头像 李华
网站建设 2026/8/14 1:26:14

175、LLC谐振变换器的PCB设计实战(热管理)

175、LLC谐振变换器的PCB设计实战(热管理) 一、一块烧焦的板子教会我的事 去年夏天,客户那边反馈说我们的LLC电源模块在满载老化时,MOSFET温度飙到了125℃,散热器烫得能煎鸡蛋。拆开看,谐振电容旁边的PCB板已经发黄,焊盘边缘有细微的碳化痕迹。我拿着热成像仪一照,好…

作者头像 李华
网站建设 2026/8/14 1:25:47

spring内置线程池用法

生产线程池疑问生产最近出现一些mq堆积事情&#xff0c;分析代码发现mq消费内部还使用了spring内置线程池。但生产上只配置了spring.task.execution.thread.pool.core-size16&#xff0c;原以为核心线程数是16&#xff0c;但实际日志发现只到task-8&#xff0c;然后加日志打印线…

作者头像 李华
网站建设 2026/8/14 1:25:32

超详细 | 麒麟操作系统安装流程(附安装包下载)

一、国产操作系统技术流派 Linux 的发行版本可以大体分为两类&#xff1a;一类是商业公司维护的发行版本&#xff0c;以著名的 红帽Redhat&#xff08;RHEL&#xff09;为代表&#xff1b;另一类是社区组织维护的发行版本&#xff0c;以 Debian 为代表。 基于不同版本国产操作…

作者头像 李华