如何用Logstash连接MongoDB与Elasticsearch?同步问题及替代方案咨询
Hey there! Let's break down your questions and walk through practical solutions, just like we do on Stack Overflow.
Setting up this sync is straightforward with the right configuration. Here's a step-by-step guide:
Step 1: Install the logstash-input-mongodb Plugin
First, install the official MongoDB input plugin via Logstash's plugin manager. Run this command in your Logstash directory:
# For Linux/macOS bin/logstash-plugin install logstash-input-mongodb # For Windows bin\logstash-plugin.bat install logstash-input-mongodb
Step 2: Write a Logstash Pipeline Configuration
Create a pipeline file (e.g., mongo-to-es.conf) with the following structure. I'll explain each section so you can tweak it to your needs:
input { mongodb { uri => "mongodb://localhost:27017/your_target_db" # Replace with your MongoDB URI collection => "your_target_collection" # Name of the collection to sync batch_size => 1000 # Number of docs to fetch per batch # For incremental sync (critical for updates later): # Option 1: Use a timestamp field (e.g., `updated_at`) in your MongoDB docs initial_timestamp_field => "updated_at" # Option 2: Use MongoDB Change Streams (requires MongoDB 3.6+ and plugin v4.0+) use_change_stream => true } } filter { # Optional: Clean up or transform data before sending to ES mutate { remove_field => ["@version", "@timestamp"] # Remove default Logstash fields rename => { "_id" => "mongo_id" } # Avoid conflicts with ES's internal _id } } output { elasticsearch { hosts => ["http://localhost:9200"] # Your Elasticsearch endpoint index => "your_es_index_name" # Name of the ES index to sync to document_id => "%{mongo_id}" # Use MongoDB's _id as ES document ID to prevent duplicates } stdout { codec => rubydebug } # Optional: Print sync details to console for debugging }
Step 3: Start Logstash with Your Pipeline
Run Logstash using your configuration file:
# Linux/macOS bin/logstash -f mongo-to-es.conf # Windows bin\logstash.bat -f mongo-to-es.conf
If you're using logstash-input-mongodb but updates aren't syncing, here's why and what you can do:
Why Updates Aren't Syncing
The most common culprit is missing incremental sync configuration or an outdated plugin version. Older versions of logstash-input-mongodb don't support MongoDB Change Streams, so they only sync new documents, not updates. Even if you use a timestamp field, you need to ensure the field is updated every time a document changes.
Is it a plugin issue? Sometimes yes—this plugin isn't the most robust for real-time update sync. But before switching, try these quick fixes:
- Update the plugin to the latest version:
bin/logstash-plugin update logstash-input-mongodb - Verify
use_change_stream => trueis set (if your MongoDB version supports it) - Double-check that your
initial_timestamp_fieldexists and is updated on every document change
Alternative Tools/Plugins for Reliable Sync
If the Logstash plugin still isn't cutting it, here are better options tailored for MongoDB-to-ES sync:
1. MongoDB Connector for Elasticsearch (Official Tool)
This is the most reliable option—it's built and maintained by MongoDB specifically for syncing to Elasticsearch. It supports:
- Full initial sync of existing data
- Real-time incremental sync via Change Streams
- Automatic schema mapping between MongoDB and ES
- Filtering and transformation rules
Setup is simple: Download the connector, create a config file specifying your MongoDB URI, ES endpoint, and sync rules, then start the connector. No Logstash required!
2. Debezium (CDC Tool)
Debezium is a powerful Change Data Capture (CDC) tool that works with MongoDB (and many other databases). It captures every change in your MongoDB collection, streams events to Kafka, and then syncs them to Elasticsearch via Kafka Connect. This is ideal for large-scale, high-throughput environments where you need guaranteed delivery of updates.
3. Custom Script + Logstash HTTP Input
If you want to stick with Logstash but avoid the MongoDB plugin, write a simple script (Node.js, Python, etc.) to listen to MongoDB Change Streams, then send the updated documents to Logstash's HTTP input. Here's a quick example of the Logstash input config for this:
input { http { port => 8080 codec => json } }
Your script would send POST requests to http://localhost:8080 with the updated document data, and Logstash would forward it to Elasticsearch as usual.
内容的提问来源于stack exchange,提问作者Edit Axha

