You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Logstash sql_last_value大于条件异常:同步含等于值数据排查

问题

搭建了由RDS(Aurora DB)、Logstash和AWS OpenSearch组成的数据管道,为保证OpenSearch索引的数据一致性,通过Logstash JDBC输入配置使用sql_last_value与updated_at的大于条件过滤重复数据,但实际同步到OpenSearch的数据中包含了updated_at等于sql_last_value的数据。

配置文件

input{
  jdbc {
    jdbc_driver_library => "/home/ubuntu/logstash-7.16.2/bin/mysql-connector-java-8.0.27.jar"
    jdbc_driver_class => "com.mysql.cj.jdbc.Driver"
    jdbc_connection_string => "jdbc:mysql://~~~~?useSSL=false"
    jdbc_user => "root"
    jdbc_password => "~~~~~~"
    jdbc_paging_enabled => true
    tracking_column => "updated_at"
    use_column_value => true
    record_last_run => true
    tracking_column_type => "timestamp"
    schedule => "*/10 * * * * *"
    statement => "select * from my_table where updated_at > :sql_last_value order by updated_at ASC"
    jdbc_default_timezone => "Asia/Seoul"
  }
}

output {
  opensearch{
    hosts => "https://~~~~~:443"
    user => "admin"
    password => "~~~~~"
    index => "index"
    ecs_compatibility => disabled
    ssl_certificate_verification => false
  }
}

Logstash查询日志

[2022-12-29T16:48:40,299][INFO ][logstash.inputs.jdbc     ][main][0bb20d034a10be3c1a48635cda2cc7dfcb97e29fb63940352f5380ec253dfe48] (0.001157s) SELECT version()
[2022-12-29T16:48:40,302][INFO ][logstash.inputs.jdbc     ][main][0bb20d034a10be3c1a48635cda2cc7dfcb97e29fb63940352f5380ec253dfe48] (0.001163s) SELECT version()
[2022-12-29T16:48:40,306][INFO ][logstash.inputs.jdbc     ][main][0bb20d034a10be3c1a48635cda2cc7dfcb97e29fb63940352f5380ec253dfe48] (0.001221s) SELECT count(*) AS `count` FROM (select * from my_table where updated_at > '2022-12-28 22:31:05' order by updated_at ASC) AS `t1` LIMIT 1
[2022-12-29T16:48:40,309][INFO ][logstash.inputs.jdbc     ][main][0bb20d034a10be3c1a48635cda2cc7dfcb97e29fb63940352f5380ec253dfe48] (0.001166s) SELECT * FROM (select * from my_table where updated_at > '2022-12-28 22:31:05' order by updated_at ASC) AS `t1` LIMIT 100000 OFFSET 0
[2022-12-29T16:48:50,172][INFO ][logstash.inputs.jdbc     ][main][0bb20d034a10be3c1a48635cda2cc7dfcb97e29fb63940352f5380ec253dfe48] (0.001303s) SELECT version()
[2022-12-29T16:48:50,174][INFO ][logstash.inputs.jdbc     ][main][0bb20d034a10be3c1a48635cda2cc7dfcb97e29fb63940352f5380ec253dfe48] (0.001152s) SELECT version()
[2022-12-29T16:48:50,178][INFO ][logstash.inputs.jdbc     ][main][0bb20d034a10be3c1a48635cda2cc7dfcb97e29fb63940352f5380ec253dfe48] (0.001382s) SELECT count(*) AS `count` FROM (select * from my_table where updated_at > '2022-12-28 22:31:05' order by updated_at ASC) AS `t1` LIMIT 1
[2022-12-29T16:48:50,182][INFO ][logstash.inputs.jdbc     ][main][0bb20d034a10be3c1a48635cda2cc7dfcb97e29fb63940352f5380ec253dfe48] (0.001153s) SELECT * FROM (select * from my_table where updated_at > '2022-12-28 22:31:05' order by updated_at ASC) AS `t1` LIMIT 100000 OFFSET 0

异常场景示例

sql_last_value = 2022.12.08 12:12:12

first data : 2022.12.08 12:12:11
second data : 2022.12.08 12:12:12
third data : 2022.12.08 12:12:13

按预期应该只同步第三条数据,但实际第二条和第三条都被查询出来,且下一次查询时sql_last_value被更新为2022.12.08 12:12:13。其中updated_at列是由NestJS的Sequelize模块生成的timestamp类型字段。

分析与解决

问题出在时间精度丢失和时区处理两个关键点上:

  1. timestamp精度不匹配
    Aurora MySQL的timestamp类型默认支持毫秒级精度,但Logstash在存储sql_last_value时,会将时间截断到秒级。而Sequelize生成的updated_at可能包含毫秒部分(比如2022-12-08 12:12:12.123),当Logstash将sql_last_value存储为2022-12-08 12:12:12时,数据库中updated_at等于2022-12-08 12:12:12.xxx的记录会满足updated_at > '2022-12-08 12:12:12'的条件,被错误同步。

  2. 时区转换的隐式问题
    虽然配置了jdbc_default_timezone => "Asia/Seoul",但Logstash在读取数据库的timestamp字段并更新sql_last_value时,可能存在时区转换的精度损失,导致存储的sql_last_value比实际记录的updated_at少了部分精度,最终使得等于原秒级时间的记录被误判为大于条件。

修复方案

  • 调整SQL语句,补充精度判断
    将查询语句改为结合主键的复合判断,避免仅靠时间精度区分:
statement => "select * from my_table where updated_at > :sql_last_value or (updated_at = :sql_last_value and id > :sql_last_value_id) order by updated_at ASC, id ASC"

同时需要新增tracking_column配置,或者改用id作为辅助跟踪字段,确保唯一判断。

  • 提升Logstash的跟踪字段精度
    将tracking_column_type改为timestamp_usec,确保Logstash存储sql_last_value时保留毫秒/微秒精度:
tracking_column_type => "timestamp_usec"
  • 统一数据库与Logstash的时间精度
    在Sequelize中配置updated_at的精度为秒级,或者确保数据库的timestamp字段与Logstash的存储精度一致,避免因精度差异导致的比较错误。

内容的提问来源于stack exchange,提问作者Jun

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.06 20:05:13