由于网络网闸问题目前采用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_pluginmysql_native_password2.创建同步使用账号CREATE USER canal% IDENTIFIED BY canal113.; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO canal%; -- (可能不需要)MySQL8.0 必须重置认证方式为 native否则canal连接报错 ALTER USER canal% IDENTIFIED WITH mysql_native_password BY canal113.; FLUSH PRIVILEGES; --验证是否有权限 SHOW GRANTS FOR canal%; ---结果如下说明成功了 GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO canal%3.验证是否开启binlog,并且是ROW---验证 binlog 是否开启 看到是 log_bin On binlog_formatROW SHOW VARIABLES LIKE log_bin; SHOW VARIABLES LIKE binlog_format;二、安装所需镜像1. canal-serverdocker pull canal/canal-server:v1.1.82. kafkadocker pull apache/kafka:3.9.03. rabbitmqdocker pull rabbitmq:4.3.2-management三、创建所需要的文件夹本文演示路径是/run/media/nvme0n1/rlzk1.canal-servermkdir canal chmod 777 canal cd canal/ mkdir logs mkdir data mkdir conf chmod 777 data chmod 777 logs chmod 777 conf2.rabbitmqmkdir rabbitmq chmod 777 rabbitmq3.kafkamkdir kafka chmod 777 kafka四、容器创建及使用1. rabbitmq1.1 创建容器(注意其中的账号密码docker run \ --name rabbitMQ \ --network xx-network \ --add-hosthost.docker.internal:host-gateway \ -e TZAsia/Shanghai \ -p 7008:5672 \ -p 15672:15672 \ -v /run/media/nvme0n1/rlzk/rabbitmq:/var/lib/rabbitmq \ -e RABBITMQ_DEFAULT_USERrabbitmq \ -e RABBITMQ_DEFAULT_PASSrabbitmq123654 \ --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-server2.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/confdocker cp canal-tmp:/home/admin/canal-server/conf .2.1.3 移除临时容器docker rm -f canal-tmp2.2 修改配置文件canal.propertiescd /run/media/nvme0n1/rlzk/canal/conf vi canal.properties --找到 添加 admin 和密码 canal.admin.user admin canal.admin.passwd admin123654 --找到 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 rabbitmq123654 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.addresshost.docker.internal:3306 --找到设置连接数据库的账号和密码设置这个账号在一中创建如下 canal.instance.dbUsernamecanal canal.instance.dbPasswordcanal113. --找到过滤只要所需的数据库 canal.instance.filter.regexxxx_db\\..* ---如果要从指定位置开始同步 找到其中的binlog 和 pos 进行修改以下 两个值是假的 # 指定binlog文件名 canal.instance.master.journal.namebinlog.000326 # 指定偏移pos canal.instance.master.position89100413 # 时间戳留空二选一用文件pos更精准 canal.instance.master.timestamp 修改后进行保存2.4 创建容器(canal-server)docker run \ --name canal-server \ --network xx-network \ --add-hosthost.docker.internal:host-gateway \ -e TZAsia/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.kafka3.1 创建容器(注意其中的172.0.10.140这个到时客户端连接会使用到)docker run \ --name kafka \ -p 9092:9092 \ --network xx-network \ --add-hosthost.docker.internal:host-gateway \ -e TZAsia/Shanghai \ --restart always \ -e KAFKA_NODE_ID1 \ -e KAFKA_PROCESS_ROLESbroker,controller \ -e KAFKA_CONTROLLER_QUORUM_VOTERS1kafka:9093 \ -e KAFKA_LISTENERSPLAINTEXT://0.0.0.0:9092,CONTROLLER://kafka:9093 \ -e KAFKA_ADVERTISED_LISTENERSPLAINTEXT://172.0.10.140:9092 \ -e KAFKA_CONTROLLER_LISTENER_NAMESCONTROLLER \ -e KAFKA_LOG_DIRS/opt/kraft-data \ -e KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR1 \ -e KAFKA_DEFAULT_REPLICATION_FACTOR1 \ -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.05.2 通用类using System; using System.Collections.Generic; using System.Linq; using System.Text; using System.Threading.Tasks; namespace CanalSync { public class CanalMsg { public ListDictionarystring, 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 Dictionarystring, string mysqlType { get; set; } public ListDictionarystring, object old { get; set; } public Liststring pkNames { get; set; } public string sql { get; set; } //public Dictionarystring,string sqlType { get; set; } /// summary /// /// /summary public string table { get; set; } public long ts { get; set; } public string type { get; set; } } }5.3 rabbitmqusing 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 rabbitmq123654, 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.DeserializeObjectCanalMsg(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 Liststring(); 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, Dictionarystring, object row, Dictionarystring, 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, Dictionarystring, object row, Dictionarystring, object old, Liststring pkNames, Dictionarystring, 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, Dictionarystring, object row, Liststring pkNames, Dictionarystring, 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 HashSetstring { TINYINT, SMALLINT, INT, BIGINT, FLOAT, DOUBLE, DECIMAL, NUMERIC, BIT }; if (numberTypes.Contains(baseType)) { return str; } return ${str}; } #endregion }5.4 kafkausing 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 ConsumerBuilderIgnore, 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.DeserializeObjectCanalMsg(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 Liststring(); 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, Dictionarystring, object row, Dictionarystring, 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, Dictionarystring, object row, Dictionarystring, object old, Liststring pkNames, Dictionarystring, 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, Dictionarystring, object row, Liststring pkNames, Dictionarystring, 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 HashSetstring { TINYINT, SMALLINT, INT, BIGINT, FLOAT, DOUBLE, DECIMAL, NUMERIC, BIT }; if (numberTypes.Contains(baseType)) { return str; } return ${str}; } #endregion } }