
1.准备
1.1.组件
JDK:1.8版本及以上;
ElasticSearch:6.x版本,目前貌似不支持7.x版本;
Kibana:6.x版本;
**Canal.deployer:**1.1.4
**Canal.Adapter:**1.1.4
1.2.配置
- 需要先开启MySQL的 binlog 写入功能,配置 binlog-format 为 ROW 模式
找到my.cnf文件,我的目录是/etc/my.cnf,添加以下配置:
1log-bin=mysql-bin # 开启 binlog 2binlog-format=ROW # 选择 ROW 模式 3server_id=1 # 配置 MySQL replaction 需要定义,不要和 canal 的 slaveId 重复

然后重启mysql,用以下命令检查一下binlog是否正确启动:
1mysql> show variables like 'log_bin%'; 2+---------------------------------+----------------------------------+ 3| Variable_name | Value | 4+---------------------------------+----------------------------------+ 5| log_bin | ON | 6| log_bin_basename | /data/mysql/data/mysql-bin | 7| log_bin_index | /data/mysql/data/mysql-bin.index | 8| log_bin_trust_function_creators | OFF | 9| log_bin_use_v1_row_events | OFF | 10+---------------------------------+----------------------------------+ 115 rows in set (0.00 sec) 12mysql> show variables like 'binlog_format%'; 13+---------------+-------+ 14| Variable_name | Value | 15+---------------+-------+ 16| binlog_format | ROW | 17+---------------+-------+ 181 row in set (0.00 sec)
-
授权 canal 链接 MySQL 账号具有作为 MySQL slave 的权限, 如果已有账户可直接 grant
CREATE USER canal IDENTIFIED BY 'Aa123456.'; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON . TO 'canal'@'%'; FLUSH PRIVILEGES;
2.安装
2.1.ElasticSearch
安装配置方法:https://www.cnblogs.com/caoweixiong/p/11826295.html
2.2.canal.deployer
2.2.1.下载解压
1直接下载 2 3访问:https://github.com/alibaba/canal/releases ,会列出所有历史的发布版本包 下载方式,比如以1.1.4版本为例子: 4wget https://github.com/alibaba/canal/releases/download/canal-1.1.4/canal.deployer-1.1.4.tar.gz 5 6自己编译 7git clone git@github.com:alibaba/canal.git 8cd canal; 9mvn clean install -Dmaven.test.skip -Denv=release 10编译完成后,会在根目录下产生target/canal.deployer-$version.tar.gz 11 12mkdir /usr/local/canal 13tar zxvf canal.deployer-1.1.4.tar.gz -C /usr/local/canal
解压完成后,进入 /usr/local/canal目录,可以看到如下结构:

2.2.2.配置
-
配置server
cd /usr/local/canal/conf vi canal.properties
标红的需要我们重点关注的,也是平常修改最多的参数:
1################################################# 2######### common argument ############# 3################################################# 4# tcp bind ip 5canal.ip = 6# register ip to zookeeper 7canal.register.ip = #运行canal-server服务的主机IP,可以不用配置,他会自动绑定一个本机的IP 8canal.port = 11111 #canal-server监听的端口(TCP模式下,非TCP模式不监听1111端口) 9canal.metrics.pull.port = 11112 #canal-server metrics.pull监听的端口 10# canal instance user/passwd 11# canal.user = canal 12# canal.passwd = E3619321C1A937C46A0D8BD1DAC39F93B27D4458 13 14# canal admin config 15#canal.admin.manager = 127.0.0.1:8089 16canal.admin.port = 11110 17canal.admin.user = admin 18canal.admin.passwd = 4ACFE3202A5FF5CF467898FC58AAB1D615029441 19 20canal.zkServers = #集群模式下要配置zookeeper进行协调配置,单机模式可以不用配置 21# flush data to zk 22canal.zookeeper.flush.period = 1000 23canal.withoutNetty = false 24# tcp, kafka, RocketMQ 25canal.serverMode = tcp #canal-server运行的模式,TCP模式就是直连客户端,不经过中间件。kafka和mq是消息队列的模式 26# flush meta cursor/parse position to file 27canal.file.data.dir = ${canal.conf.dir} #存放数据的路径 28canal.file.flush.period = 1000 29## memory store RingBuffer size, should be Math.pow(2,n) 30canal.instance.memory.buffer.size = 16384 31## memory store RingBuffer used memory unit size , default 1kb #下面是一些系统参数的配置,包括内存、网络等 32canal.instance.memory.buffer.memunit = 1024 33## meory store gets mode used MEMSIZE or ITEMSIZE 34canal.instance.memory.batch.mode = MEMSIZE 35canal.instance.memory.rawEntry = true 36 37## detecing config #这里是心跳检查的配置,做HA时会用到 38canal.instance.detecting.enable = false 39#canal.instance.detecting.sql = insert into retl.xdual values(1,now()) on duplicate key update x=now() 40canal.instance.detecting.sql = select 1 41canal.instance.detecting.interval.time = 3 42canal.instance.detecting.retry.threshold = 3 43canal.instance.detecting.heartbeatHaEnable = false 44 45# support maximum transaction size, more than the size of the transaction will be cut into multiple transactions delivery 46canal.instance.transaction.size = 1024 47# mysql fallback connected to new master should fallback times 48canal.instance.fallbackIntervalInSeconds = 60 49 50# network config 51canal.instance.network.receiveBufferSize = 16384 52canal.instance.network.sendBufferSize = 16384 53canal.instance.network.soTimeout = 30 54 55# binlog filter config #binlog过滤的配置,指定过滤那些SQL 56canal.instance.filter.druid.ddl = true 57canal.instance.filter.query.dcl = false 58canal.instance.filter.query.dml = false 59canal.instance.filter.query.ddl = false 60canal.instance.filter.table.error = false 61canal.instance.filter.rows = false 62canal.instance.filter.transaction.entry = false 63 64# binlog format/image check #binlog格式检测,使用ROW模式,非ROW模式也不会报错,但是同步不到数据 65canal.instance.binlog.format = ROW,STATEMENT,MIXED 66canal.instance.binlog.image = FULL,MINIMAL,NOBLOB 67 68# binlog ddl isolation 69canal.instance.get.ddl.isolation = false 70 71# parallel parser config 72canal.instance.parser.parallel = true #并行解析配置,如果是单个CPU就把下面这个true改为false 73## concurrent thread number, default 60% available processors, suggest not to exceed Runtime.getRuntime().availableProcessors() 74#canal.instance.parser.parallelThreadSize = 16 75## disruptor ringbuffer size, must be power of 2 76canal.instance.parser.parallelBufferSize = 256 77 78# table meta tsdb info 79canal.instance.tsdb.enable = true 80canal.instance.tsdb.dir = ${canal.file.data.dir:../conf}/${canal.instance.destination:} 81canal.instance.tsdb.url = jdbc:h2:${canal.instance.tsdb.dir}/h2;CACHE_SIZE=1000;MODE=MYSQL; 82canal.instance.tsdb.dbUsername = canal 83canal.instance.tsdb.dbPassword = canal 84# dump snapshot interval, default 24 hour 85canal.instance.tsdb.snapshot.interval = 24 86# purge snapshot expire , default 360 hour(15 days) 87canal.instance.tsdb.snapshot.expire = 360 88 89# aliyun ak/sk , support rds/mq 90canal.aliyun.accessKey = 91canal.aliyun.secretKey = 92 93################################################# 94######### destinations ############# 95################################################# 96canal.destinations = example #canal-server创建的实例,在这里指定你要创建的实例的名字,比如test1,test2等,逗号隔开 97# conf root dir 98canal.conf.dir = ../conf 99# auto scan instance dir add/remove and start/stop instance 100canal.auto.scan = true 101canal.auto.scan.interval = 5 102 103canal.instance.tsdb.spring.xml = classpath:spring/tsdb/h2-tsdb.xml 104#canal.instance.tsdb.spring.xml = classpath:spring/tsdb/mysql-tsdb.xml 105 106canal.instance.global.mode = spring 107canal.instance.global.lazy = false 108canal.instance.global.manager.address = ${canal.admin.manager} 109#canal.instance.global.spring.xml = classpath:spring/memory-instance.xml 110canal.instance.global.spring.xml = classpath:spring/file-instance.xml 111#canal.instance.global.spring.xml = classpath:spring/default-instance.xml 112 113################################################## 114######### MQ ############# 115################################################## 116canal.mq.servers = 127.0.0.1:6667 117canal.mq.retries = 0 118canal.mq.batchSize = 16384 119canal.mq.maxRequestSize = 1048576 120canal.mq.lingerMs = 100 121canal.mq.bufferMemory = 33554432 122canal.mq.canalBatchSize = 50 123canal.mq.canalGetTimeout = 100 124canal.mq.flatMessage = true 125canal.mq.compressionType = none 126canal.mq.acks = all 127#canal.mq.properties. = 128canal.mq.producerGroup = test 129# Set this value to "cloud", if you want open message trace feature in aliyun. 130canal.mq.accessChannel = local 131# aliyun mq namespace 132#canal.mq.namespace = 133 134################################################## 135######### Kafka Kerberos Info ############# 136################################################## 137canal.mq.kafka.kerberos.enable = false 138canal.mq.kafka.kerberos.krb5FilePath = "../conf/kerberos/krb5.conf" 139canal.mq.kafka.kerberos.jaasFilePath = "../conf/kerberos/jaas.conf"
- 配置example
在根配置文件中创建了实例名称之后,需要在根配置的同级目录下创建该实例目录,canal-server为我们提供了一个示例的实例配置,因此我们可以直接复制该示例,举个例子吧:根配置配置了如下实例:
1[root@aliyun conf]# vim canal.properties 2... 3canal.destinations = user_order,delivery_info 4... 5 6我们需要在根配置的同级目录下创建这两个实例 7[root@aliyun conf]# pwd 8/usr/local/canal-server/conf 9[root@aliyun conf]# cp -a example/ user_order 10[root@aliyun conf]# cp -a example/ delivery_info
这里只举例1个example的配置:
vi /usr/local/canal/conf/example/instance.properties
标红的需要我们重点关注的,也是平常修改最多的参数:
1################################################### mysql serverId , v1.0.26+ will autoGencanal.instance.mysql.slaveId=11# enable gtid use true/false 2canal.instance.gtidon=false 3 4# position info 5canal.instance.master.address=172.16.10.26:3306 #指定要读取binlog的MySQL的IP地址和端口 6canal.instance.master.journal.name= #从指定的binlog文件开始读取数据 7canal.instance.master.position= #指定偏移量,做过主从复制的应该都理解这两个参数。 #tips:binlog和偏移量也可以不指定,则canal-server会从当前的位置开始读取。我建议不设置canal.instance.master.timestamp= 8canal.instance.master.gtid= 9 10# rds oss binlog 11canal.instance.rds.accesskey= 12canal.instance.rds.secretkey= 13canal.instance.rds.instanceId= 14 15# table meta tsdb info 16canal.instance.tsdb.enable=true 17#canal.instance.tsdb.url=jdbc:mysql://127.0.0.1:3306/canal_tsdb 18#canal.instance.tsdb.dbUsername=canal 19#canal.instance.tsdb.dbPassword=canal 20#这几个参数是设置高可用配置的,可以配置mysql从库的信息 21#canal.instance.standby.address = 22#canal.instance.standby.journal.name = 23#canal.instance.standby.position = 24#canal.instance.standby.timestamp = 25#canal.instance.standby.gtid= 26 27# username/password 28canal.instance.dbUsername=canal #指定连接mysql的用户密码 29canal.instance.dbPassword=Aa123456. 30canal.instance.connectionCharset = UTF-8 #字符集 31# enable druid Decrypt database password 32canal.instance.enableDruid=false 33#canal.instance.pwdPublicKey=MFwwDQYJKoZIhvcNAQEBBQADSwAwSAJBALK4BUxdDltRRE5/zXpVEVPUgunvscYFtEip3pmLlhrWpacX7y7GCMo2/JM6LeHmiiNdH1FWgGCpUfircSwlWKUCAwEAAQ== 34 35# table regex 36#canal.instance.filter.regex=.*\\..* 37canal.instance.filter.regex=risk.canal,risk.cwx #这个是比较重要的参数,匹配库表白名单,比如我只要test库的user表的增量数据,则这样写 test.user 38# table black regex 39canal.instance.filter.black.regex= 40# table field filter(format: schema1.tableName1:field1/field2,schema2.tableName2:field1/field2) 41#canal.instance.filter.field=test1.t_product:id/subject/keywords,test2.t_company:id/name/contact/ch 42# table field black filter(format: schema1.tableName1:field1/field2,schema2.tableName2:field1/field2) 43#canal.instance.filter.black.field=test1.t_product:subject/product_image,test2.t_company:id/name/contact/ch 44 45# mq config 46canal.mq.topic=example 47# dynamic topic route by schema or table regex 48#canal.mq.dynamicTopic=mytest1.user,mytest2\\..*,.*\\..* 49canal.mq.partition=0 50# hash partition config 51#canal.mq.partitionsNum=3 52#canal.mq.partitionHash=test.table:id^name,.*\\..* 53#################################################
2.2.3.启动
bin/startup.sh
-
查看 server 日志
vi logs/canal/canal.log
2013-02-05 22:45:27.967 [main] INFO com.alibaba.otter.canal.deployer.CanalLauncher - ## start the canal server. 2013-02-05 22:45:28.113 [main] INFO com.alibaba.otter.canal.deployer.CanalController - ## start the canal server[10.1.29.120:11111] 2013-02-05 22:45:28.210 [main] INFO com.alibaba.otter.canal.deployer.CanalLauncher - ## the canal server is running now ......
-
查看 instance 的日志
vi logs/example/example.log
2013-02-05 22:50:45.636 [main] INFO c.a.o.c.i.spring.support.PropertyPlaceholderConfigurer - Loading properties file from class path resource [canal.properties] 2013-02-05 22:50:45.641 [main] INFO c.a.o.c.i.spring.support.PropertyPlaceholderConfigurer - Loading properties file from class path resource [example/instance.properties] 2013-02-05 22:50:45.803 [main] INFO c.a.otter.canal.instance.spring.CanalInstanceWithSpring - start CannalInstance for 1-example 2013-02-05 22:50:45.810 [main] INFO c.a.otter.canal.instance.spring.CanalInstanceWithSpring - start successful....
-
关闭
bin/stop.sh
2.3.canal.adapter
2.3.1.下载解压
1访问:https://github.com/alibaba/canal/releases ,会列出所有历史的发布版本包 下载方式,比如以1.1.4版本为例子: 2wget https://github.com/alibaba/canal/releases/download/canal-1.1.4/canal.adapter-1.1.4.tar.gz 3 4mkdir /usr/local/canal-adapter 5tar zxvf canal.adapter-1.1.4.tar.gz -C /usr/local/canal-adapter
2.3.2.配置
-
adapter配置
cd /usr/local/canal-adapter vim conf/application.yml
标红的需要我们重点关注的,也是平常修改最多的参数:
1server: 2 port: 8081 3spring: 4 jackson: 5 date-format: yyyy-MM-dd HH:mm:ss 6 time-zone: GMT+8 7 default-property-inclusion: non_null 8 9canal.conf: 10 mode: tcp #模式 11 canalServerHost: 127.0.0.1:11111 #指定canal-server的地址和端口 12# zookeeperHosts: slave1:2181 13# mqServers: 127.0.0.1:9092 #or rocketmq 14# flatMessage: true 15 batchSize: 500 16 syncBatchSize: 1000 17 retries: 0 18 timeout: 19 accessKey: 20 secretKey: 21 srcDataSources: #数据源配置,从哪里获取数据 22 defaultDS: #指定一个名字,在ES的配置中会用到,唯一 23 url: jdbc:mysql://172.16.10.26:3306/risk?useUnicode=true 24 username: root 25 password: 123456 26 canalAdapters: 27 - instance: example # canal instance Name or mq topic name #指定在canal-server配置的实例 28 groups: 29 - groupId: g1 #默认就好,组标识 30 outerAdapters: 31 - name: logger 32# - name: rdb 33# key: mysql1 34# properties: 35# jdbc.driverClassName: com.mysql.jdbc.Driver 36# jdbc.url: jdbc:mysql://127.0.0.1:3306/mytest2?useUnicode=true 37# jdbc.username: root 38# jdbc.password: 121212 39# - name: rdb 40# key: oracle1 41# properties: 42# jdbc.driverClassName: oracle.jdbc.OracleDriver 43# jdbc.url: jdbc:oracle:thin:@localhost:49161:XE 44# jdbc.username: mytest 45# jdbc.password: m121212 46# - name: rdb 47# key: postgres1 48# properties: 49# jdbc.driverClassName: org.postgresql.Driver 50# jdbc.url: jdbc:postgresql://localhost:5432/postgres 51# jdbc.username: postgres 52# jdbc.password: 121212 53# threads: 1 54# commitSize: 3000 55# - name: hbase 56# properties: 57# hbase.zookeeper.quorum: 127.0.0.1 58# hbase.zookeeper.property.clientPort: 2181 59# zookeeper.znode.parent: /hbase 60 - name: es #输出到哪里,指定es 61 hosts: 172.16.99.2:40265 #指定es的地址,注意端口为es的传输端口9300 62 properties: 63 # mode: transport # or rest 64 # security.auth: test:123456 # only used for rest mode 65 cluster.name: log-es-cluster #指定es的集群名称
-
es配置
[root@aliyun es]# pwd /usr/local/canal-adapter/conf/es [root@aliyun es]# ll total 12 -rwxrwxrwx 1 root root 466 Apr 4 10:27 biz_order.yml #这三个配置文件是自带的,可以删除,不过最好不要删除,因为可以参考他的格式 -rwxrwxrwx 1 root root 855 Apr 4 10:27 customer.yml -rwxrwxrwx 1 root root 416 Apr 4 10:27 mytest_user.yml
创建canal.yml文件:
cp customer.yml canal.ymlvim conf/es/canal.yml
标红的需要我们重点关注的,也是平常修改最多的参数:
1dataSourceKey: defaultDS #指定数据源,这个值和adapter的application.yml文件中配置的srcDataSources值对应。 2destination: example #指定canal-server中配置的某个实例的名字,注意:我们可能配置多个实例,你要清楚的知道每个实例收集的是那些数据,不要瞎搞。 3groupId: g1 #组ID,默认就好 4esMapping: #ES的mapping(映射) 5 _index: canal #要同步到的ES的索引名称(自定义),需要自己在ES上创建哦! 6 _type: _doc #ES索引的类型名称(自定义) 7 _id: _id #ES标示文档的唯一标示,通常对应数据表中的主键ID字段,注意我这里写成的是"_id",有个下划线哦! #pk: id #如果不需要_id, 则需要指定一个属性为主键属性 sql: "select t.id as _id, t.name, t.sex, t.age, t.amount, t.email, t.occur_time from canal t" #这里就是数据表中的每个字段到ES索引中叫什么名字的sql映射,注意映射到es中的每个字段都要是唯一的,不能重复。 #etlCondition: "where t.occur_time>='{0}'" commitBatch: 3000
sql映射文件写完之后,要去ES上面创建对应的索引和映射,映射要求要和sql文件的映射保持一致,即sql映射中有的字段在ES的索引映射中必须要有,否则同步会报字段错误,导致失败。
2.3.3.创建mysql表和es索引
1CREATE TABLE `canal` ( 2id int(11) NOT NULL AUTO_INCREMENT, 3name varchar(20) NULL COMMENT '名称', 4sex varchar(2) NULL COMMENT '性别', 5age int NULL COMMENT '年龄', 6amount decimal(12,2) NULL COMMENT '资产', 7email varchar(50) NULL COMMENT '邮箱', 8occur_time timestamp NOT NULL DEFAULT CURRENT_TIMESTAMP, 9PRIMARY KEY (`ID`) 10) ENGINE=InnoDB AUTO_INCREMENT=3 DEFAULT CHARSET=utf8; 11 12{ 13 "mappings": { 14 "_doc": { 15 "properties": { 16 "id": { 17 "type": "long" 18 }, 19 "name": { 20 "type": "text" 21 }, 22 "sex": { 23 "type": "text" 24 }, 25 "age": { 26 "type": "long" 27 }, 28 "amount": { 29 "type": "text" 30 }, 31 "email": { 32 "type": "text" 33 }, 34 "occur_time": { 35 "type": "date" 36 } 37 } 38 } 39 } 40}
2.3.4.启动
1cd /usr/local/canal-adapter 2./bin/startup.sh
查看日志:
cat logs/adapter/adapter.log

2.4.Kibana
安装配置方法:https://www.cnblogs.com/caoweixiong/p/11826655.html
3.验证
- 没有数据时:

-
插入1条数据:
insert into canal(id,name,sex,age,amount,email,occur_time) values(null,'cwx','男',18,100000000,'249299170@qq.com',now());


-
更新1条数据:
update canal set name='cwx1',sex='女',age=28,amount=200000,email='asdf',occur_time=now() where id=16;


-
删除1条数据:
delete from canal where id=16;


4.总结
4.1.全量更新不能实现,但是增删改都是可以的;
4.2.一定要提前创建好es索引;
4.3.es配置的是tcp端口,比如默认的9300;
4.4.目前es貌似支持6.x版本,不支持7.x版本;