Logstash exec插件无法导入Twitter Decahose流至Elastic,求替代方案
Hey there, sorry to hear the exec plugin isn't working out for your Twitter Decahose import to Elastic—let's break down why that's not the right tool for the job, then walk through far better Logstash-based solutions tailored for streaming Twitter data.
First, quick context: The exec plugin is designed for running one-off or scheduled external commands, not handling persistent, real-time data streams like Decahose. It can't maintain long-lived connections, handle automatic reconnections if the stream drops, or natively manage Twitter's authentication requirements—so it's no surprise it's failing here.
Solution 1: Use the Official X (Twitter) Input Plugin
This is the most straightforward approach, as the plugin is purpose-built for streaming Twitter data (including Enterprise-level feeds like Decahose). It handles authentication, persistent connections, automatic retries, and parses the stream natively.
Step 1: Install the Plugin
First, install the plugin via Logstash's CLI:
bin/logstash-plugin install logstash-input-twitter
Step 2: Sample Pipeline Configuration
Here's a basic config tailored for Decahose (adjust the endpoint and credentials to match your Enterprise account):
input { twitter { # OAuth credentials for your Enterprise account consumer_key => "YOUR_CONSUMER_KEY" consumer_secret => "YOUR_CONSUMER_SECRET" access_token => "YOUR_ACCESS_TOKEN" access_token_secret => "YOUR_ACCESS_TOKEN_SECRET" # Decahose stream endpoint (replace with your account-specific URL) endpoint => "https://gnip-api.twitter.com/stream/powertrack/accounts/YOUR_ACCOUNT_NAME/publishers/twitter/streams/track/dev.json" # Optional: Add track keywords if you need to filter the stream keywords => ["your-target-terms"] } } filter { # Parse the raw Twitter JSON into structured fields json { source => "message" } # Clean up unnecessary fields and add custom metadata mutate { remove_field => ["message", "@version"] add_field => { "ingested_via" => "logstash-twitter-plugin" } } # Ensure timestamps are properly parsed (Twitter's created_at field) date { match => ["created_at", "EEE MMM dd HH:mm:ss ZZZ yyyy"] target => "@timestamp" } } output { # Send to Elasticsearch with a daily rolling index elasticsearch { hosts => ["http://localhost:9200"] index => "twitter-decahose-%{+YYYY.MM.dd}" # Use Twitter's unique ID as the document ID to avoid duplicates document_id => "%{id_str}" } # Optional: Print to stdout for debugging stdout { codec => rubydebug } }
Solution 2: HTTP Poller Input (For Custom/API v2 Use Cases)
If the official plugin doesn't support your specific Decahose setup (e.g., newer API v2 endpoints), use Logstash's http_poller input to connect directly to the Enterprise stream. This gives you more control over headers, authentication, and stream handling.
Sample Configuration
input { http_poller { urls => { twitter_decahose => { url => "https://api.twitter.com/2/tweets/search/stream" # Adjust to your Decahose endpoint headers => { # Use Bearer Token for API v2 authentication "Authorization" => "Bearer YOUR_ENTERPRISE_BEARER_TOKEN" "Content-Type" => "application/json" } # Maintain a persistent connection for streaming keepalive => true request_timeout => 3600 # Keep connection open for 1 hour (adjust as needed) } } # For streaming, the cron schedule is less critical—this ensures the input runs continuously schedule => { cron => "* * * * *" } # Decahose returns newline-delimited JSON, so use this codec codec => json_lines } } filter { # Add your custom parsing/transform logic here mutate { rename => { "data" => "tweet" } # Adjust based on API v2's response structure add_field => { "ingest_time" => "%{@timestamp}" } } } output { elasticsearch { hosts => ["http://localhost:9200"] index => "twitter-decahose-%{+YYYY.MM.dd}" document_id => "%{tweet.id}" } }
Pro Tips for Success
- Debug First: Start with the
stdoutoutput to verify data is being received and parsed correctly before connecting to Elasticsearch. - Check Logs: Enable debug logging in Logstash (
--log.level debug) to troubleshoot connection issues or authentication errors. - Optimize Performance: For high-volume Decahose streams, adjust
pipeline.batch.sizeandpipeline.batch.delayinlogstash.ymlto balance throughput and latency. - Index Templates: Create an Elasticsearch index template to define field types (e.g.,
keywordfor hashtags,datefor timestamps) to avoid mapping conflicts.
内容的提问来源于stack exchange,提问作者Arun

