由于网络网闸问题,目前采用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_password2.创建同步使用账号
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.82. kafka
docker pull apache/kafka:3.9.03. 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 conf2.rabbitmq
mkdir rabbitmq chmod 777 rabbitmq3.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-management1.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.82.1.2 复制容器内文件到本地目录(注意需要先定位到/run/media/nvme0n1/rlzk/canal/conf)
docker cp canal-tmp:/home/admin/canal-server/conf .2.1.3 移除临时容器
docker rm -f canal-tmp2.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.82.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.03.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 } }