如何使用fs模块与bulk请求读取多个JSON文件导入Elasticsearch不同索引
多JSON文件批量导入Elasticsearch不同索引实现方案
将原有单文件导入逻辑抽象为可复用函数,通过配置映射关联不同JSON文件和对应的ES索引,避免重复代码,实现逻辑如下:
实现步骤
- 首先定义配置列表,统一维护3个JSON文件的路径、对应索引名、类型、字段转换规则
- 封装通用的单文件导入异步函数,接收配置参数即可完成单个文件到对应索引的批量导入
- 可选择并发或串行执行三个导入任务,根据自身服务器性能调整即可
注意事项
- 如果使用的是Elasticsearch 7.x及以上版本,官方已废弃
_type配置,可直接删除代码中对应的_type字段 - 批量插入的单次批次大小(示例中为1000)可根据你的服务器性能调整,避免单次请求过大超时
- 示例使用
fs.promises的Promise风格API替代原有回调,逻辑更清晰,避免回调嵌套问题
完整可运行代码
const fs = require('fs').promises; // 保留你原有初始化好的ES客户端即可 const client = require('./你的ES客户端初始化文件路径'); // 第一步:定义3个文件对应的配置,按需修改参数即可 const importConfigs = [ { filePath: './DocRes.json', // 第一个JSON文件路径 indexName: 'ncar_index', // 对应ES索引名 indexType: 'ncar', // ES7+版本可删除该字段 // 字段转换规则,和你原有逻辑保持一致 transform: (obj) => ({ id: obj.id, name: obj.name, summary: obj.summary, image: obj.image, approvetool: obj.approvetool, num: obj.num, date: obj.date, }) }, { filePath: './第二个文件.json', indexName: '第二个索引名', indexType: '第二个类型', // ES7+版本可删除 transform: (obj) => ({ // 按需填写第二个文件的字段转换规则 id: obj.id, // 其他自定义字段... }) }, { filePath: './第三个文件.json', indexName: '第三个索引名', indexType: '第三个类型', // ES7+版本可删除 transform: (obj) => ({ // 按需填写第三个文件的字段转换规则 id: obj.id, // 其他自定义字段... }) } ]; // 第二步:封装通用单文件导入异步函数 async function importJsonToEs(config) { const { filePath, indexName, indexType, transform } = config; // 读取JSON文件内容 const data = await fs.readFile(filePath, { encoding: 'utf-8' }); // 构建bulk请求体 let bulkRequest = data.split('\n').reduce((acc, line) => { let obj; try { obj = JSON.parse(line); } catch (e) { console.log(`${filePath} 文件读取完成`); return acc; } const formattedData = transform(obj); const indexAction = { index: { _index: indexName, _id: formattedData.id } }; if (indexType) indexAction.index._type = indexType; acc.push(indexAction); acc.push(formattedData); return acc; }, []); // 批量插入逻辑,和原有逻辑兼容 return new Promise((resolve, reject) => { let busy = false; const callback = (err, resp) => { if (err) { console.error(`${indexName} 导入出错:`, err); busy = false; return reject(err); } busy = false; }; const perhapsInsert = () => { if (!busy) { busy = true; client.bulk({ body: bulkRequest.slice(0, 1000) }, callback); bulkRequest = bulkRequest.slice(1000); console.log(`${indexName} 剩余待插入条数:`, bulkRequest.length); } if (bulkRequest.length > 0) { setTimeout(perhapsInsert, 100); } else { console.log(`${indexName} 所有记录插入完成`); resolve(); } }; perhapsInsert(); }); } // 第三步:执行三个导入任务(并发版本) async function runAllImports() { try { // 并发执行三个导入任务,若需要串行执行可参考下方的串行实现 await Promise.all(importConfigs.map(config => importJsonToEs(config))); console.log('所有文件全部导入完成'); } catch (err) { console.error('导入过程出现错误:', err); } } // 启动导入 runAllImports();
可选调整:串行导入
如果需要避免并发导入占用过多IO/网络资源,可以改成串行执行,替换上面的runAllImports函数即可:
async function runAllImports() { try { for (const config of importConfigs) { await importJsonToEs(config); } console.log('所有文件全部导入完成'); } catch (err) { console.error('导入过程出现错误:', err); } }
内容的提问来源于stack exchange,提问作者Drsaud
相关产品推荐
相关产品推荐

