如何在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
相关产品推荐
相关产品推荐

