You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

咨询:如何基于Logstash将多个MySQL表迁移至同一Elasticsearch索引?

Got it, let's walk through how to get your multiple MySQL tables into a single Elasticsearch index using Logstash. Your initial input setup is solid—here's how to expand it for multi-table support, with clear, actionable steps:

Complete Logstash Config for Multi-Table Migration to One ES Index

1. Input Section: Fetch Data from Multiple Tables

We'll use separate JDBC input blocks for each table (flexible for handling unique table needs) and add a source identifier to track where each record comes from. Don't forget to specify the MySQL driver class (easy to miss!) and use pagination to avoid memory overload.

input {
  # First table: e.g., table_movies
  jdbc {
    jdbc_connection_string => "jdbc:mysql://localhost:3306/cinestop?useUnicode=yes&characterEncoding=UTF-8&serverTimezone=UTC"
    jdbc_user => "root"
    jdbc_password => "abc"
    jdbc_driver_library => "/home/limafoxtrottango/Downloads/mysql-connector-java-8.0.30.jar" # Ensure full, valid path
    jdbc_driver_class => "com.mysql.cj.jdbc.Driver"
    jdbc_paging_enabled => true
    jdbc_page_size => 1000
    # Pull all data from the table; add WHERE clauses here for incremental sync later
    statement => "SELECT * FROM table_movies"
    # Add a field to identify which table this record came from
    add_field => {"source_table" => "table_movies"}
    # Optional incremental sync setup (uncomment if you need ongoing updates):
    # use_column_value => true
    # tracking_column => "update_time"
    # tracking_column_type => "timestamp"
    # schedule => "* * * * *"
  }

  # Second table: e.g., table_users
  jdbc {
    jdbc_connection_string => "jdbc:mysql://localhost:3306/cinestop?useUnicode=yes&characterEncoding=UTF-8&serverTimezone=UTC"
    jdbc_user => "root"
    jdbc_password => "abc"
    jdbc_driver_library => "/home/limafoxtrottango/Downloads/mysql-connector-java-8.0.30.jar"
    jdbc_driver_class => "com.mysql.cj.jdbc.Driver"
    jdbc_paging_enabled => true
    jdbc_page_size => 1000
    statement => "SELECT * FROM table_users"
    add_field => {"source_table" => "table_users"}
  }

  # Add more JDBC blocks for additional tables here...
}

2. Filter Section: Unify & Clean Your Data

Since your tables share common fields, use this section to standardize data: fix conflicting field names, remove irrelevant data, and correct data types (like converting string dates to proper timestamps).

filter {
  # Example: Unify conflicting field names (e.g., one table uses user_id, another uses id)
  mutate {
    rename => {
      "user_id" => "id" # Map to a single consistent field name
    }
    # Remove fields you don't need in Elasticsearch
    remove_field => ["@version"]
  }

  # Example: Convert string date fields to Elasticsearch-compatible timestamps
  date {
    match => ["create_time", "yyyy-MM-dd HH:mm:ss"]
    target => "@timestamp"
  }
}

3. Output Section: Send All Data to a Single Index

Route all cleaned records to your unified Elasticsearch index. Use a unique document ID (combining source table and primary key) to avoid overwriting records.

output {
  elasticsearch {
    hosts => ["http://localhost:9200"] # Your Elasticsearch server address
    index => "cinestop_unified" # Name your single unified index here
    # Create a unique ID for each document to prevent duplicates
    document_id => "%{source_table}_%{id}"
  }

  # Optional: Print output to console for debugging
  stdout {
    codec => rubydebug
  }
}

Key Tips for Your Small Dataset (<3G)

  • Pagination: Enabled here as a best practice, but you can disable it if you're confident your server can handle loading all data into memory at once.
  • Schema Consistency: Ensure common fields have matching data types across tables (e.g., create_time should be a timestamp everywhere). Use the mutate/convert filter to fix mismatches.
  • Alternative Single Input: If your tables have nearly identical schemas, you can simplify with a single JDBC input using UNION ALL:
    SELECT *, 'table_movies' as source_table FROM table_movies 
    UNION ALL 
    SELECT *, 'table_users' as source_table FROM table_users
    
    This cuts down on config duplication but only works if tables are structurally compatible.

Once your config is ready, save it as cinestop_migration.conf and run Logstash with:

bin/logstash -f cinestop_migration.conf

内容的提问来源于stack exchange,提问作者Prashant Pandey

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.20 10:40:28