mysql到pg怎么高效_干货 | Debezium实现Mysql到Elasticsearch高效实时同步(示例代码)
題記
來(lái)自Elasticsearch中文社區(qū)的問(wèn)題——
MySQL中表無(wú)唯一遞增字段,也無(wú)唯一遞增時(shí)間字段,該怎么使用logstash實(shí)現(xiàn)MySQL實(shí)時(shí)增量導(dǎo)數(shù)據(jù)到es中?
logstash和kafka_connector都僅支持基于自增id或者時(shí)間戳更新的方式增量同步數(shù)據(jù)。
回到問(wèn)題本身:如果庫(kù)表里沒有相關(guān)字段,該如何處理呢?
本文給出相關(guān)探討和解決方案。
1、 binlog認(rèn)知
1.1 啥是 binlog?
binlog是Mysql sever層維護(hù)的一種二進(jìn)制日志,與innodb引擎中的redo/undo log是完全不同的日志;其主要是用來(lái)記錄對(duì)mysql數(shù)據(jù)更新或潛在發(fā)生更新的SQL語(yǔ)句,并以"事務(wù)"的形式保存在磁盤中。
作用主要有:
1)復(fù)制:達(dá)到master-slave數(shù)據(jù)一致的目的。
2)數(shù)據(jù)恢復(fù):通過(guò)mysqlbinlog工具恢復(fù)數(shù)據(jù)。
3)增量備份。
1.2 阿里的Canal實(shí)現(xiàn)了增量Mysql同步
一圖勝千言,canal是用java開發(fā)的基于數(shù)據(jù)庫(kù)增量日志解析、提供增量數(shù)據(jù)訂閱&消費(fèi)的中間件。
目前,canal主要支持了MySQL的binlog解析,解析完成后才利用canal client 用來(lái)處理獲得的相關(guān)數(shù)據(jù)。目的:增量數(shù)據(jù)訂閱&消費(fèi)。
綜上,使用binlog可以突破logstash或者kafka-connector沒有自增id或者沒有時(shí)間戳字段的限制,實(shí)現(xiàn)增量同步。
2、基于binlog的同步方式
1)基于kafka Connect的Debezium 開源工程,地址:. https://debezium.io/
2)不依賴第三方的獨(dú)立應(yīng)用: Maxwell開源項(xiàng)目,地址:http://maxwells-daemon.io/
由于已經(jīng)部署過(guò)conluent(kafka的企業(yè)版本,自帶zookeeper、kafka、ksql、kafka-connector等),本文僅針對(duì)Debezium展開。
3、Debezium介紹
Debezium是捕獲數(shù)據(jù)實(shí)時(shí)動(dòng)態(tài)變化的開源的分布式同步平臺(tái)。能實(shí)時(shí)捕獲到數(shù)據(jù)源(Mysql、Mongo、PostgreSql)的:新增(inserts)、更新(updates)、刪除(deletes)操作,實(shí)時(shí)同步到Kafka,穩(wěn)定性強(qiáng)且速度非常快。
特點(diǎn):
1)簡(jiǎn)單。無(wú)需修改應(yīng)用程序。可對(duì)外提供服務(wù)。
2)穩(wěn)定。持續(xù)跟蹤每一行的每一處變動(dòng)。
3)快速。構(gòu)建于kafka之上,可擴(kuò)展,經(jīng)官方驗(yàn)證可處理大容量的數(shù)據(jù)。
4、同步架構(gòu)
如圖,Mysql到ES的同步策略,采取“曲線救國(guó)”機(jī)制。
步驟1: 基Debezium的binlog機(jī)制,將Mysql數(shù)據(jù)同步到Kafka。
步驟2: 基于Kafka_connector機(jī)制,將kafka數(shù)據(jù)同步到Elasticsearch。
5、Debezium實(shí)現(xiàn)Mysql到ES增刪改實(shí)時(shí)同步
軟件版本:
confluent:5.1.2;
Debezium:0.9.2_Final;
Mysql:5.7.x.
Elasticsearch:6.6.1
5.1 Debezium安裝
Debezium的安裝只需要把debezium-connector-mysql的壓縮包解壓放到Confluent的解壓后的插件目錄(share/java)中。
MySQL Connector plugin 壓縮包的下載地址:
注意重啟一下confluent,以使得Debezium生效。
5.2 Mysql binlog等相關(guān)配置。
Debezium使用MySQL的binlog機(jī)制實(shí)現(xiàn)數(shù)據(jù)動(dòng)態(tài)變化監(jiān)測(cè),所以需要Mysql提前配置binlog。
核心配置如下,在Mysql機(jī)器的/etc/my.cnf的mysqld下添加如下配置。
1[mysqld]
2
3server-id = 223344
4log_bin = mysql-bin
5binlog_format = row
6binlog_row_image = full
7expire_logs_days = 10
然后,重啟一下Mysql以使得binlog生效。
1systemctl start mysqld.service
5.3 配置connector連接器。
配置confluent路徑目錄 : /etc
創(chuàng)建文件夾命令 :
1mkdir kafka-connect-debezium
在mysql2kafka_debezium.json存放connector的配置信息 :
1[root@localhost kafka-connect-debezium]# cat mysql2kafka_debezium.json
2{
3 "name" : "debezium-mysql-source-0223",
4 "config":
5 {
6 "connector.class" : "io.debezium.connector.mysql.MySqlConnector",
7 "database.hostname" : "192.168.1.22",
8 "database.port" : "3306",
9 "database.user" : "root",
10 "database.password" : "XXXXXX",
11 "database.whitelist" : "kafka_base_db",
12 "table.whitlelist" : "accounts",
13 "database.server.id" : "223344",
14 "database.server.name" : "full",
15 "database.history.kafka.bootstrap.servers" : "192.168.1.22:9092",
16 "database.history.kafka.topic" : "account_topic",
17 "include.schema.changes" : "true" ,
18 "incrementing.column.name" : "id",
19 "database.history.skip.unparseable.ddl" : "true",
20 "transforms": "unwrap,changetopic",
21 "transforms.unwrap.type": "io.debezium.transforms.UnwrapFromEnvelope",
22 "transforms.changetopic.type":"org.apache.kafka.connect.transforms.RegexRouter",
23 "transforms.changetopic.regex":"(.*)",
24 "transforms.changetopic.replacement":"$1-smt"
25 }
26}
注意如下配置:
"database.server.id",對(duì)應(yīng)Mysql中的server-id的配置。
"database.whitelist" : 待同步的Mysql數(shù)據(jù)庫(kù)名。
"table.whitlelist" :待同步的Mysq表名。
重要:“database.history.kafka.topic”:存儲(chǔ)數(shù)據(jù)庫(kù)的Shcema的記錄信息,而非寫入數(shù)據(jù)的topic、
"database.server.name":邏輯名稱,每個(gè)connector確保唯一,作為寫入數(shù)據(jù)的kafka topic的前綴名稱。
坑一:transforms相關(guān)5行配置作用是寫入數(shù)據(jù)格式轉(zhuǎn)換。
如果沒有,輸入數(shù)據(jù)會(huì)包含:before、after記錄修改前對(duì)比信息以及元數(shù)據(jù)信息(source,op,ts_ms等)。
這些信息在后續(xù)數(shù)據(jù)寫入Elasticsearch是不需要的。(注意結(jié)合自己業(yè)務(wù)場(chǎng)景)。
5.4 啟動(dòng)connector
1curl -X POST -H "Content-Type:application/json"
2--data @mysql2kafka_debezium.json.json
3http://192.168.1.22:18083/connectors | jq
5.5 驗(yàn)證寫入是否成功。
5.5.1 查看kafka-topic
1 kafka-topics --list --zookeeper localhost:2181
此處會(huì)看到寫入數(shù)據(jù)topic的信息。
注意新寫入數(shù)據(jù)topic的格式:database.schema.table-smt 三部分組成。
本示例topic名稱:
full.kafka_base_db.account-smt
5.5.2 消費(fèi)數(shù)據(jù)驗(yàn)證寫入是否正常
1./kafka-avro-console-consumer --topic full.kafka_base_db.account-smt --bootstrap-server 192.168.1.22:9092 --from-beginning
至此,Debezium實(shí)現(xiàn)mysql同步kafka完成。
6、kafka-connector實(shí)現(xiàn)kafka同步Elasticsearch
6.1、Kafka-connector介紹
Kafka Connect是一個(gè)用于連接Kafka與外部系統(tǒng)(如數(shù)據(jù)庫(kù),鍵值存儲(chǔ),檢索系統(tǒng)索引和文件系統(tǒng))的框架。
連接器實(shí)現(xiàn)公共數(shù)據(jù)源數(shù)據(jù)(如Mysql、Mongo、Pgsql等)寫入Kafka,或者Kafka數(shù)據(jù)寫入目標(biāo)數(shù)據(jù)庫(kù),也可以自己開發(fā)連接器。
6.2、kafka到ES connector同步配置
配置路徑:
1/home/confluent-5.1.0/etc/kafka-connect-elasticsearch/quickstart-elasticsearch.properties
配置內(nèi)容:
1"connector.class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector",
2"tasks.max": "1",
3"topics": "full.kafka_base_db.account-smt",
4"key.ignore": "true",
5"connection.url": "http://192.168.1.22:9200",
6"type.name": "_doc",
7"name": "elasticsearch-sink-test"
6.3 kafka到ES啟動(dòng)connector
啟動(dòng)命令
1confluent load elasticsearch-sink-test
2-d /home/confluent-5.1.0/etc/kafka-connect-elasticsearch/quickstart-elasticsearch.properties
6.4 Kafka-connctor RESTFul API查看
Mysql2kafka,kafka2ES的connector詳情信息可以借助postman或者瀏覽器或者命令行查看。
1curl -X GET http://localhost:8083/connectors
7、坑復(fù)盤。
坑2: 同步的過(guò)程中可能出現(xiàn)錯(cuò)誤,比如:kafka topic沒法消費(fèi)到數(shù)據(jù)。
排解思路如下:
1)確認(rèn)消費(fèi)的topic是否是寫入數(shù)據(jù)的topic;
2)確認(rèn)同步的過(guò)程中沒有出錯(cuò)。可以借助connector如下命令查看。
1curl -X GET http://localhost:8083/connectors-xxx/status
坑3: Mysql2ES出現(xiàn)日期格式不能識(shí)別。
是Mysql jar包的問(wèn)題,解決方案:在my.cnf中配置時(shí)區(qū)信息即可。
坑4: kafka2ES,ES沒有寫入數(shù)據(jù)。
排解思路:
1)建議:先創(chuàng)建同topic名稱一致的索引,注意:Mapping靜態(tài)自定義,不要?jiǎng)討B(tài)識(shí)別生成。
2)通過(guò)connetor/status排查出錯(cuò)原因,一步步分析。
8、小結(jié)
binlog的實(shí)現(xiàn)突破了字段的限制,實(shí)際上業(yè)界的go-mysql-elasticsearch已經(jīng)實(shí)現(xiàn)。
對(duì)比:logstash、kafka-connector,雖然Debezium“曲線救國(guó)”兩步實(shí)現(xiàn)了實(shí)時(shí)同步,但穩(wěn)定性+實(shí)時(shí)性能相對(duì)不錯(cuò)。
推薦大家使用。大家有好的同步方式也歡迎留言討論交流。
推薦閱讀:
重磅 | 死磕Elasticsearch方法論認(rèn)知清單(2019春節(jié)更新版)
Elasticsearch基礎(chǔ)、進(jìn)階、實(shí)戰(zhàn)第一公眾號(hào)
總結(jié)
以上是生活随笔為你收集整理的mysql到pg怎么高效_干货 | Debezium实现Mysql到Elasticsearch高效实时同步(示例代码)的全部?jī)?nèi)容,希望文章能夠幫你解決所遇到的問(wèn)題。
- 上一篇: 因为懂得所以慈悲张爱玲经典语录(因为懂得
- 下一篇: 怎么避免买到假润康呢?