如何在Logstash与Elasticsearch中实现数据库更新后仅删除特定行
Hey there! Let's walk through exactly how to set up your Oracle-to-ELK sync with the status-based logic you need—including using upsert for PENDING requests and deleting records when they flip to OK. We'll use Logstash (with its JDBC input plugin) since it's the standard tool for this kind of relational-database-to-ELK pipeline.
Step 1: Set Up Logstash JDBC Input (Pull Data from Oracle)
First, we need to configure Logstash to pull incremental updates from Oracle (using last_update to avoid reprocessing old data). This ensures we only grab records that have changed since the last sync.
Here's the input config snippet:
input { jdbc { jdbc_connection_string => "jdbc:oracle:thin:@//your-oracle-host:1521/your-sid" jdbc_user => "your_db_user" jdbc_password => "your_db_password" jdbc_driver_library => "/path/to/ojdbc8.jar" # Make sure you have the Oracle JDBC driver jdbc_driver_class => "oracle.jdbc.OracleDriver" schedule => "* * * * *" # Run every minute (adjust based on your real-time needs) statement => "SELECT request_number, request_status, last_update FROM your_request_table WHERE last_update > :sql_last_value" tracking_column => "last_update" use_column_value => true tracking_column_type => "timestamp" clean_run => false # Keep track of the last sync time between runs } }
Step 2: Filter Logic to Handle Status Rules
Next, we'll add a filter to tag each record with the action it needs in Elasticsearch:
- For
request_status = 1(PENDING): Mark for upsert (insert if new, update if existing) - For
request_status = 10(OK): Mark for deletion - Drop any other statuses (adjust if you need to handle others)
filter { # First, make sure request_status is treated as a number (avoid string comparisons) mutate { convert => { "request_status" => "integer" } } # Tag actions based on status if [request_status] == 10 { mutate { add_field => { "[@metadata][action]" => "delete" } } } elsif [request_status] == 1 { mutate { add_field => { "[@metadata][action]" => "upsert" } } } else { # Drop records with other statuses (customize if needed) drop {} } }
Step 3: Elasticsearch Output with Upsert/Delete Logic
Finally, configure the Elasticsearch output to act on the tags we set. We'll use request_number as the document ID (since it's your primary key) to ensure we target the correct record for upsert or delete.
output { elasticsearch { hosts => ["http://your-es-host:9200"] index => "pending_requests" # Name your index whatever makes sense document_id => "%{request_number}" # Use primary key to match records action => "%{[@metadata][action]}" # Use the action tag from the filter doc_as_upsert => true # For upserts, use the entire incoming document as the update/insert content upsert => '{}' # Required placeholder for doc_as_upsert } # Optional: Log to stdout for debugging stdout { codec => rubydebug } }
How This Works for Your Status Change Rules
- When a request stays PENDING (status 1): Logstash pulls the updated
last_updatevalue, runs an upsert, and updates the existing document in ELK (keeping it in the index). - When a request changes from PENDING to OK (status 10): Logstash pulls the updated record, tags it for deletion, and removes the matching document from ELK using the
request_numberID. - New PENDING requests: The upsert action inserts them into the index since no existing document with that
request_numberexists.
Key Notes to Avoid Issues
- Oracle JDBC Driver: Make sure you download the correct driver (ojdbc8 for Oracle 12c+) and place it in a path Logstash can access.
- Time Zones: Ensure Logstash and Oracle use the same time zone for
last_updateto avoid missing records. - Testing: Start with
schedule => "0 0 * * *"(daily) for testing, then adjust to higher frequency once you confirm it works. - Index Mapping: Pre-create your ELK index with proper mappings (e.g.,
request_statusas integer,last_updateas date) to avoid auto-mapping issues.
内容的提问来源于stack exchange,提问作者Manlio

