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类型字段。
问题出在时间精度丢失和时区处理两个关键点上:
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'的条件,被错误同步。时区转换的隐式问题
虽然配置了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

