Logstash同步PostgreSQL与Elasticsearch时:sql_last_value更新异常排查
PostgreSQL与Elasticsearch数据同步:Logstash JDBC输入时间类型不匹配问题解决
问题场景
使用Logstash的JDBC输入插件同步PostgreSQL的products表到Elasticsearch时,触发类型不匹配错误;尝试通过SQL函数转换查询条件后,sql_last_value始终为0,无法实现增量同步。
原配置与错误日志
初始Logstash配置
input { jdbc { jdbc_connection_string => "jdbc:postgresql://postgres:5432/onlinestore" jdbc_user => "postgres" jdbc_password => "1234" jdbc_driver_class => "org.postgresql.Driver" jdbc_driver_library => "/usr/share/logstash/libs/postgresql-42.7.3.jar" statement => " SELECT * FROM products WHERE updated_date > :sql_last_value ORDER BY updated_date ASC LIMIT 1000 " use_column_value => true tracking_column => "updated_date" last_run_metadata_path => "/usr/share/logstash/.logstash_jdbc_last_run" schedule => "* * * * *" # clean_run => true } } filter { ruby { code => " event.set('sql_last_value_log', event.get('sql_last_value')) event.set('tracking_column_log', event.get('updated_date')) " } mutate { add_field => { "timestamp" => "%{@timestamp}" } } } output { elasticsearch{ hosts => ["http://elasticsearch:9200"] index => "products" document_id => "%{id}" } }
错误日志
/usr/share/logstash/vendor/bundle/jruby/2.5.0/gems/rufus-scheduler-3.0.9/lib/rufus/scheduler/cronline.rb:77: warning: constant ::Fixnum is deprecated [2024-07-07T10:06:04,030][ERROR][logstash.inputs.jdbc ][main][c96da82d0fe625ff4baf1477649ffdc670902119c45fc017f20bc336f857bcfe] Java::OrgPostgresqlUtil::PSQLException: ERROR: operator does not exist: timestamp without time zone > integer Hint: No operator matches the given name and argument types. You might need to add explicit type casts. Position: 82: SELECT * FROM products WHERE updated_date > 0 ORDER BY updated_date ASC LIMIT 1000 [2024-07-07T10:06:04,117][WARN ][logstash.inputs.jdbc ][main][c96da82d0fe625ff4baf1477649ffdc670902119c45fc017f20bc336f857bcfe] Exception when executing JDBC query {:exception=>Sequel::DatabaseError, :message=>"Java::OrgPostgresqlUtil::PSQLException: ERROR: operator does not exist: timestamp without time zone > integer\n Hint: No operator matches the given name and argument types. You might need to add explicit type casts.\n Position: 82", :cause=>"org.postgresql.util.PSQLException: ERROR: operator does not exist: timestamp without time zone > integer\n Hint: No operator matches the given name and argument types. You might need to add explicit type casts.\n Position: 82"}
问题根源
- 首次运行类型不匹配:当
last_run_metadata_path文件不存在或为空时,Logstash默认将sql_last_value设为整数0,而products表的updated_date是timestamp without time zone类型,直接比较触发类型错误。 - 修改后无法更新
sql_last_value:使用COALESCE+NULLIF转换查询条件后,虽然查询能执行,但Logstash内部仍将sql_last_value以整数类型存储,无法识别并更新为updated_date的timestamp值,导致增量同步失效。
解决方案
方案1:手动初始化元数据文件(推荐)
- 创建并初始化
last_run_metadata_path指定的文件:echo "1970-01-01 00:00:00" > /usr/share/logstash/.logstash_jdbc_last_run - 恢复原JDBC查询语句(无需额外类型转换):
statement => " SELECT * FROM products WHERE updated_date > :sql_last_value ORDER BY updated_date ASC LIMIT 1000 " - 确保
use_column_value和tracking_column配置正确:use_column_value => true tracking_column => "updated_date" - 重启Logstash,后续
sql_last_value会自动更新为每次查询到的最大updated_date值。
方案2:指定跟踪列类型
在JDBC输入配置中添加tracking_column_type参数,明确告诉Logstash跟踪列的类型为timestamp,自动处理初始值的类型转换:
input { jdbc { # 其他原有配置不变 tracking_column_type => "timestamp" statement => " SELECT * FROM products WHERE updated_date > :sql_last_value ORDER BY updated_date ASC LIMIT 1000 " use_column_value => true tracking_column => "updated_date" last_run_metadata_path => "/usr/share/logstash/.logstash_jdbc_last_run" } }
该参数会让Logstash将sql_last_value以timestamp类型处理,首次运行时自动将默认的0转换为1970-01-01 00:00:00,同时后续正常更新sql_last_value。
验证
- 重启Logstash后,查看
/usr/share/logstash/.logstash_jdbc_last_run文件,内容应更新为本次同步的最大updated_date值。 - 检查Logstash日志,无类型不匹配错误,Elasticsearch的
products索引可正常同步增量数据。
内容的提问来源于stack exchange,提问作者Joker 15
相关产品推荐
相关产品推荐

