如何实现MySQL到Logstash的数据更新自动化?
Hey there! Let’s work through this webhook-based sync problem you’re having with Logstash and MySQL. I get why relying on scheduled runs or restarting Logstash isn’t ideal—you want near-real-time sync without unnecessary overhead. Let’s break down two approaches to solve this: a webhook-triggered Logstash pipeline, and a more robust CDC (Change Data Capture) alternative that’s better suited for production.
Option 1: Webhook-Triggered Logstash Sync
The core idea here is to set up Logstash to listen for HTTP requests (your webhook) and run a targeted JDBC sync when it gets a trigger. Here’s how to configure it properly:
Step 1: Logstash Pipeline Configuration
Create a pipeline that uses the http input to listen for triggers, then uses jdbc_streaming to pull only new/updated data from MySQL. Replace placeholders with your actual values:
input { http { port => 8080 host => "0.0.0.0" # Add basic auth to block unauthorized triggers user => "logstash_webhook_user" password => "your_strong_password" } } filter { # Parse the JSON payload sent by the webhook json { source => "message" } # Fetch only changed data from MySQL using dynamic parameters from the webhook jdbc_streaming { jdbc_connection_string => "jdbc:mysql://localhost:3306/testdb" jdbc_user => "root" jdbc_password => "ankit" jdbc_driver_library => "/home/ankit/Downloads/mysql-connector-java-5.1.38.jar" jdbc_driver_class => "com.mysql.jdbc.Driver" # Query uses table name and last sync time from the webhook payload statement => "SELECT * FROM %{[table]} WHERE created_at > :last_sync OR updated_at > :last_sync" parameters => { "last_sync" => "%{[last_sync]}" } target => "synced_records" # Clean up unnecessary fields from the original HTTP request remove_field => ["message", "headers", "host", "@version"] } # Split the array of fetched records into individual Logstash events split { field => "synced_records" target => "" } # Remove the now-empty synced_records field mutate { remove_field => ["synced_records"] } } output { elasticsearch { hosts => ["http://localhost:9200"] # Use the table name from the payload to target the correct ES index index => "%{[table]}" # Use your table's primary key as the ES document ID to avoid duplicate entries document_id => "%{[id]}" } # Optional: Debug output to console for testing stdout { codec => rubydebug } }
Step 2: Trigger the Webhook from MySQL
MySQL doesn’t have built-in HTTP support, so we’ll use a trigger with a stored procedure that calls curl via a UDF (User-Defined Function). Note: This has security risks—only use it in a controlled environment, or skip to Option 2 for a safer approach.
- Install the
lib_mysqludf_sysUDF (this lets MySQL execute system commands). - Create a stored procedure to send the webhook request:
DELIMITER // CREATE PROCEDURE send_sync_trigger(IN table_name VARCHAR(255), IN last_sync DATETIME) BEGIN SET @curl_command = CONCAT( 'curl -X POST http://your_logstash_host:8080 -u logstash_webhook_user:your_strong_password -H "Content-Type: application/json" -d ''{"table": "', table_name, '", "last_sync": "', DATE_FORMAT(last_sync, '%Y-%m-%d %H:%i:%s'), '"}''' ); SELECT sys_exec(@curl_command); END // DELIMITER ; - Add triggers to your MySQL tables to call this procedure on insert/update:
DELIMITER // -- Triggers for ghijkl table CREATE TRIGGER after_ghijkl_insert AFTER INSERT ON ghijkl FOR EACH ROW BEGIN -- Sync records from the last minute (adjust the interval as needed) CALL send_sync_trigger('ghijkl', NOW() - INTERVAL 1 MINUTE); END // CREATE TRIGGER after_ghijkl_update AFTER UPDATE ON ghijkl FOR EACH ROW BEGIN CALL send_sync_trigger('ghijkl', NOW() - INTERVAL 1 MINUTE); END // -- Repeat for abcdef table CREATE TRIGGER after_abcdef_insert AFTER INSERT ON abcdef FOR EACH ROW BEGIN CALL send_sync_trigger('abcdef', NOW() - INTERVAL 1 MINUTE); END // CREATE TRIGGER after_abcdef_update AFTER UPDATE ON abcdef FOR EACH ROW BEGIN CALL send_sync_trigger('abcdef', NOW() - INTERVAL 1 MINUTE); END // DELIMITER ;
Option 2: Robust CDC with Debezium (Recommended for Production)
If you need reliable, low-overhead sync without the security risks of MySQL UDFs, Debezium is the way to go. It captures MySQL binlog events (row-level changes) and sends them directly to a message broker like Kafka, which Logstash can then consume to sync with Elasticsearch.
Key Benefits:
- No triggers or webhooks needed—Debezium listens directly to MySQL’s binlog.
- No missed changes, even if Logstash is temporarily down.
- Better performance for high-volume databases.
Quick Setup Steps:
- Enable MySQL binlog: Configure your MySQL server to use row-based replication (edit
my.cnf/my.ini):log_bin = mysql-bin binlog_format = ROW server_id = 1 - Deploy Debezium Connector: Use Kafka Connect to run the Debezium MySQL Connector. It will capture changes from MySQL and send them to Kafka topics.
- Configure Logstash to consume from Kafka: Add a Kafka input to your Logstash pipeline to read the change events and sync them to Elasticsearch.
Important Notes
- Avoid Duplicates: Make sure your MySQL tables have
created_atandupdated_attimestamp fields (auto-set on insert/update) so you only sync new/changed data. - Webhook Reliability: If you stick with webhooks, add a retry mechanism (like a MySQL log table for failed triggers) since HTTP requests can fail if Logstash is down.
- Security: For the webhook approach, always use HTTPS and strong authentication to prevent abuse. The
sys_execUDF is powerful but risky—restrict MySQL user permissions if you use it.
内容的提问来源于stack exchange,提问作者ankitkhandelwal185

