You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Node.js中监听Elasticsearch新增文档并完成数据拉取

Elasticsearch Node.js 监听指定索引新增文档实现方案

以下提供两种可直接落地的实现方案,可根据你使用的Elasticsearch版本选择:


方案1:时间戳轮询方案(兼容所有ES版本)

实现逻辑简单无额外依赖,所有ES版本都可以使用,原理是基于文档创建时间定期拉取上一次查询之后新增的文档。

  • 第一步:安装官方ES客户端
    npm install @elastic/elasticsearch
  • 第二步:业务代码实现
const { Client } = require('@elastic/elasticsearch')
// 初始化ES客户端
const esClient = new Client({
  node: 'http://你的ES服务地址:9200',
  // 有身份认证时补充以下配置
  // auth: { username: 'elastic', password: '你的ES密码' }
})

// 自定义配置
const MONITOR_INDEX = 'assets' // 要监听的索引名
const POLL_INTERVAL = 1000 // 轮询间隔,单位毫秒,可根据延迟要求调整
// 记录上次拉取的最大时间戳,初始值可按需求调整为更早时间
let lastFetchTimestamp = Date.now()

// 拉取新增文档方法
async function fetchNewDocs() {
  try {
    const res = await esClient.search({
      index: MONITOR_INDEX,
      query: {
        range: {
          // 字段替换为你实际维护的文档创建时间字段,比如@timestamp
          create_time: {
            gt: lastFetchTimestamp
          }
        }
      },
      sort: [{ create_time: 'asc' }] // 按创建时间正序,保证数据顺序正确
    })

    const newDocs = res.hits.hits
    if (newDocs.length > 0) {
      // 更新最新时间戳为最后一条新增文档的创建时间
      const lastDocTime = newDocs[newDocs.length - 1]._source.create_time
      lastFetchTimestamp = lastDocTime
      // 此处编写你的业务逻辑,处理拿到的新增文档
      console.log('获取到新增文档:', newDocs.map(item => item._source))
    }
  } catch (err) {
    console.error('拉取ES新增文档失败:', err)
  }
}

// 启动轮询任务
setInterval(fetchNewDocs, POLL_INTERVAL)

如果你没有手动维护文档的创建时间字段,可以配置ES的Ingest Pipeline自动为新增文档添加@timestamp字段,直接用这个字段做轮询的范围查询条件即可。

该方案优点是兼容性强、逻辑可控,缺点是存在轮询间隔级别的延迟。


方案2:ES原生Change API监听(ES 7.16+版本支持)

ES 7.16及以上版本原生提供了变更监听接口,不需要自己维护轮询逻辑,延迟更低,是高版本ES的首选方案。

const { Client } = require('@elastic/elasticsearch')
const esClient = new Client({
  node: 'http://你的ES服务地址:9200',
  // auth: { username: 'elastic', password: '你的ES密码' }
})

const MONITOR_INDEX = 'assets' // 要监听的索引名

async function startChangeListener() {
  try {
    // 初始化变更流,仅监听文档创建操作
    const changeStream = esClient.changes({
      index: MONITOR_INDEX,
      filters: {
        operation_types: ['create'] // 只监听创建事件,去掉该配置可同时监听更新、删除等所有变更
      }
    })

    // 有新增文档时触发该事件
    changeStream.on('data', (change) => {
      const newDoc = change.document
      // 此处编写你的业务逻辑,处理拿到的新增文档
      console.log('监听到新增文档:', newDoc)
    })

    // 监听错误事件,实现自动重连
    changeStream.on('error', (err) => {
      console.error('ES监听异常:', err)
      setTimeout(startChangeListener, 3000)
    })
  } catch (err) {
    console.error('启动ES监听失败:', err)
    setTimeout(startChangeListener, 3000)
  }
}

// 启动监听
startChangeListener()

该方案优点是延迟低、不需要自己维护轮询逻辑,缺点是对ES版本有要求。


通用注意事项

  • 如果业务对数据准确性要求很高,建议在处理完文档后记录已处理的文档ID或者时间戳,避免服务重启、重连时重复处理数据
  • 要监听多个索引时,直接把索引名配置为逗号分隔的多个索引名或者通配符即可

内容的提问来源于stack exchange,提问作者Nothing

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.27 07:06:08