すべてのプロダクト
Search
ドキュメントセンター

Elasticsearch:Logstash によるセルフマネージド Elasticsearch からのデータ移行

最終更新日:Aug 21, 2026

このトピックでは、ECS インスタンスに Logstash をデプロイし、移行パイプラインを設定して、セルフマネージド Elasticsearch クラスターから Alibaba Cloud Elasticsearch へ全量データまたは増分データを移行する方法を説明します。

考慮事項

  • Logstash をホストする Elastic Compute Service (ECS) インスタンスは、Alibaba Cloud Elasticsearch クラスターと同じ Virtual Private Cloud (VPC) 内にあり、ソースクラスターと宛先クラスターの両方へのネットワークアクセスが必要です。

  • アプリケーションが継続的にデータを書き込みまたは更新する場合は、まず完全移行を実行し、その後タイムスタンプまたは他の識別フィールドに基づいて増分移行を実行してください。そうしないと、宛先クラスターで古いデータが新しいデータを上書きする可能性があります。宛先にすべての既存データがすでにある場合は、増分移行のみが必要です。

手順

  1. ステップ 1:環境とインスタンスの準備

    Alibaba Cloud Elasticsearch クラスターを作成し、ECS インスタンスにセルフマネージド Elasticsearch と Logstash をデプロイし、移行データを準備します。

  2. ステップ 2 (オプション):インデックスメタデータ (設定とマッピング) の移行

    ECS インスタンスで Python スクリプトを実行して、インデックスメタデータを移行します。

  3. ステップ 3:完全データ移行

    Logstash を使用して、セルフマネージドクラスターから Alibaba Cloud Elasticsearch にすべてのデータを移行します。

  4. ステップ 4:増分データ移行

  5. ステップ 5:移行結果の検証

ステップ 1:環境とインスタンスの準備

  1. Alibaba Cloud Elasticsearch インスタンスを作成します。

    Alibaba Cloud Elasticsearch インスタンスを作成します。テスト環境では、次の設定を使用します。

    パラメータ

    説明

    リージョン

    中国 (杭州)

    エディション

    Standard Edition 7.10.0

    インスタンス仕様

    3 つのゾーン、3 つのデータノード。各ノードには 4 vCPU、16 GB のメモリ、100 GB の拡張 SSD (ESSD) が搭載されています。

  2. セルフマネージド Elasticsearch、Kibana、Logstash インスタンス用の ECS インスタンスを作成します。

    ウィザードを使用してインスタンスを作成します。テスト環境では、次の設定を使用します。

    パラメータ

    説明

    リージョン

    中国 (杭州)

    インスタンスタイプ

    4 vCPU、16 GiB のメモリ

    イメージ

    パブリックイメージ、CentOS 7.9 64 ビット

    ストレージ

    システムディスク、100 GiB 拡張 SSD (ESSD)

    ネットワーク

    Alibaba Cloud Elasticsearch クラスターと同じ仮想プライベートクラウド (VPC) を選択します。パブリック IPv4 の割り当て を選択し、課金方法をトラフィック課金に、ピーク帯域幅を 100 Mbit/s に設定します。

    セキュリティグループ

    インバウンドルールを追加して、ポート 5601 (デフォルトの Kibana ポート) へのアクセスを許可します。承認オブジェクトをクライアントの IP アドレスに設定します。

    重要
    • クライアントが自宅または企業のネットワーク上にある場合は、コンピューターのプライベート IP ではなく、ネットワークのパブリックアウトバウンド IP を使用してください。パブリック IP は https://www.whatismyip.com で確認できます。

    • 承認オブジェクトとして 0.0.0.0/0 を設定すると、すべての IPv4 アドレスが許可されますが、ECS インスタンスがパブリックインターネットに公開されます。本番環境では使用しないでください。

  3. セルフマネージド Elasticsearch クラスターをデプロイします。

    このドキュメントでは、1 つのデータノードを持つセルフマネージド Elasticsearch 7.10.0 クラスターを使用します。

    1. ECS インスタンスに接続します。

      Workbench を使用して Linux インスタンスに接続します。

    2. root ユーザーとして、elastic という名前の新しいユーザーを作成します。

      useradd elastic
    3. elastic ユーザーのパスワードを設定します。

      passwd elastic

      プロンプトに従って、新しいパスワードを入力し、確認のため再入力します。

    4. elastic ユーザーに切り替えます。

      su -l elastic
    5. Elasticsearch インストールパッケージをダウンロードして展開します。

      wget https://artifacts.elastic.co/downloads/elasticsearch/elasticsearch-7.10.0-linux-x86_64.tar.gz
      tar -zvxf elasticsearch-7.10.0-linux-x86_64.tar.gz
    6. Elasticsearch を起動します。

      Elasticsearch インストールディレクトリに移動し、バックグラウンドでサービスを起動します。

      cd elasticsearch-7.10.0
      ./bin/elasticsearch -d
    7. Elasticsearch サービスが実行されていることを確認します。

      cd ~ 
      curl localhost:9200

      成功した場合、Elasticsearch のバージョン番号とタグライン "You Know, for Search" を含むレスポンスが返されます。

      [elastic@vm01 ~]$ curl localhost:9200
      {
        "name" : "vm01",
        "cluster_name" : "elasticsearch",
        "cluster_uuid" : "SRB4pnk4SmS-YHzsrxxx",
        "version" : {
          "number" : "7.10.0",
          "build_flavor" : "default",
          "build_type" : "tar",
          "build_hash" : "ef48eb35cf30adf4db14086e8aabd07ef6xxx",
          "build_date" : "2020-03-26T06:34:37.794943Z",
          "build_snapshot" : false,
          "lucene_version" : "8.4.0",
          "minimum_wire_compatibility_version" : "6.8.0",
          "minimum_index_compatibility_version" : "6.0.0-beta1"
        },
        "tagline" : "You Know, for Search"
      }
  4. セルフマネージド Kibana インスタンスをデプロイし、サンプルデータを準備します。

    このドキュメントでは、セルフマネージド Kibana 7.10.0 インスタンスを使用します。

    1. ECS インスタンスに接続します。

      Workbench を使用して Linux インスタンスに接続します。

      説明

      本ドキュメントの手順では、特に指定がない限り、コマンドは非 root ユーザーとして実行します。

    2. Kibana インストールパッケージをダウンロードして展開します。

      wget https://artifacts.elastic.co/downloads/kibana/kibana-7.10.0-linux-x86_64.tar.gz
      tar -zvxf kibana-7.10.0-linux-x86_64.tar.gz
    3. Kibana 設定ファイル config/kibana.yml を編集し、server.host: "0.0.0.0" を追加してリモートアクセスを有効にします。

      Kibana のインストールディレクトリに移動し、kibana.yml を編集します。

      cd kibana-7.10.0
      vi config/kibana.yml

      server.host の値を "0.0.0.0" に設定して、リモート接続を許可します。更新後のファイルの主要な設定は次のとおりです。

      # Kibana はバックエンドサーバーによって提供されます。この設定では、使用するポートを指定します。
      #server.port: 5601
      
      # Kibana サーバーがバインドするアドレスを指定します。IP アドレスとホスト名の両方が有効な値です。
      # デフォルトは 'localhost' で、通常はリモートマシンからの接続はできません。
      # リモートユーザーからの接続を許可するには、このパラメーターをループバック以外のアドレスに設定します。
      #server.host: "localhost"
      server.host: "0.0.0.0"
      # プロキシの背後で実行している場合に Kibana をマウントするパスを指定できます。
      # `server.rewriteBasePath` 設定は、受信リクエストから basePath を削除するかどうかを Kibana に指示し、起動時の非推奨警告を回避するために使用します。
      # この設定の末尾をスラッシュにすることはできません。
      #server.basePath: ""
    4. 非 root ユーザーとして Kibana を起動します。

      nohup ./bin/kibana &
    5. Kibana コンソールにログオンし、サンプルデータを追加します。

      1. ECS インスタンスのパブリック IP アドレスを使用して Kibana コンソールにアクセスします。

        URL の形式は次のとおりです: [http://<your_ecs_instance_public_ip>:5601/app/kibana#/home]。

      2. Kibana ホームページで、[サンプルデータをお試しください] をクリックします。

      3. [サンプルデータ] タブで、[サンプル Web ログ] カードを見つけ、カードの下部にある [データの追加] をクリックしてサンプルデータを追加します。

  5. セルフマネージド Logstash インスタンスをデプロイします。

    このドキュメントでは、1 つのノードを持つセルフマネージド Logstash 7.10.0 インスタンスを使用します。

    1. ECS インスタンスに接続します。

      Workbench を使用して Linux インスタンスに接続します。

      説明

      本ドキュメントの手順では、コマンドは非 root ユーザーとして実行します。

    2. ホームディレクトリに戻り、Logstash インストールパッケージをダウンロードして展開します。

      cd ~
      wget https://artifacts.elastic.co/downloads/logstash/logstash-7.10.0-linux-x86_64.tar.gz
      tar -zvxf logstash-7.10.0-linux-x86_64.tar.gz
    3. Logstash ヒープサイズを調整します。

      デフォルトのヒープサイズは 1 GB です。移行パフォーマンスを向上させるために、ECS インスタンス仕様に基づいて調整してください。

      Logstash インストールディレクトリに移動し、config/jvm.options を編集して、初期ヒープサイズと最大ヒープサイズの両方を 8 GB (-Xms8g および -Xmx8g) に設定します。

      cd logstash-7.10.0
      sudo vi config/jvm.options
      ## JVM 設定
      
      # Xms はヒープ領域全体の初期サイズを表します
      # Xmx はヒープ領域全体の最大サイズを表します
      
      -Xms8g
      -Xmx8g
      
      ################################################################
      ## エキスパート向け設定
      ################################################################
      ##
      ## このセクション以降のすべての設定はエキスパート向け設定です。
      ## 内容を理解せずに変更しないでください。
      ##
      ################################################################
      
      ## GC 設定
      -XX:+UseConcMarkSweepGC
      -XX:CMSInitiatingOccupancyFraction=75
      -XX:+UseCMSInitiatingOccupancyOnly
      
      ## ロケール
    4. Logstash バッチサイズを変更します。

      5 MB から 15 MB のバッチでデータを書き込むと、データ移行が高速化されます。

      config/pipelines.yml を編集し、pipeline.batch.size を 125 から 5000 に変更します。

      vi config/pipelines.yml
      #    # 設定テキストを読み込むパス
      #    path.config: "/etc/conf.d/logstash/myconfig.cfg"
      #
      #    # パイプラインの Filters+Outputs ステージを実行するワーカースレッド数
      #    pipeline.workers: 1 (実際のデフォルトは CPU 数)
      #
      #    # filters+workers に送信する前に入力から取得するイベント数
           pipeline.batch.size: 5000
      #
      #    # サイズが小さいバッチを filters+outputs にディスパッチする前に、次のイベントをポーリングする際の待機時間 (ミリ秒)
      #    pipeline.batch.delay: 50
      #
      #    # 内部キューイングモデル。「memory」はレガシーなメモリベースのキューイング、「persisted」はディスクベースの ACK 付きキューイングです。デフォルトは memory です。
      #    queue.type: memory
    5. Logstash が正常に機能していることを確認します。

      1. 標準入力を受け取り、標準出力に送信する簡単なパイプラインを実行します。

        bin/logstash -e 'input { stdin { } } output { stdout {} }'
      2. パイプラインが起動したら、"Hello world!" と入力して Enter キーを押します。

        Logstash が正常に動作している場合、"Hello world!" を含む構造化ログメッセージがコンソールに出力されます。

        [elastic@vm01 logstash-7.10.0]$ bin/logstash -e 'input { stdin { } } output { stdout {} }'
        Using bundled JDK: /home/elastic/logstash-7.10.0/jdk
        OpenJDK 64-Bit Server VM warning: Option UseConcMarkSweepGC was deprecated in version 9.0 a
        WARNING: An illegal reflective access operation has occurred
        WARNING: Illegal reflective access by org.jruby.ext.openssl.SecurityHelper (file:/tmp/jruby
        WARNING: Please consider reporting this to the maintainers of org.jruby.ext.openssl.Securit
        WARNING: Use --illegal-access-warn to enable warnings of further illegal reflective access
        WARNING: All illegal access operations will be denied in a future release
        Sending Logstash logs to /home/elastic/logstash-7.10.0/logs which is now configured via log
        [2022-03-21T15:39:24,470][INFO ][logstash.runner          ] Starting Logstash {"logstash.ve
        inux-x86_64]"}
        [2022-03-21T15:39:24,606][INFO ][logstash.setting.writabledirectory] Creating directory {:s
        [2022-03-21T15:39:24,618][INFO ][logstash.setting.writabledirectory] Creating directory {:s
        [2022-03-21T15:39:24,845][WARN ][logstash.config.source.multilocal] Ignoring the 'pipelines
        [2022-03-21T15:39:24,865][INFO ][logstash.agent           ] No persistent UUID file found.
        [2022-03-21T15:39:25,961][INFO ][org.reflections.Reflections] Reflections took 36 ms to sca
        [2022-03-21T15:39:26,356][INFO ][logstash.javapipeline    ][main] Starting pipeline {:pipel
        s"=>["config string"], :thread=>"#<Thread:0x75693a9 run>"}
        [2022-03-21T15:39:26,997][INFO ][logstash.javapipeline    ][main] Pipeline Java execution i
        [2022-03-21T15:39:27,032][INFO ][logstash.javapipeline    ][main] Pipeline started {"pipeli
        The stdin plugin is now waiting for input:
        [2022-03-21T15:39:27,073][INFO ][logstash.agent           ] Pipelines running {:count=>1, :
        [2022-03-21T15:39:27,211][INFO ][logstash.agent           ] Successfully started Logstash A
        Hello world!
        {
               "host" => "vm01",
            "@version" => "1",
            "message" => "\"Hello world!\"",
          "@timestamp" => 2022-03-21T07:39:46.598Z
        }

ステップ 2 (オプション):インデックスメタデータの移行

Logstash は、移行先クラスターにインデックスが存在しない場合、自動的にインデックスを作成しますが、自動生成された設定とマッピングはソースと異なる場合があります。インデックス構造の一貫性を確保するために、移行前に移行先インデックスを手動で作成してください。

次の Python スクリプトを使用して、移行先インデックスを作成します。

  1. ECS インスタンスに接続します。

    ワークベンチを使用して Linux インスタンスに接続します。

    説明

    このドキュメントの手順では、非 root ユーザーとしてコマンドを実行することを前提としています。

  2. Python スクリプトファイルを作成して開きます。このドキュメントでは、ファイル名として indiceCreate.py を使用します。

    sudo vi indiceCreate.py
  3. 次のコードを Python スクリプトファイルにコピーし、クラスターエンドポイント、ユーザー名、パスワードのプレースホルダー値を実際の認証情報に置き換えます。

    #!/usr/bin/python
    # -*- coding: UTF-8 -*-
    # ファイル名: indiceCreate.py
    import sys
    import base64
    import time
    import httplib
    import json
    ## ソースクラスターのホスト。
    oldClusterHost = "localhost:9200"
    ## ソースクラスターのユーザー名。空欄のままでも構いません。
    oldClusterUserName = "elastic"
    ## ソースクラスターのパスワード。空欄のままでも構いません。
    oldClusterPassword = "xxxxxx"
    ## 移行先クラスターのホスト。Alibaba Cloud Elasticsearch インスタンスの基本情報ページで確認できます。
    newClusterHost = "es-cn-zvp2m4bko0009****.elasticsearch.aliyuncs.com:9200"
    ## 移行先クラスターのユーザー名。
    newClusterUser = "elastic"
    ## 移行先クラスターのパスワード。
    newClusterPassword = "xxxxxx"
    DEFAULT_REPLICAS = 0
    def httpRequest(method, host, endpoint, params="", username="", password=""):
        conn = httplib.HTTPConnection(host)
        headers = {}
        if (username != "") :
            'Hello {name}, your age is {age} !'.format(name = 'Tom', age = '20')
            base64string = base64.encodestring('{username}:{password}'.format(username = username, password = password)).replace('\n', '')
            headers["Authorization"] = "Basic %s" % base64string;
        if "GET" == method:
            headers["Content-Type"] = "application/x-www-form-urlencoded"
            conn.request(method=method, url=endpoint, headers=headers)
        else :
            headers["Content-Type"] = "application/json"
            conn.request(method=method, url=endpoint, body=params, headers=headers)
        response = conn.getresponse()
        res = response.read()
        return res
    def httpGet(host, endpoint, username="", password=""):
        return httpRequest("GET", host, endpoint, "", username, password)
    def httpPost(host, endpoint, params, username="", password=""):
        return httpRequest("POST", host, endpoint, params, username, password)
    def httpPut(host, endpoint, params, username="", password=""):
        return httpRequest("PUT", host, endpoint, params, username, password)
    def getIndices(host, username="", password=""):
        endpoint = "/_cat/indices"
        indicesResult = httpGet(oldClusterHost, endpoint, oldClusterUserName, oldClusterPassword)
        indicesList = indicesResult.split("\n")
        indexList = []
        for indices in indicesList:
            if (indices.find("open") > 0):
                indexList.append(indices.split()[2])
        return indexList
    def getSettings(index, host, username="", password=""):
        endpoint = "/" + index + "/_settings"
        indexSettings = httpGet(host, endpoint, username, password)
        print (index + "  元の設定:\n" + indexSettings)
        settingsDict = json.loads(indexSettings)
        ## シャード数はデフォルトでソースインデックスと一致します。
        number_of_shards = settingsDict[index]["settings"]["index"]["number_of_shards"]
        ## デフォルトのレプリカ数は 0 です。
        number_of_replicas = DEFAULT_REPLICAS
        newSetting = "\"settings\": {\"number_of_shards\": %s, \"number_of_replicas\": %s}" % (number_of_shards, number_of_replicas)
        return newSetting
    def getMapping(index, host, username="", password=""):
        endpoint = "/" + index + "/_mapping"
        indexMapping = httpGet(host, endpoint, username, password)
        print (index + " 元のマッピング:\n" + indexMapping)
        mappingDict = json.loads(indexMapping)
        mappings = json.dumps(mappingDict[index]["mappings"])
        newMapping = "\"mappings\" : " + mappings
        return newMapping
    def createIndexStatement(oldIndexName):
        settingStr = getSettings(oldIndexName, oldClusterHost, oldClusterUserName, oldClusterPassword)
        mappingStr = getMapping(oldIndexName, oldClusterHost, oldClusterUserName, oldClusterPassword)
        createstatement = "{\n" + str(settingStr) + ",\n" + str(mappingStr) + "\n}"
        return createstatement
    def createIndex(oldIndexName, newIndexName=""):
        if (newIndexName == "") :
            newIndexName = oldIndexName
        createstatement = createIndexStatement(oldIndexName)
        print ("新しいインデックス " + newIndexName + " の設定とマッピング:\n" + createstatement)
        endpoint = "/" + newIndexName
        createResult = httpPut(newClusterHost, endpoint, createstatement, newClusterUser, newClusterPassword)
        print ("新しいインデックス " + newIndexName + " の作成結果: " + createResult)
    ## メイン
    indexList = getIndices(oldClusterHost, oldClusterUserName, oldClusterPassword)
    systemIndex = []
    for index in indexList:
        if (index.startswith(".")):
            systemIndex.append(index)
        else :
            createIndex(index, index)
    if (len(systemIndex) > 0) :
        for index in systemIndex:
            print (index + " はシステムインデックスの可能性があるため、再作成されません。必要に応じて個別に処理してください。")
  4. Python スクリプトを実行して、移行先インデックスを作成します。

    sudo /usr/bin/python indiceCreate.py
  5. 宛先クラスターのKibana コンソールにログオンし、インデックスが作成されたことを確認します。

    GET /_cat/indices?v

ステップ 3:完全データ移行の実行

  1. ECS インスタンスに接続します。

  2. config ディレクトリで、Logstash 設定ファイルを作成して開きます。

    cd logstash-7.10.0/config
    vi es2es_all.conf
  3. ファイルに次の設定を追加します。

    説明
    • Logstash の設定パラメーターはバージョン 8.5 で変更されました。このドキュメントでは、バージョン 7.10.0 とバージョン 8.5.1 の両方の設定例を提供します。

    • データの正確性を確保するために、個別の Logstash パイプライン設定ファイルを作成し、データをバッチで移行してください。

    バージョン 7.10.0

    input{
        elasticsearch{
            # ソース Elasticsearch クラスターのエンドポイント。
            hosts =>  ["http://localhost:9200"]
            # ソースクラスターのユーザー名とパスワード。
            user => "xxxxxx"
            password => "xxxxxx"
            # 移行するインデックスのリスト。複数のインデックスはカンマ (,) で区切ります。
            index => "kibana_sample_data_*"
            # 次の 3 つの項目はデフォルトのままで構いません。これらはスレッド数、移行データサイズ、Logstash JVM 設定に関するものです。
            docinfo=>true
            slices => 5
            size => 5000
        }
    }
    
    filter {
      # Logstash が追加したメタデータフィールドを削除します。
      mutate {
        remove_field => ["@timestamp", "@version"]
      }
    }
    
    output{
        elasticsearch{
            # 移行先クラスターのエンドポイント。 Alibaba Cloud Elasticsearch インスタンスの基本情報ページで確認できます。
            hosts => ["http://es-cn-zvp2m4bko0009****.elasticsearch.aliyuncs.com:9200"]
            # 移行先クラスターのユーザー名とパスワード。
            user => "elastic"
            password => "xxxxxx"
            # 移行先インデックスの名前。 この設定では、インデックス名をソースと同じに保ちます。
            index => "%{[@metadata][_index]}"
            # 移行先インデックスのタイプ。 この設定では、インデックスタイプをソースと同じに保ちます。
            document_type => "%{[@metadata][_type]}"
            # 移行先クラスターのデータの ID。 元のドキュメント ID を保持する必要がない場合は、パフォーマンスを向上させるためにこの行を削除できます。
            document_id => "%{[@metadata][_id]}"
            ilm_enabled => false
            manage_template => false
        }
    }

    バージョン 8.5.1

    input{
        elasticsearch{
            # ソース Elasticsearch クラスターのエンドポイント。
            hosts =>  ["http://es-cn-uqm3811160002***.elasticsearch.aliyuncs.com:9200"]
            # ソースクラスターのユーザー名とパスワード。
            user => "elastic"
            password => ""
            # 移行するインデックスのリスト。複数のインデックスはカンマ (,) で区切ります。
            index => "test_ecommerce"
            # 次の項目はデフォルトのままで構いません。 これらは移行データサイズと Logstash JVM 設定に関するものです。
            docinfo => true
            size => 10000
            docinfo_target => "[@metadata]"
        }
    }
    
    filter {
      # Logstash が追加したメタデータフィールドを削除します。
      mutate {
        remove_field => ["@timestamp","@version"]
      }
    }
    
    output{
        elasticsearch{
            # 移行先クラスターのエンドポイント。 Alibaba Cloud Elasticsearch インスタンスの基本情報ページで確認できます。
            hosts => ["http://es-cn-nwy38aixp0001****.elasticsearch.aliyuncs.com:9200"]
            # 移行先クラスターのユーザー名とパスワード。
            user => "elastic"
            password => ""
            # 移行先インデックスの名前。 この設定では、インデックス名をソースと同じに保ちます。
            index => "%{[@metadata][_index]}"
            # 移行先クラスターのデータの ID。 元のドキュメント ID を保持する必要がない場合は、パフォーマンスを向上させるためにこの行を削除できます。
            document_id => "%{[@metadata][_id]}"
            ilm_enabled => false
            manage_template => false
        }
    }

    Elasticsearch 入力プラグインは、すべてのデータを読み取った後に停止します。一部の環境では、Logstash が自動的に再起動し、重複書き込みが発生する可能性があります。これを防ぐには、schedule パラメーターに Cron 式を指定して、特定の時刻にタスクを実行します (スケジューリング)。

    たとえば、3 月 5 日の午後 1 時 20 分にタスクを実行する場合は、次のようにします。

    schedule => "20 13 5 3 *"
  4. Logstash ディレクトリに移動します。

    cd ~/logstash-7.10.0
  5. 完全データ移行タスクを開始します。

    nohup bin/logstash -f config/es2es_all.conf >/dev/null 2>&1 &

ステップ 4:増分データの移行

  1. ECS インスタンスに接続します。config ディレクトリで、増分移行用の新しい Logstash 設定ファイルを作成して開きます。

    cd config
    vi es2es_kibana_sample_data_logs.conf
    説明

    このドキュメントの手順では、非 root ユーザーとしてコマンドを実行することを前提としています。

  2. ファイルに次の設定を追加します。

    以下は、バージョン 7.10.0 の設定例です。

    説明
    • Logstash 8.5 以降では、ドキュメントタイプが非推奨のため、document_type => "%{[@metadata][_type]}" の行を削除する必要があります。

    • ファイルを設定した後、スケジュールされた Logstash タスクを開始すると、増分移行がトリガーされます。

    input{
        elasticsearch{
            # ソース Elasticsearch クラスターのエンドポイント。
            hosts =>  ["http://localhost:9200"]
            # ソースクラスターのユーザー名とパスワード。
            user => "xxxxxx"
            password => "xxxxxx"
            # 移行するインデックスのリスト。複数のインデックスはカンマ (,) で区切ります。
            index => "kibana_sample_data_logs"
            # 時間範囲内の増分データをクエリします。次の設定では、過去 5 分間のデータをクエリします。
            query => '{"query":{"range":{"@timestamp":{"gte":"now-5m","lte":"now/m"}}}}'
            # スケジュールされたタスク。次の設定では、タスクを毎分実行します。
            schedule => "* * * * *"
            scroll => "5m"
            docinfo=>true
            size => 5000
        }
    }
    
    filter {
      # Logstash が追加したメタデータフィールドを削除します。
      mutate {
        remove_field => ["@timestamp", "@version"]
      }
    }
    
    
    output{
        elasticsearch{
            # 宛先クラスターのエンドポイント。Alibaba Cloud Elasticsearch インスタンスの基本情報ページで確認できます。
            hosts => ["http://es-cn-zvp2m4bko0009****.elasticsearch.aliyuncs.com:9200"]
            # 宛先クラスターのユーザー名とパスワード。
            user => "elastic"
            password => "xxxxxx"
            # 宛先インデックスの名前。この設定では、ソースと同じインデックス名にします。
            index => "%{[@metadata][_index]}"
            # 宛先ドキュメントタイプ。この設定では、ソースと同じドキュメントタイプにします。
            document_type => "%{[@metadata][_type]}"
            # 宛先クラスターのデータの ID。元のドキュメント ID を保持する必要がない場合は、パフォーマンスを向上させるためにこの行を削除できます。
            document_id => "%{[@metadata][_id]}"
            ilm_enabled => false
            manage_template => false
        }
    }
    重要
    • Logstash は UTC タイムスタンプを使用します。ソースデータが別のタイムゾーンを使用している場合は、それに応じてクエリ範囲を調整してください。@timestamp フィールドの now-5m は、サーバーの UTC クロックに基づきます。

    • ソースインデックスには、増分同期用の時間フィールドを含める必要があります。含まれていない場合は、_ingest.timestamp メタデータフィールドを使用した インジェストパイプラインを使用して、インデックス作成時にドキュメントに @timestamp を追加してください。

  3. Logstash ディレクトリに移動します。

    cd ~/logstash-7.10.0
  4. 増分データ移行タスクを開始します。

    sudo nohup bin/logstash -f config/es2es_kibana_sample_data_logs.conf >/dev/null 2>&1 &
  5. 宛先 Elasticsearch クラスターの Kibana コンソールで、最新のレコードをクエリして、増分データが同期されていることを確認します。

    次のクエリは、kibana_sample_data_logs インデックスから過去 5 分間のレコードを検索します。

    GET kibana_sample_data_logs/_search
    {
      "query": {
        "range": {
          "@timestamp": {
            "gte": "now-5m",
            "lte": "now/m"
          }
        }
      },
      "sort": [
        {
          "@timestamp": {
            "order": "desc"
          }
        }
      ]
    }
                            

ステップ 5:移行結果の検証

  1. 完全なデータ移行を検証します。

    1. 自己管理型ソースクラスターでインデックスとドキュメント数の情報を確認します。

      GET _cat/indices?v

      以下に結果の例を示します。

      GET _cat/indices?v
      
      health status index                    uuid                   pri rep docs.count docs.deleted store.size pri.store.size
      green  open   .kibana_task_manager_1   CxAx5J2sT0qHPsWV       1   0   2          0            6.6kb      6.6kb
      green  open   .apm-agent-configuration dYz5bh4dTomjtDP3       1   0   0          0            283b       283b
      green  open   kibana_sample_data_logs  PUBQrSkJRMGyI-cV       1   0   14074      0            11.6mb     11.6mb
      green  open   .kibana_1                MXhG2XbYTYSORB8G       1   0   49         4            139.5kb    139.5kb
    2. 移行前に、Alibaba Cloud の移行先クラスターでインデックスとドキュメント数を確認してください。

      以下は、移行前の Alibaba Cloud Elasticsearch デスティネーションクラスターのインデックス情報の一例です。

      GET _cat/indices?v
      
      health status index                          uuid                 pri rep docs.count docs.deleted store.size pri.store.size
      green  open   .aliyun-limiter-group          5K4N8YNUSxeJZCXPxxx   1   1          0            0       522b           261b
      green  open   .apm-agent-configuration       vaVC28KVQMCsABwuxxx   1   1          0            0       522b           261b
      green  open   .monitoring-es-7-2022.03.19    9NUdZCaAQw-426Zrxxx   1   1     207485        15328    229.8kb          9.9kb
      green  open   highlight_unified              PubNS7HIRR2B5FIfxxx   1   1          2            0     19.8kb          9.9kb
      green  open   .monitoring-es-7-2022.03.18    kEP-0LeeSh01-kg2xxx   1   1     117792            0    132.3mb         60.3mb
      green  open   .aliyun-limiter-config         6SJImN0bRoap3fYMxxx   1   1          0            0       522b           261b
      green  open   .kibana_1                      0RRrLWLCT4aaT-1fxxx   1   1         27            4     20.8mb         10.4mb
      green  open   .security-7                    D7Ux5eq7S5WtYH_Yxxx   1   1         55            0    397.7kb        198.4kb
      green  open   .monitoring-es-7-2022.03.21    n6DZS66KRmW1zaN7xxx   1   1      85969         5244    102.5mb         51.5mb
      green  open   .apm-custom-link               SBnBUOojSd-Vt3xxxxx   1   1          0            0       522b           261b
      green  open   .monitoring-kibana-7-2022.03.20 eHPFB1h4Q8yxYbxAxxx  1   1      17278            0      5.8mb          2.8mb
      green  open   .kibana_task_manager_1         iDK1EK-iR22Gkhfxxxx   1   1          6           68      157kb         66.1kb
      green  open   .monitoring-kibana-7-2022.03.21 YIivw66dSBi0_Rwuxxx  1   1       6062            0      4.5mb          2.2mb
      green  open   kibana_sample_data_logs        1zaN5Ji7RWqbFwKZxxx   1   0          0            0       208b           208b
      green  open   .kibana-event-log-7.16.0-000001 ImZU-V4KRq2K3EUxxxx 1   1          1            0     11.4kb          5.7kb
      green  open   highlight_fvh                  sErtUXXpToiiPSaSxxx   1   1          2            0     23.5kb         11.7kb
      green  open   .monitoring-es-7-2022.03.20    SyOns3d-QU6ysbFDxxx   1   1     224781        43812    246.6mb        124.1mb
      green  open   .monitoring-kibana-7-2022.03.18 gOvcKvRlQ9O-PipQxxx  1   1      10700            0      3.4mb          1.7mb
      green  open   .monitoring-kibana-7-2022.03.19 IwSi_UIYQ5eFSyUJxxx  1   1      17280            0      5.8mb          2.9mb
    3. 全データ移行後、Alibaba Cloud の移行先クラスターでインデックスとドキュメント数の情報を再度確認します。

      ドキュメント数はソースクラスターの数と一致する必要があります。Kibana 開発ツールで GET _cat/indices?v コマンドを実行します。結果には、すべてのクラスターインデックスの ヘルスgreenステータスopen であると表示されます。kibana_sample_data_logs インデックスの docs.count は 14074、store.size は 11.6 MB であり、データがデスティネーションクラスターに正常に移行されたことを確認できます。

  2. 増分データ移行を検証します。

    セルフマネージドソースクラスターで最新のレコードを確認します。

    GET kibana_sample_data_logs/_search
    {
      "query": {
        "range": {
          "@timestamp": {
            "gte": "now-5m",
            "lte": "now/m"
          }
        }
      },
      "sort": [
        {
          "@timestamp": {
            "order": "desc"
          }
        }
      ]
    }

    以下に結果の例を示します。

    {
      "_source" : {
        "agent" : "Mozilla/5.0 (X11; Linux x86_64; rv:6.0a1) Gecko/20110421 Firefox/6.0a1",
        "bytes" : 658,
        "clientip" : "171.66.xxx",
        "extension" : "",
        "geo" : {
          "srcdest" : "CN:US",
          "src" : "CN",
          "dest" : "US",
          "coordinates" : {
            "lat" : 45.54039389,
            "lon" : -122.9498258
          }
        },
        "host" : "www.elastic.co",
        "index" : "kibana_sample_data_logs",
        "ip" : "171.66.xxx",
        "machine" : {
          "ram" : 3221225xxx,
          "os" : "win 7"
        },
        "memory" : null,
        "message" : "171.66.xxx - - [2022-03-21T09:23:11.012Z] \"GET /security-analytics Gecko/20110421 Firefox/6.0a1\"",
        "phpmemory" : null,
        "referer" : "http://www.elastic-elastic-elastic.com/success/albert-sacco",
        "request" : "/security-analytics",
        "response" : 200,
        "tags" : [
          "success",
          "security"
        ],
        "timestamp" : "2022-03-21T09:23:11.012Z",
        "url" : "https://www.elastic.co/solutions/security-analytics",
        "utc_time" : "2022-03-21T09:23:11.012Z",
        "event" : {
          "dataset" : "sample_web_logs"
        }
      },
      "sort" : [
        1647854591012
      ]
    }

    デスティネーションクラスターの Kibana コンソールで同じクエリを実行します。結果が一致すれば、増分同期が成功したことを確認できます。