一次同步多张表是开发中的一般需求。之前研究了很久找到方法,但没有详细总结。 博友前天在线提问,说明这块理解的还不够透彻。 我整理下, 一是为了尽快解决博友问题, 二是加深记忆,便于未来产品开发中快速上手。
1、同步原理
原有ES专栏中有详解,不再赘述。详细请参考我的专栏: 深入详解Elasticsearch 以下是通过ES5.4.0, logstash5.4.1 验证成功。 可以确认的是2.X版本同样可以验证成功。
2、核心配置文件
input {
stdin {
}
jdbc {
type =>
"cxx_article_info"
jdbc_connection_string =>
"jdbc:mysql://110.10.15.37:3306/cxxwb"
jdbc_user =>
"root"
jdbc_password =>
"xxxxx"
record_last_run =>
"true"
use_column_value =>
"true"
tracking_column =>
"id"
last_run_metadata_path =>
"/opt/logstash/bin/logstash_xxy/cxx_info"
clean_run =>
"false"
jdbc_driver_library =>
"/opt/elasticsearch/lib/mysql-connector-java-5.1.38.jar"
jdbc_driver_class =>
"com.mysql.jdbc.Driver"
jdbc_paging_enabled =>
"true"
jdbc_page_size =>
"500"
statement =>
"select * from cxx_article_info where id > :sql_last_value"
schedule =>
"* * * * *"
}
jdbc {
type =>
"cxx_user"
jdbc_connection_string =>
"jdbc:mysql://110.10.15.37:3306/cxxwb"
jdbc_user =>
"root"
jdbc_password =>
"xxxxxx"
record_last_run =>
"true"
use_column_value =>
"true"
tracking_column =>
"id"
last_run_metadata_path =>
"/opt/logstash/bin/logstash_xxy/cxx_user_info"
clean_run =>
"false"
jdbc_driver_library =>
"/opt/elasticsearch/lib/mysql-connector-java-5.1.38.jar"
jdbc_driver_class =>
"com.mysql.jdbc.Driver"
jdbc_paging_enabled =>
"true"
jdbc_page_size =>
"500"
statement =>
"select * from cxx_user_info where id > :sql_last_value"
schedule =>
"* * * * *"
}
}
filter {
mutate {
convert => [
"publish_time",
"string" ]
}
date {
timezone =>
"Europe/Berlin"
match => [
"publish_time" ,
"ISO8601",
"yyyy-MM-dd HH:mm:ss"]
}
json {
source =>
"message"
remove_field => [
"message"]
}
}
output {
if [type]==
"cxxarticle_info" {
elasticsearch {
hosts =>
"10.100.11.231:9200"
index =>
"cxx_info_index"
}
}
if [type]==
"cxx_user" {
elasticsearch {
hosts =>
"10.100.11.231:9200"
index =>
"cxx_user_index"
}
}
}
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104
3、同步成功结果
[
2017-07-19T15:08:05,438][
INFO ][
logstash.pipeline ] Pipeline main started
The stdin plugin is now waiting for input:
[
2017-07-19T15:08:05,491][
INFO ][
logstash.agent ] Successfully started Logstash API endpoint {:port=>9600}
[
2017-07-19T15:09:00,721][
INFO ][
logstash.inputs.jdbc ] (0.007000s) SELECT count(
*) AS `count` FROM (select * from cxx
_article_info where id > 0) AS
`t1` LIMIT 1
[
2017-07-19T15:09:00,721][
INFO ][
logstash.inputs.jdbc ] (0.008000s) SELECT count(
*) AS `count` FROM (select * from cxx
_user_info where id > 0) AS
`t1` LIMIT 1
[
2017-07-19T15:09:00,730][
INFO ][
logstash.inputs.jdbc ] (0.004000s) SELECT
* FROM (select * from cxx
_user_info where id > 0) AS
`t1` LIMIT 500 OFFSET 0
[
2017-07-19T15:09:00,731][
INFO ][
logstash.inputs.jdbc ] (0.007000s) SELECT
* FROM (select * from cxx
_article_info where id > 0) AS
`t1` LIMIT 500 OFFSET 0
[
2017-07-19T15:10:00,173][
INFO ][
logstash.inputs.jdbc ] (0.002000s) SELECT count(
*) AS `count` FROM (select * from cxx
_article_info where id > 3) AS
`t1` LIMIT 1
[
2017-07-19T15:10:00,174][
INFO ][
logstash.inputs.jdbc ] (0.003000s) SELECT count(
*) AS `count` FROM (select * from cxx
_user_info where id > 2) AS
`t1` LIMIT 1
[
2017-07-19T15:11:00,225][
INFO ][
logstash.inputs.jdbc ] (0.001000s) SELECT count(
*) AS `count` FROM (select * from cxx
_article_info where id > 3) AS
`t1` LIMIT 1
[
2017-07-19T15:11:00,225][
INFO ][
logstash.inputs.jdbc ] (0.002000s) SELECT count(
*) AS `count` FROM (select * from cxx
_user_info where id > 2) AS
`t1` LIMIT 1
12345678910111213
4、扩展
1)多个表无非就是在input里面多加几个类型,在output中多加基础 类型判定。 举例:
if [
type]==
"cxx_user"
1
2)input里的type和output if判定的type**保持一致**,该type对应ES中的type。
后记
死磕ES,有问题欢迎大家提问探讨!
—————————————————————————————————— 更多ES相关实战干货经验分享,请扫描下方【铭毅天下】微信公众号二维码关注。 (每周至少更新一篇!)
2017年07月19日 23:32 于家中床前
作者:铭毅天下 转载请标明出处,原文地址: http://blog.csdn.net/laoyang360/article/details/75452953 如果感觉本文对您有帮助,请点击‘顶’支持一下,您的支持是我坚持写作最大的动力,谢谢!