【问题标题】:Logstash input jdbc is duplicating resultsLogstash 输入 jdbc 重复结果
【发布时间】:2015-12-06 20:30:34
【问题描述】:

我正在使用 logstash 输入 jdbc 插件来读取两个(或更多)数据库并将数据发送到 elasticsearch,并使用 kibana 4 将这些数据可视化。

这是我的 logstash 配置:

input {
  jdbc {
    type => "A"
    jdbc_driver_library => "C:\DEV\elasticsearch-1.7.1\plugins\elasticsearch-jdbc-1.7.1.0\lib\jtds-1.3.1.jar"
    jdbc_driver_class => "Java::net.sourceforge.jtds.jdbc.Driver"
    jdbc_connection_string => "jdbc:jtds:sqlserver://dev_data_base_server:1433/dbApp1;domain=CORPDOMAIN;useNTLMv2=true"
    jdbc_user => "user"
    jdbc_password => "pass"
    schedule => "5 * * * *"
    statement => "SELECT id, date, content, status from test_table"
  }

jdbc {
    type => "B"
    jdbc_driver_library => "C:\DEV\elasticsearch-1.7.1\plugins\elasticsearch-jdbc-1.7.1.0\lib\jtds-1.3.1.jar"
    jdbc_driver_class => "Java::net.sourceforge.jtds.jdbc.Driver"
    jdbc_connection_string => "jdbc:jtds:sqlserver://dev_data_base_server:1433/dbApp2;domain=CORPDOMAIN;useNTLMv2=true"
    jdbc_user => "user"
    jdbc_password => "pass"
    schedule => "5 * * * *"
    statement => "SELECT id, date, content, status from test_table"
  }
}
filter {

}
output {

    if [type] == "A" {
        elasticsearch {
            host => "localhost"
            protocol => http
            index => "logstash-servera-%{+YYYY.MM.dd}"
        }    
    }
    if [type] == "B" {
        elasticsearch {
            host => "localhost"
            protocol => http
            index => "logstash-serverb-%{+YYYY.MM.dd}"
        }    
    }

  stdout { codec => rubydebug }
}

问题是每次运行logstash,它都会开始保存所有已经在elasticsearch中的数据。

使用 where 子句 = date > '2015-09-10' 运行后,我停止了 logstash 并使用“特殊参数”再次运行(使用 --debug):sql_last_date。 logstash启动后开始在日志中显示这个:

←[36mExecuting JDBC query {:statement=>"SELECT \n\tSUBSTRING(R.RECEBEDOR, 1, 2)
AS 'DDD',\nCASE WHEN R.STATUS <>  'RCON' AND R.COD_RESPOSTA in (428,429,230,425,
430,427,418,422,415,424,214,433,435,207,426) THEN 'REGRA DE NEGÓCIO'  \n       W
HEN R.STATUS = 'RCON' THEN 'SUCESSO'\n\t   ELSE 'ERRO'\n   END AS 'TIPO_MENSAGEM
',\nAP.ALIAS as 'CANAL', R.ID_RECARGA, R.VALOR, R.STATUS, R.COD_RESPOSTA, R.DESC
_RESPOSTA, R.DT_RECARGA as '@timestamp', R.ID_CLIENTE, R.ID_DEPENDENTE, R.ID_APL
ICACAO, RECEBEDOR, R.ID_OPERADORA, R.TIPO_PRODUTO \n\nFROM RECARGA R (NOLOCK)\nJ
OIN APLICACAO AP ON R.ID_APLICACAO = AP.ID_APLICACAO \nwhere R.DT_RECARGA > :sql
_last_start\nORDER BY R.DT_RECARGA ASC", :parameters=>{:sql_last_start=>2015-09-
10 18:48:00 UTC}, :level=>:debug, :file=>"/DEV/logstash-1.5.4/vendor/bundle/jrub
y/1.9/gems/logstash-input-jdbc-1.0.0/lib/logstash/plugin_mixins/jdbc.rb", :line=
>"107", :method=>"execute_statement"}←[0m

这一次我用“真实”的语句运行:

SELECT 
    SUBSTRING(R.RECEBEDOR, 1, 2) AS 'DDD',
CASE WHEN R.STATUS <>  'RCON' AND R.COD_RESPOSTA in (428,429,230,425,430,427,418,422,415,424,214,433,435,207,426) THEN 'REGRA DE NEGÓCIO'  
       WHEN R.STATUS = 'RCON' THEN 'SUCESSO'
       ELSE 'ERRO'
   END AS 'TIPO_MENSAGEM',
AP.ALIAS as 'CANAL', R.ID_RECARGA, R.VALOR, R.STATUS, R.COD_RESPOSTA, R.DESC_RESPOSTA, R.DT_RECARGA as '@timestamp', R.ID_CLIENTE, R.ID_DEPENDENTE, R.ID_APLICACAO, RECEBEDOR, R.ID_OPERADORA

FROM RECARGA R (NOLOCK)
JOIN APLICACAO AP ON R.ID_APLICACAO = AP.ID_APLICACAO 
where R.DT_RECARGA > :sql_last_start
ORDER BY R.DT_RECARGA ASC

有人知道怎么解决吗?

谢谢!

【问题讨论】:

    标签: elasticsearch logstash kibana-4


    【解决方案1】:

    sql_last_start 现在是 sql_last_valueplease check here 特殊参数 sql_last_start 现在重命名为 sql_last_value 以更清晰,因为它不仅限于日期时间,还可能具有其他列类型。 所以现在的解决方案可能是这样的

    input {
    jdbc {
         type => "A"
         jdbc_driver_library => "C:\DEV\elasticsearch-1.7.1\plugins\elasticsearch-  jdbc-1.7.1.0\lib\jtds-1.3.1.jar"
         jdbc_driver_class => "Java::net.sourceforge.jtds.jdbc.Driver"
         jdbc_connection_string => "jdbc:jtds:sqlserver://dev_data_base_server:1433/dbApp1;domain=CORPDOMAIN;useNTLMv2=true"
         jdbc_user => "user"
         jdbc_password => "pass"
         schedule => "5 * * * *"
         use_column_value => true
         tracking_column => date
         statement => "SELECT id, date, content, status from test_table WHERE date >:sql_last_value"
        #clean_run true means it will reset sql_last_value to zero or initial value if datatype is date(default is also false)
         clean_run =>false
       }
    jdbc{
      #for type B....
      }
    }
    

    我已经用 sql Server DB 测试过

    请使用 clean_run=>ture 首次运行以避免数据类型错误,而在开发过程中,我们可能在 sql_last_value 变量中存储了不同的数据类型值

    【讨论】:

    • 请提供必要的信息来回答网站上的问题。
    • 早些时候我只是一个评论,因为我没有足够的排名来评论上述答案。现在我写了完整的答案,所以它可以帮助别人..:)
    【解决方案2】:

    默认情况下,jdbc 输入将执行配置的 SQL 语句。在您的情况下,您的语句选择 test_table 中的所有内容。您需要通过在 SQL 查询中使用预定义的 sql_last_start 参数来指示您的 SQL 语句仅加载上次运行 jdbc 输入的数据。

    input {
      jdbc {
        type => "A"
        jdbc_driver_library => "C:\DEV\elasticsearch-1.7.1\plugins\elasticsearch-jdbc-1.7.1.0\lib\jtds-1.3.1.jar"
        jdbc_driver_class => "Java::net.sourceforge.jtds.jdbc.Driver"
        jdbc_connection_string => "jdbc:jtds:sqlserver://dev_data_base_server:1433/dbApp1;domain=CORPDOMAIN;useNTLMv2=true"
        jdbc_user => "user"
        jdbc_password => "pass"
        schedule => "5 * * * *"
        statement => "SELECT id, date, content, status from test_table WHERE date > :sql_last_start"
      }
    
    jdbc {
        type => "B"
        jdbc_driver_library => "C:\DEV\elasticsearch-1.7.1\plugins\elasticsearch-jdbc-1.7.1.0\lib\jtds-1.3.1.jar"
        jdbc_driver_class => "Java::net.sourceforge.jtds.jdbc.Driver"
        jdbc_connection_string => "jdbc:jtds:sqlserver://dev_data_base_server:1433/dbApp2;domain=CORPDOMAIN;useNTLMv2=true"
        jdbc_user => "user"
        jdbc_password => "pass"
        schedule => "5 * * * *"
        statement => "SELECT id, date, content, status from test_table WHERE date > :sql_last_start"
      }
    }
    

    此外,如果碰巧同一条记录从您的数据库加载了两次,并且您不希望在您的 ES 服务器中创建副本,您还可以指定使用记录 ID 作为您的 elasticsearch 中的文档 ID输出,这样文档将在 ES 中更新而不是重复。

    output {
    
        if [type] == "A" {
            elasticsearch {
                host => "localhost"
                protocol => http
                index => "logstash-servera-%{+YYYY.MM.dd}"
                document_id => "%{id}"       <--- same id as in DB
            }    
        }
        if [type] == "B" {
            elasticsearch {
                host => "localhost"
                protocol => http
                index => "logstash-serverb-%{+YYYY.MM.dd}"
                document_id => "%{id}"       <--- same id as in DB
            }    
        }
    
      stdout { codec => rubydebug }
    }
    

    【讨论】:

    • 我试过用这个。我第一次使用 WHERE 日期 > '2015-09-09'。从昨天到今天上午 11:30,logstash 获取了所有数据。我第二次跑的时候(11:50),Logstash 没有得到任何数据。是的,数据库中有更多数据。
    • 日志中唯一出现的是:Logstash启动完成Logstash关闭完成
    • 你能用--debug 运行logstash 并提供更多输出吗?
    • 我会在这个问题上添加这个输出。
    • 在您的输出中,我看到sql_last_start2015-09-10 18:48:00 UTC。不确定是否是问题所在,但这并不是真正接近上午 11:50。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2019-01-08
    • 1970-01-01
    • 1970-01-01
    • 2021-12-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多