Logstash経由でMySQLとElasticsearchを統合:データ自動同期パイプラインの構築

MySQL tutorial - IT technology blog
MySQL tutorial - IT technology blog

現実の課題:データ肥大化によるMySQLの過負荷

products テーブルのレコード数が50万件に達すると、検索処理でパフォーマンス問題が発生し始めます。ユーザーは曖昧なキーワードや誤字を含む入力を行ったり、カテゴリ、価格帯、評価、在庫状況など4〜5つの条件を同時に絞り込んだりします。

もしMySQLで LIKE '%keyword%' に加え3〜4個の JOIN を含むクエリを実行し続けると、データベースサーバーのCPU使用率は瞬く間に100%に達します。レスポンスタイムは50ミリ秒から3〜5秒へと急増し、コネクションプール全体が枯渇して、注文作成や決済処理のAPIまでもが停止してしまいます。

なぜMySQLは全文検索(Full-Text Search)で力尽きるのか?

MySQLが高負荷な高度検索処理に向いていないのには、技術的に3つの理由があります。

  • B-Treeインデックスは先頭ワイルドカードに対応できない: MySQLは完全一致検索や前方一致検索(LIKE 'iphone%')には最適化されています。しかし、先頭に % を付与する中間・後方一致検索(LIKE '%iphone%')ではインデックスがまったく機能せず、フルテーブルスキャン(Full Table Scan)が強制されます。
  • 読み取りと書き込みのリソース競合: MySQLはOLTP(ACIDトランザクション)を目的として設計されています。データベースに関連度スコアリングの計算をさせながら同時に注文データの書き込みを行わせると、テーブルロックのリスクが高まり、ディスクI/Oが危険域に達します。
  • Elasticsearchの転置インデックス(Inverted Index)アーキテクチャ: Elasticsearchはテキストを個別のトークンに分割し、転置インデックスを構築します。キーワード検索時にはトークン一覧を参照して該当するDocument IDを取得するだけで済むため、数千万件規模のデータセットでも通常30ミリ秒未満で応答可能です。

代表的な3つのデータ同期アプローチ

システムの規模に応じて、以下のいずれかのアプローチを選択できます。

  1. アプリケーション層でのDual-Write(二重書き込み): MySQLへの追加・更新のたびに、バックエンドコードがRabbitMQやKafka経由でイベントを発行してElasticsearchを更新します。ほぼリアルタイムで同期できますが、コードロジックが複雑化し、ネットワーク障害時にデータ不整合が発生しやすくなります。
  2. DebeziumによるChange Data Capture(CDC): MySQLのbinlogを直接読み取るソリューションです。極めて高い精度でリアルタイム性を担保できますが、インフラコストと運用負荷が高く、大規模なマイクロサービスアーキテクチャで真価を発揮します。
  3. JDBC Inputプラグインを使用したLogstash: Logstashが中間ワーカーとして機能し、MySQLから新規または更新されたレコードを定期的にクエリしてElasticsearchへ転送します。

現実的な最適解:Logstash JDBCによる自動同期パイプライン

中小規模の大半のアプリケーション(レコード数が数百万件未満)において、Logstash JDBCは最もバランスの取れた選択肢です。バックエンドのコードを1行も変更する必要がなく、設定を一元管理でき、わずか30分程度でデプロイできます。

以下は、実際に構築するための4つのステップです。

ステップ1:MySQLデータベーステーブルの準備

Logstashが数百万行を再スキャンすることなく同期対象レコードを特定できるように、テーブルにはインデックス付きの updated_at カラムが必須となります。

CREATE DATABASE IF NOT EXISTS shop_db;
USE shop_db;

CREATE TABLE products (
    id INT AUTO_INCREMENT PRIMARY KEY,
    title VARCHAR(255) NOT NULL,
    description TEXT,
    price DECIMAL(10, 2) NOT NULL,
    is_deleted TINYINT(1) DEFAULT 0,
    created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
    updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
    INDEX idx_updated_at (updated_at)
);

-- サンプルデータの挿入
INSERT INTO products (title, description, price) VALUES
('ワイヤレスメカニカルキーボード', '75%レイアウト、Bluetoothおよび2.4GHz接続対応', 1500000),
('エルゴノミクスゲーミングマウス', '重量60g、26000 DPIセンサー搭載', 850000),
('27インチ 4K IPSモニター', 'sRGBカバー率100%、グラフィック制作用', 7200000);

ステップ2:MySQL Connector/J(JDBCドライバー)のダウンロード

LogstashがMySQLと通信するにはJDBCドライバーが必要です。Logstashサーバーに .jar ファイルをダウンロードします。

# ドライバー配置用ディレクトリの作成
sudo mkdir -p /etc/logstash/drivers
cd /etc/logstash/drivers

# MySQL Connector/J(バージョン 8.0.33)のダウンロード
sudo wget https://repo1.maven.org/maven2/mysql/mysql-connector-java/8.0.33/mysql-connector-java-8.0.33.jar

ステップ3:Logstashパイプラインの設定

/etc/logstash/conf.d/mysql_to_es.conf に設定ファイルを作成します。

input {
  jdbc {
    jdbc_driver_library => "/etc/logstash/drivers/mysql-connector-java-8.0.33.jar"
    jdbc_driver_class => "com.mysql.cj.jdbc.Driver"
    jdbc_connection_string => "jdbc:mysql://localhost:3306/shop_db?useSSL=false&serverTimezone=UTC"
    jdbc_user => "db_user"
    jdbc_password => "Secret_Password_123"
    
    # 実行間隔:1分ごと
    schedule => "* * * * *"
    
    # 更新日時のタイムスタンプを追跡
    use_column_value => true
    tracking_column => "updated_at"
    tracking_column_type => "timestamp"
    last_run_metadata_path => "/var/lib/logstash/.logstash_products_last_run"
    
    # 前回実行時より新しいデータのみを取得
    statement => "SELECT id, title, description, price, is_deleted, updated_at FROM products WHERE updated_at > :sql_last_value ORDER BY updated_at ASC"
  }
}

filter {
  mutate {
    remove_field => ["@version", "@timestamp"]
  }
}

output {
  elasticsearch {
    hosts => ["http://localhost:9200"]
    index => "products"
    document_id => "%{id}"
    action => "index"
  }
  
  # 必要に応じてデバッグ用にコンソールへログ出力
  stdout {
    codec => rubydebug
  }
}

ステップ4:パイプラインの起動と動作検証

まず、設定ファイルが正しいか構文チェックを行います。

# 設定ファイルの構文チェック
sudo /usr/share/logstash/bin/logstash -f /etc/logstash/conf.d/mysql_to_es.conf --config.test_and_exit

# Logstashサービスの起動と自動起動の有効化
sudo systemctl start logstash
sudo systemctl enable logstash

約1分後、cURLコマンドを実行してデータがElasticsearchに同期されたか確認します。

curl -X GET "http://localhost:9200/products/_search?pretty" -H 'Content-Type: application/json' -d'
{
  "query": {
    "match": {
      "title": "メカニカルキーボード"
    }
  }
}'

本番運用における実践的な注意点

テスト環境でパイプラインが安定して動作していても、本番環境で問題が発生しないとは限りません。特に注意すべき3つのポイントを挙げます。

  • データ削除の取り扱い(物理削除 vs 論理削除): Logstashは定期的に SELECT クエリを実行する仕組みのため、物理削除(Hard Delete)されたレコードを検知できません。そのため、論理削除(Soft Delete)(is_deleted = 1 のフラグ付与)を採用してください。Elasticsearch側では、検索クエリ時にこのフラグを除外するようフィルタリングします。
  • タイムゾーンの統一: MySQL、Linuxサーバー、Logstash間で必ずUTCに統一してください。タイムゾーンのズレによって :sql_last_value がデータを取得漏れするのを防ぐため、JDBC接続文字列の serverTimezone=UTC パラメータは必須です。
  • MySQLを信頼できる唯一の情報源(Single Source of Truth)として維持: Elasticsearchはあくまで検索用のキャッシュ層として位置づけてください。決済処理や在庫引き当てなどの重要トランザクションでElasticsearchのデータを直接読み込むことは避けるべきです。MySQLの定期バックアップを継続し、マスター・スレーブのデータフローを明確に定義しておきましょう。

Logstashを介してMySQLとElasticsearchを組み合わせる構成により、アプリケーションアーキテクチャのシンプルさと整合性を維持したまま、Elasticsearchの圧倒的な検索スピードを最大限に活用できます。

Share: