使用LogStash将数据同步至PolarSearch

更新时间:
复制 MD 格式

当您需要将自建ElasticsearchOpenSearch集群中的数据迁移至PolarSearch时,可以使用LogStash实现数据同步。本文介绍如何配置LogStash,以完成全量数据迁移和持续的增量数据同步两种场景。

功能简介

LogStash是一款开源的数据处理管道工具,能够从多种来源采集数据,经过转换后,再发送到多种目标存储。通过为其安装针对不同数据源(如Elasticsearch、OpenSearch)和目标(如PolarSearch)的插件,您可以构建一条高效、灵活的数据同步链路,以满足数据迁移或实时同步的需求。

准备工作

在开始操作前,请确保满足以下环境和权限要求:

  • 环境依赖:已准备一台可以同时访问源库和目标库网络的服务器用于运行LogStash。

  • 源库权限:用于连接源库的账号需具备待同步索引的readread_metadata权限。

  • 目标库权限:用于连接PolarSearch的账号需具备writecreate_index等数据写入权限。

  • 增量同步要求:若执行增量同步,源索引中必须包含一个时间戳字段(如@timestamp),用于标识文档的写入或更新时间。

步骤一:准备LogStash环境

LogStash支持不同的源和目标,本文使用的目标库为PolarSearch,源库为Elasticsearch 7.10PolarSearch。

  1. 下载并解压LogStash:请根据您的环境选择对应的安装包。本文以Linux x86_64环境8.8.2版本为例。

    说明
    • 不同版本的LogStash在使用上略有差异但核心逻辑一致。

    • 更多LogStash版本,请参见LogStash版本列表

  2. 安装插件:不同的源库与目标库需要安装不同的插件,请根据您的实际情况选择安装。插件仓库地址,请参见插件仓库

    • 源库插件:

      • 若源库为Elasticsearch,则需要安装input-elasticsearch插件。当前插件已经预安装,您可通过命令查看是否已经安装。

        # 进入 LogStash 根目录
        cd /path/to/logstash-8.8.2
        
        # 检查插件是否已经安装
        ./bin/logstash-plugin list
        
        # 如果没有安装,则执行插件安装命令
        ./bin/logstash-plugin install logstash-input-elasticsearch
      • 若源库为OpenSearch,则需要安装input-opensearch插件。

        # 进入 LogStash 根目录
        cd /path/to/logstash-8.8.2
        
        # 执行插件安装命令
        ./bin/logstash-plugin install logstash-input-opensearch
    • 目标库插件:安装output-opensearch插件,PolarSearch兼容OpenSearch接口,导入PolarSearch需安装output-opensearch插件。

      # 进入 LogStash 根目录
      cd /path/to/logstash-8.8.2
      
      # 执行插件安装命令
      ./bin/logstash-plugin install logstash-output-opensearch

步骤二:(可选)为源数据准备时间戳字段

若您计划进行增量同步(包括全量完成后继续增量,或对持续写入的数据进行同步),但源索引中缺少可用的时间戳字段(如@timestamp),请参照以下步骤为新写入的数据自动注入时间戳。

重要
  • 若源索引中已有此类字段(如@timestamp),可直接在后续增量备份过程中通过query引用,跳过当前步骤。

  • 数据更新捕获限制:若您的业务使用POST /_update接口局部更新文档,该操作不会触发Ingest Pipeline刷新时间戳。因此,这类更新无法被基于时间戳的增量同步捕获。建议在应用层维护一个随每次变更而更新的时间戳字段。

  1. 创建Ingest Pipeline:
    在源库执行以下命令,创建一个名为add-timestamp-pipeline的管道,它会将文档的写入时间_ingest.timestamp存入@timestamp 字段。

    说明

    _ingest.timestampOpenSearch/Elasticsearch在处理文档时由服务端自动生成的内部元数据,代表文档实际写入的时间,精度为毫秒,时区为UTC,不受客户端时钟影响。在处理器中需通过三重括号语法{{{_ingest.timestamp}}}引用,以避免Mustache模板对特殊字符进行转义。

    PUT _ingest/pipeline/add-timestamp-pipeline
    {
      "description": "为文档自动注入写入时间戳至 @timestamp 字段",
      "processors": [
        {
          "set": {
            "field": "@timestamp",
            "value": "{{{_ingest.timestamp}}}",
            "override": false 
          }
        }
      ]
    }
    

    参数说明

    参数

    说明

    field

    目标字段名,此处写入@timestamp,也可根据实际需要自定义,如ingest_time

    value

    使用三重括号{{{...}}}语法引用_ingest.timestamp,确保时间戳原样写入,不被Mustache转义处理。

    override

    • false:表示若文档中已存在@timestamp字段则保留原值不覆盖。

    • true(默认值):始终以写入时间覆盖。

  2. 为索引设置默认Pipeline。
    将此Pipeline应用于目标索引,确保所有写入都自动添加时间戳。

    • 方式一(推荐,用于新索引):通过索引模板配置。

      PUT _index_template/my-index-template
      {
        "index_patterns": ["your-business-logs-*"],
        "template": {
          "settings": {
            "index.default_pipeline": "add-timestamp-pipeline"
          },
          "mappings": {
            "properties": {
              "@timestamp": {
                "type": "date"
              }
            }
          }
        }
      }
      
    • 方式二(用于已有索引):直接修改索引设置。

      说明

      通过_settings API仅对新写入的数据生效,已存量文档不会追溯注入时间戳。如需为存量数据补充时间戳,请使用Reindex API 配合pipeline参数重新索引。

      PUT /your-business-logs-index/_settings
      {
        "index.default_pipeline": "add-timestamp-pipeline"
      }
  3. 验证Pipeline:
    使用_simulate API测试效果,确认返回结果的_source中已包含@timestamp字段。

    POST _ingest/pipeline/add-timestamp-pipeline/_simulate
    {
      "docs": [
        {
          "_source": {
            "message": "test log entry",
            "level": "INFO"
          }
        }
      ]
    }

    预期结果:

    {
      "docs": [
        {
          "doc": {
            "_source": {
              "message": "test log entry",
              "level": "INFO",
              "@timestamp": "2026-04-08T02:59:32.249592033Z"
            },
            "_ingest": {
              "timestamp": "2026-04-08T02:59:32.249592033ZZ"
            }
          }
        }
      ]
    }
    

步骤三:创建同步配置文件

LogStash根目录下创建一个synchronization.conf配置文件,并根据您的场景选择相应的配置。

场景一:全量数据迁移

此配置用于一次性将源库的指定索引完整复制到目标库。

相关参数说明

  • input源库中的index字段支持使用通配符*以同步多个索引。然而,不建议使用全通配逻辑*来复制所有索引,因为这可能会导致不必要的内部索引复制。

  • input源库中的docinfo字段,可以获取原始索引名和文档ID。

  • output目标库中index名字等可以通过记录的metadata读取,从而保持索引名等不变或在原有名字等元信息基础上进行定制。

  • output目标库中可以增加stdout调试选项可输出调试信息。

更多信息,请参见同步Logstash事件至OpenSearch

调试输出

如下所示,通过该输出可以观察到记录的_indexmetadata及其结构层次。例如,_index字段的结构层次表明其在output等后续流程中的提取逻辑为:[@metadata][input][opensearch][_index]

重要

不同input插件的metadata提取逻辑存在差异。

[2025-10-23T01:36:52,765][INFO ][logstash.inputs.opensearch][main][bb6f7ddd51cb0a42a08aa287556a451e8939ecbc0ce17c51ddffe0ae0b0d217e] Slice complete {:slice_id=>2, :slices=>4}
[2025-10-23T01:36:52,787][INFO ][logstash.inputs.opensearch][main][bb6f7ddd51cb0a42a08aa287556a451e8939ecbc0ce17c51ddffe0ae0b0d217e] Slice complete {:slice_id=>0, :slices=>4}
{
      "@timestamp" => 2025-10-22T17:36:52.780903870Z,
           "price" => 29.99,
        "@metadata" => {
        "input" => {
            "opensearch" => {
                "_index" => "lptest",
                   "_id" => "sQXXDJoBvnBnbbz3pKb9",
                 "_type" => nil
            }
        }
    },
        "@version" => "1",
     "description" => "这是一个示例文档",
           "title" => "示例标题",
       "timestamp" => "2024-01-15T10:30:00Z"
}
[2025-10-23T01:36:57,204][INFO ][logstash.javapipeline    ][main] Pipeline terminated {"pipeline.id"=>"main"}
[2025-10-23T01:36:57,263][INFO ][logstash.pipelinesregistry] Removed pipeline from registry successfully {:pipeline_id=>:main}

配置文件示例

重要
  • 为避免同步不必要的系统内部索引(如.kibana),在配置中指定索引时,不建议使用*.*等全通配符。

  • 为确保目标库的索引配置(如分词器、字段类型)与源端一致,建议在同步前,手动在PolarSearch中创建好索引模板。

# synchronization.conf
# 一个从 Elasticsearch/OpenSearch 向 PolarSearch 同步数据的完整配置示例。

input {
  # 如果源集群是 Elasticsearch,请将 'opensearch' 替换为 'elasticsearch'。
  # 两个插件的配置参数基本相同,但元数据路径可能存在差异。
  opensearch {
    # 【必填】源集群的连接地址,建议使用 HTTPS 协议。
    hosts => ["https://source-cluster-endpoint:9200"]
    # 【必填】源集群的认证凭据。
    user => "your_source_user"
    password => "your_source_password"

    # 【必填】指定需要同步的索引,支持通配符。
    # 为避免同步不必要的内部索引(如 .kibana),不建议使用 "*" 或 ".*"。
    index => "your-business-logs-*"

    # --- 安全配置 ---
    # 如果源集群启用了 SSL/TLS,请设置为 true。
    ssl => true
    # 如果源集群使用自签名证书,取消注释并指定 CA 证书路径。
    # cacert => "/path/to/source_ca.crt"

    # --- 性能与元数据 ---
    # 开启此选项以获取原始索引名和文档 ID。
    docinfo => true
    # 设置并发读取数,建议设置为源索引的主分片数量以最大化读取性能。
    slices => 4
    # 每次批量获取的文档数量。
    size => 1000
    # 滚动查询的存活时间,确保长任务不会因超时而中断。
    scroll => "5m"
  }
}

output {
  opensearch {
    # 【必填】目标 PolarSearch 集群的连接地址。
    hosts => ["https://polarsearch-cluster-endpoint:9200"]
    # 【必填】目标集群的认证凭据。
    user => "your_target_user"
    password => "your_target_password"

    # --- 安全配置 ---
    ssl => true
    # 如果目标集群使用自签名证书,取消注释并指定 CA 证书路径。
    # cacert => "/path/to/polarsearch_ca.crt"

    # --- 索引与文档 ID ---
    # 从元数据中动态读取并设置索引名,以保持与源端一致。
    # 注意:元数据路径因 input 插件而异,需通过调试输出确认。
    index => "%{[@metadata][input][opensearch][_index]}"
    # 从元数据中读取并设置文档 ID,以保持文档的唯一性。
    document_id => "%{[@metadata][input][opensearch][_id]}"
  }

  # --- 调试输出(可选) ---
  # 在开发和测试阶段,取消此段注释可在控制台打印数据流信息。
  # 正式同步时,注释掉此段以获得最佳性能。
  # stdout {
  #   codec => rubydebug {
  #     metadata => true
  #   }
  # }
}

参考写入速率

场景

参考写入速率

多索引、6分片、8pipeline worker、ECS 816 GB、文档平均260字节。

18.67 MB/s, 37883 docs/s。

场景二:增量数据同步

此配置通过定时调度,周期性地拉取源库中在指定时间窗口内新增或更新的数据。

重要
  • 准备工作:请确保已完成步骤二:(可选)为源数据准备时间戳字段的配置。

  • 启动时机:建议先完成一次全量数据迁移,再启动增量同步任务。并将增量同步的时间查询起点对齐到全量迁移完成的时刻,避免历史数据遗漏或新旧数据交叉覆盖。

  • 调度配置:LogStash 的schedule调度是串行执行的。请确保单次同步任务的执行耗时小于调度周期,否则可能导致部分时间窗口的数据未被采集,造成数据遗漏。建议将query的时间范围设置得略大于调度周期(例如,每5分钟调度一次,查询近6分钟的数据),以利用幂等写入特性来覆盖边界情况。

配置文件示例

# synchronization.conf
# 一个从 Elasticsearch/OpenSearch 向 PolarSearch 增量同步数据的完整配置示例。
# 通过定时执行带时间范围过滤的查询,每个调度周期仅拉取指定时间窗口内
# 新写入或更新的文档,实现增量数据同步。

input {
  # 如果源集群是 Elasticsearch,请将 'opensearch' 替换为 'elasticsearch'。
  elasticsearch {
    # 【必填】源集群的连接地址。
    hosts => ["https://source-cluster-endpoint:9200"]
    # 【必填】源集群的认证凭据。
    user     => "your_source_user"
    password => "your_source_password"

    # 【必填】指定需要增量同步的索引,支持通配符。
    index => "your-business-logs-*"

    # --- 增量同步核心配置 ---
    # 时间范围查询:仅拉取最近 5 分钟内写入的文档。
    # gte: 大于等于当前时间减去 5 分钟。
    # lte: 小于等于当前分钟整点(now/m 表示向下取整到分钟,用于避免拉取尚未写入完毕的数据)。
    # 时间窗口大小需与 schedule 调度周期保持匹配,防止两次调度之间出现数据空档。
    query => '{"query":{"range":{"@timestamp":{"gte":"now-5m","lte":"now/m"}}}}'

    # 定时调度:使用 Cron 语法定义执行周期,以下配置为每分钟执行一次。
    # 常见示例:
    #   "* * * * *"     每分钟执行
    #   "*/5 * * * *"   每 5 分钟执行
    #   "0 * * * *"     每小时整点执行
    schedule => "* * * * *"

    # 滚动查询的存活时间,确保分批拉取时上下文不会超时中断。
    scroll => "5m"
    # 开启此选项以获取原始索引名和文档 ID。
    docinfo => true
    # 设置并发读取数,建议设置为源索引的主分片数量以最大化读取性能。
    slices => 4
    # 每次批量获取的文档数量。
    size => 5000

    # --- 安全配置 ---
    # ssl => true
    # cacert => "/path/to/source_ca.crt"
  }
}

# filter(可选): 移除 pipeline 附加的@timestamp字段
# filter {
#   mutate {
#     remove_field => ["@timestamp"]
#   }
# }

output {
  opensearch {
    # 【必填】目标 PolarSearch 集群的连接地址。
    hosts    => ["https://polarsearch-cluster-endpoint:9200"]
    # 【必填】目标集群的认证凭据。
    user     => "your_target_user"
    password => "your_target_password"

    # --- 安全配置 ---
    ssl => true
    # cacert => "/path/to/polarsearch_ca.crt"

    # 从元数据中动态读取索引名,保持与源端一致。
    # 注意:元数据路径因 input 插件而异,需通过调试输出确认。
    index       => "%{[@metadata][input][elasticsearch][_index]}"
    # 从元数据中读取文档 ID,实现幂等写入:同一文档重复同步时执行更新而非重复插入。
    document_id => "%{[@metadata][input][elasticsearch][_id]}"
  }

  # --- 调试输出(可选) ---
  # stdout {
  #   codec => rubydebug {
  #     metadata => true
  #   }
  # }
}

步骤四:执行同步与验证

  1. 启动任务:在LogStash根目录下执行以下命令启动同步任务。

    # 启动同步任务,-f 参数指定配置文件
    ./bin/logstash -f synchronization.conf

    对于大规模数据同步,建议使用nohupLogStash作为后台服务运行:

    nohup ./bin/logstash -f synchronization.conf &
  2. 等待任务执行完成后,登录目标PolarSearch,检查索引和数据是否成功创建和写入。

常见问题

如何确认元数据(如索引名)的正确路径?

output配置中启用stdout调试输出。启动LogStash后,控制台会打印每条文档的详细信息,包括[@metadata]对象。从中找到_index_id的确切层级结构,并修改output配置中的路径。

示例调试输出

{
    // ... 其他字段
    "@metadata": {
        "input" => {
          "opensearch" => {
              "_id": "some-document-id",
              "_index": "your-business-logs-2023.01.01",
              // ... 其他元数据
          }
        }
    }
}

根据以上输出,正确的索引名路径是%{[@metadata][input][opensearch][_index]}

重要

不同input插件的元数据路径存在差异。使用logstash-input-elasticsearch时,路径中的opensearch应替换为elasticsearch,即%{[@metadata][input][elasticsearch][_index]}

同步速度很慢,如何优化?

  1. 增加slices选项:确保input配置中的slices参数值等于源索引的主分片数。

  2. 增加pipeline.workers:启动LogStash时,使用--pipeline.workers参数增加并发处理线程数,例如 ./bin/logstash -f synchronization.conf --pipeline.workers 8。此值通常可设置为服务器的CPU核心数。

  3. 增加LogStashJVM堆大小:编辑config/jvm.options文件,增大-Xms-Xmx的值,例如-Xms4g -Xmx4g

如何确保目标库的索引配置(如分词器、字段类型)与源端一致?

同步前,从源库导出索引的MappingSettings,然后在目标PolarSearch中手动创建对应的索引模板。最后,在LogStashoutput配置中设置manage_template => false,确保LogStash不会覆盖预设的模板。

增量同步时,如何避免遗漏文档更新?

增量同步的时间过滤基于文档的时间戳字段,能否捕获更新操作取决于更新方式:

  • 若业务侧使用POST /_update进行局部更新,该操作不会触发Ingest Pipeline,@timestamp不会刷新,因此增量同步无法感知此类更新。

  • 若业务侧使用完整的index(全量覆盖写入)方式更新文档,则会触发Pipeline并刷新@timestamp,增量同步可以捕获。

如果您的业务更新操作以_update为主,建议在业务层面由写入方主动维护一个随每次变更刷新的时间戳字段,并以此作为增量同步的过滤依据。

增量同步任务意外中断后,如何恢复?

增量同步基于时间窗口查询,无需额外的状态持久化。重新启动LogStash 后,任务将按照query中配置的时间范围(如now-5m)重新拉取最近一个时间窗口内的数据并继续执行,无需手动干预。如需从某一历史时间点补同步数据,可临时调整query 中的时间范围(例如将now-5m改为now-2h)启动一次性补偿同步,完成后再恢复原配置。