Elasticsearch 6.8から7.10へのクラスタ管理とタスク管理機能の進化

概要

Elasticsearch 7.10では、6.8から大幅な改善が加えられました。特にクエリの自動キャンセル、投票専用ノードロール、非同期検索、検索可能なスナップショットといった重要な機能追加により、クラスタ管理とタスク管理の柔軟性が向上しました。

クエリの自動キャンセル機能

接続が切断された際、関連するバックグラウンド処理も自動的に停止するようになります。この機能はElasticsearch 7.4で正式に導入されました。

Elasticsearchは、発信元の接続が閉じられた際に、_searchエンドポイント経由で送信されたクエリを自動的に終了します。

実験コード

以下のPainlessスクリプトを使用して、長時間実行されるクエリをシミュレートします。

GET /_search?max_concurrent_shard_requests=1
{
    "query": {
        "bool": {
            "must": [
                {
                    "script": {
                        "script": {
                            "lang": "painless",
                            "source": """
                                long total = 0;
                                for (int x = 0; x < 50000; x++) {
                                    total += x * x;
                                }
                                return true;
                            """
                        }
                    }
                },
                {
                    "script": {
                        "script": {
                            "lang": "painless",
                            "source": """
                                long result = 1;
                                for (int y = 1; y < 80000; y++) {
                                    result = result * y;
                                }
                                return true;
                            """
                        }
                    }
                },
                {
                    "script": {
                        "script": {
                            "lang": "painless",
                            "source": """
                                long value = 1;
                                for (int z = 1; z < 90000; z++) {
                                    value *= z;
                                }
                                long squaredTotal = 0;
                                for (int w = 0; w < 70000; w++) {
                                    squaredTotal += w * w;
                                }
                                return true;
                            """
                        }
                    }
                },
                {
                    "script": {
                        "script": {
                            "lang": "painless",
                            "source": """
                                long first = 0;
                                long second = 1;
                                long current;
                                for (int n = 0; n < 60000; n++) {
                                    current = first + second;
                                    first = second;
                                    second = current;
                                }
                                return true;
                            """
                        }
                    }
                }
            ]
        }
    }
}

Pythonテストスクリプト

import requests
import threading
import time
from datetime import datetime

ELASTICSEARCH_URL = "http://localhost:9201"

COMPLEX_QUERY = {
    "size": 0,
    "query": {
        "bool": {
            "must": [
                {
                    "script": {
                        "script": {
                            "lang": "painless",
                            "source": """
                                long total = 0;
                                for (int x = 0; x < 50000; x++) {
                                    total += x * x;
                                }
                                return true;
                            """
                        }
                    }
                },
                {
                    "script": {
                        "script": {
                            "lang": "painless",
                            "source": """
                                long result = 1;
                                for (int y = 1; y < 80000; y++) {
                                    result = result * y;
                                }
                                return true;
                            """
                        }
                    }
                }
            ]
        }
    }
}

def timestamp_log(message):
    now = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
    print(f"[{now}] {message}")

def execute_search():
    try:
        timestamp_log("検索開始...")
        session = requests.Session()
        response = session.post(
            f"{ELASTICSEARCH_URL}/_search",
            json=COMPLEX_QUERY,
            timeout=3
        )
        timestamp_log(f"応答コード: {response.status_code}")
    except Exception as e:
        timestamp_log(f"検索中断: {str(e)}")

def monitor_tasks():
    time.sleep(2)
    timestamp_log("タスク状態確認開始...")
    try:
        task_check = f"{ELASTICSEARCH_URL}/_tasks?detailed=true&actions=*search*"
        response = requests.get(task_check)
        if response.status_code == 200:
            nodes_data = response.json().get("nodes", {})
            if nodes_data:
                timestamp_log(f"実行中のタスクが検出されました: {len(nodes_data)}個")
            else:
                timestamp_log("実行中のタスクはありません")
        else:
            timestamp_log("タスク取得失敗")
    except Exception as e:
        timestamp_log(f"監視エラー: {str(e)}")

def run_test():
    search_thread = threading.Thread(target=execute_search)
    monitor_thread = threading.Thread(target=monitor_tasks)
    
    search_thread.start()
    monitor_thread.start()
    
    search_thread.join(timeout=5)
    monitor_thread.join(timeout=5)

if __name__ == "__main__":
    run_test()

動作確認結果

Elasticsearch 6.8では、接続を切断してもバックグラウンドタスクが残存する一方、バージョン7.10以降では自動的にタスクが終了します。

非同期検索機能

大量データを対象とした検索で、リアルタイム性よりも完了保証が重要になるケースに適した機能です。

基本的な使用方法

POST my_index/_async_search?keep_on_completion=true
{
  "query": {
    "match_all": {}
  }
}

レスポンス例

{
  "id": "RmxleGlibGUtQXN5bmMtU2VhcmNoLUlk", // 結果識別子
  "is_partial": false, // 部分結果かどうか
  "is_running": false, // 実行中かどうか
  "start_time_in_millis": 1738978637287, // 開始時刻
  "expiration_time_in_millis": 1739410637287, // 有効期限
  "response": {
    "took": 2,
    "timed_out": false,
    "_shards": {
      "total": 1,
      "successful": 1,
      "skipped": 0,
      "failed": 0
    },
    "hits": {
      "total": {
        "value": 5,
        "relation": "eq"
      },
      "max_score": 1.0,
      "hits": [...]
    }
  }
}

重要なパラメータ

  • wait_for_completion_timeout: クエリ完了までの待機時間を設定(デフォルト1秒)
  • keep_on_completion: 完了後も結果を保持するか(デフォルトfalse)
  • keep_alive: 結果の保持期間(デフォルト5日間)

投票専用マスターノード

マスターノード選出プロセスに参加するが、自身がマスターになることはない特別なノードタイプです。

役割と利点

  • 選出プロセスに参加し、票の決定的役割を果たす
  • マスターノードとしての負荷を持たないため、安定性が向上
  • 高可用性構成において、奇数ノードのバランスを保つ

検索可能なスナップショット

これは有料機能であり、ストレージコストを削減しつつ、データ復旧能力を維持できる仕組みです。

動作原理

スナップショットからインデックスをマウントすることで、完全なコピーを作成することなく検索可能にします。必要なシャードデータのみをローカルに展開し、検索パフォーマンスを確保します。

実装手順

  1. S3互換ストレージを設定(MinIOなど)
  2. スナップショットリポジトリを登録
  3. スナップショットからインデックスをマウント
# リポジトリ登録
PUT _snapshot/my-backup-repo
{
  "type": "s3",
  "settings": {
    "bucket": "my-es-backups",
    "endpoint": "http://127.0.0.1:9000",
    "access_key": "your-access-key",
    "secret_key": "your-secret-key"
  }
}

# スナップショットマウント
POST /_snapshot/my-backup-repo/daily-snapshot-2024/_mount
{
  "index": "source-index",
  "renamed_index": "searchable-snapshot-index",
  "index_settings": {
    "index.number_of_replicas": 0
  }
}

タグ: Elasticsearch cluster-management async-search searchable-snapshots voting-only-nodes

7月24日 23:27 投稿