import DBService from '../db/DBService' import logger from '../logger/Logger' import BaseProvider from '../core/BaseProvider' const ROW_LIMIT = 2000 export default class MysqlProvider extends BaseProvider { async connectDataSource() { let dbService = new DBService() let inputs = this.params.inputs for (let i = 0; i < inputs.length; i++) { const input = inputs[i]; let { topic, table, conditions, fields } = input // 首先查询该条件下的数据量大小 let countSql = `SELECT COUNT(*) as num FROM ${table} WHERE ${conditions}` logger.info(countSql) let { rows, err } = await dbService.query(countSql, []) if (!err) { let num = rows[0]["num"] // 每次查询200条数据,需要查询的次数 let queryCount = (num % ROW_LIMIT) > 0 ? Math.floor(num / ROW_LIMIT) + 1 : (Math.floor(num / ROW_LIMIT)) logger.info(`总共有${num}条数据, 需要查询${queryCount}次`) let offset = 0 for (let i = 0; i < queryCount; i++) { let sql = `SELECT ${fields} FROM ${table} WHERE ${conditions} ORDER BY createTimeStamp LIMIT ${ROW_LIMIT} OFFSET ${offset}` logger.info(sql) let { rows, err } = await dbService.query(sql, []) if (err) { logger.info(err) break } for (let i = 0; i < rows.length; i++) { const row = rows[i]; this.onMessage(topic, row) } offset += ROW_LIMIT } } else { logger.info(err) } } console.log("输入结束") this.core.endInput() } }