Logstashを活用したMySQLからEasysearchへのデータ同期構築方法

ELKスタックの一環であるLogstashを利用すれば、リレーショナルデータベースから検索エンジンへ効率的にデータ移行が可能です。本稿では、MySQLからEasysearchへの増分同期パイプライン設定について解説します。

必須要件

  • 主キーの存在: テーブルには一意な識別子(例: pk_id)が必須です。これにより、ソースDBレコードとEasysearchのドキュメントIDを厳密に一致させ、更新時の上書き制御を可能にします。
  • 更新時刻カラム: レコード作成および更新時のタイムスタンプを格納する列が必要不可欠です。この値を用いてLogstashが差分検知を行うため、増分取り込みの基盤となります。

これらの条件を満たせば、定期ジョブによる新規および変更データの自動連携が実現します。

環境バージョン

  • MySQL: 5.7
  • Logstash: 7.10.2
  • Easysearch: 1.5.0

データベース準備

テーブル定義

CREATE DATABASE IF NOT EXISTS sync_source_db;
USE sync_source_db;

DROP TABLE IF EXISTS target_records;
CREATE TABLE target_records (
  pk_id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT,
  PRIMARY KEY (pk_id),
  UNIQUE KEY uk_pk_id (pk_id),
  customer_alias VARCHAR(64) NOT NULL,
  updated_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
  created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP
);

各項目の意味:

  • pk_id: ドキュメントIDとして使用される主キー兼ユニークキー。
  • updated_at: 挿入・更新イベントを検出するための基準時刻。
  • customer_alias: 転送対象の実務データ。
  • created_at: データ生成日を記録(任意項目)。

初期データ投入

INSERT INTO target_records (pk_id, customer_alias) VALUES (101, 'Sample_User_A');
INSERT INTO target_records (pk_id, customer_alias) VALUES (102, 'Sample_User_B');
INSERT INTO target_records (pk_id, customer_alias) VALUES (103, 'Sample_User_C');

Logstash パイプライン設定

以下は pipeline.conf の実装例です。JDBC Plugin を駆使して、毎度全体のスキャンではなく、updated_at に基づく増分抽出を実現しています。

input {
  jdbc {
    jdbc_driver_library => "./libs/mysql-connector-j-8.1.0.jar"
    jdbc_driver_class => "com.mysql.cj.jdbc.Driver"
    jdbc_connection_string => "jdbc:mysql://192.168.56.3:3306/sync_source_db"
    jdbc_user => "root"
    jdbc_password => "secure_pass"
    jdbc_paging_enabled => true
    
    # 増分同期のキーとなるUNIX時間戳
    tracking_column => "unix_update_ts"
    use_column_value => true
    tracking_column_type => "numeric"
    last_run_metadata_path => ".logstash_meta/target_records_state.yml"
    schedule => "*/5 * * * * *"
    
    statement => "SELECT *, UNIX_TIMESTAMP(updated_at) AS unix_update_ts FROM target_records WHERE UNIX_TIMESTAMP(updated_at) > :sql_last_value AND updated_at < NOW() ORDER BY updated_at ASC"
  }

  jdbc {
    jdbc_driver_library => "./libs/mysql-connector-j-8.1.0.jar"
    jdbc_driver_class => "com.mysql.cj.jdbc.Driver"
    jdbc_connection_string => "jdbc:mysql://192.168.56.3:3306/sync_source_db"
    jdbc_user => "root"
    jdbc_password => "secure_pass"
    schedule => "*/5 * * * * *"
    
    # インデックス状態モニタリング用クエリ
    statement => "SELECT COUNT(*) AS doc_count, 'target_records' AS source_tbl FROM target_records"
  }
}

filter {
  # 通常データ処理フロー
  if ![source_tbl] {
    mutate {
      copy => { "pk_id" => "[@metadata][dest_doc_id]" }
      rename => { "customer_alias" => "alias" }
      remove_field => ["@version", "unix_update_ts", "@timestamp"]
      add_field => { "[@metadata][target_index]" => "synced_target_records" }
    }
  # モニタリングデータ処理フロー
  } else {
    mutate {
      add_field => { "[@metadata][target_index]" => "db_monitor_counts" }
      remove_field => ["@version"]
    }
    uuid {
      target => "[@metadata][dest_doc_id]"
      overwrite => true
    }
  }
}

output {
  elasticsearch {
    hosts => ["https://localhost:9200"]
    user => "admin"
    password => "a1b2c3d4e5f6g7h8"
    ssl_certificate_verification => false
    index => "%{[@metadata][target_index]}"
    manage_template => false
    document_id => "%{[@metadata][dest_doc_id]}"
  }
  # 開発時デバッグ用
  # stdout { codec => rubydebug { metadata => true } }
}

上記設定により、5秒ごとに target_records の変更分が synced_target_records インデックスへ流れます。同時に統計用クエリの結果も db_monitor_counts に蓄積され、運用監視に利用可能です。

起動と動作確認

パイプラインファイルを読み込んでプロセスを開始します。

$ ./bin/logstash -f pipeline.conf

起動直後に3件の初期データが正常に検索エンジン側に反映されます。以降の変更操作を確認してみましょう。

新規登録テスト:

INSERT INTO target_records (pk_id, customer_alias) VALUES (104, 'New_Record_X');

Easysearch側でも即座に検索可能になっています。

更新テスト:

UPDATE target_records SET customer_alias = 'Updated_Variant_01' WHERE pk_id=101;

Logstashの増分クエリが更新時刻を検出し、既存ドキュメントを上書きするため、Easysearch側のレコード内容も即時更新されています。

削除処理の対応方針

JDBC Plugin は標準機能として「物理DELETE」をElasticsearch互換エンジン側へ直接送信する仕組みを持っていません。以下のいずれかのアーキテクチャを採用する必要があります。

  1. 論理削除(Soft Delete): テーブルに is_deleted TINYINT DEFAULT 0 などのフラグを追加します。削除時はフラグを立てるのみとし、検索クエリで除外します。バックグラウンドバッチ処理にて一定期間経過した論理削除レコードを実際に削除またはアーカイブします。
  2. アプリケーション層連携: DB側でのDELETEコマンド実行後、同一トランザクションまたはメッセージキュー経由で専用APIコールを発行し、Easysearch側でも該当 document_id を削除する処理を実装します。

同期状況の可視化

転送が定着したら、データの不整合や遅延を防ぐために監視ダッシュボードの構築が推奨されます。INFINI Console 等の管理コンソールを活用すれば、ソーステーブルの行数、ターゲットインデックスの文档数、前日比の変化率などをワンページでグラフ化でき、運用品質を担保しやすくなります。

タグ: Logstash MySQL JDBCプラグイン Easysearch 増分データ同期

8月15日 21:20 投稿