如何返回InfluxDB fluxObserver收集的queryRows查询结果
问题原因
queryApi.queryRows 是异步回调风格的API,采用观察者模式按行返回查询结果:next 每拿到一行触发一次,所有行返回完才触发complete,查询出错触发error。这几个回调的执行时机远晚于外层同步代码,在complete里写return、或者直接接queryRows的返回值,本质是用同步逻辑写异步代码,根本拿不到结果。
正确实现方案
用Promise把整个查询逻辑包裹起来,在闭包里维护存储结果的数组,行数据全部收集完成后通过resolve返回结果,出错时通过reject抛出异常即可。
封装通用查询方法
/** * 封装InfluxDB行查询逻辑,返回结构化结果数组 * @param {string} fluxQuery Flux查询语句 * @returns {Promise<Array<Object>>} 查询结果对象数组 */ function runFluxQuery(fluxQuery) { return new Promise((resolve, reject) => { // 闭包内存储查询结果 const resultList = [] const observer = { next(row, tableMeta) { // 单行数据转对象后存入结果集 const rowData = tableMeta.toObject(row) resultList.push(rowData) }, error(err) { // 查询异常直接抛出 reject(err) }, complete() { // 所有数据接收完成,返回结果 resolve(resultList) } } queryApi.queryRows(fluxQuery, observer) }) }
调用方法
因为返回的是Promise,推荐用async/await写法拿结果,逻辑更清晰:
// 示例:在异步函数中调用 async function getSensorData() { try { const myQuery = `from(bucket:"sensor_data") |> range(start: -24h) |> filter(fn: (r) => r._measurement == "temperature")` const data = await runFluxQuery(myQuery) // 这里拿到的data就是完整的查询结果数组,后续业务逻辑直接写在这里即可 console.log(`查询成功,共获取${data.length}条温度数据`) return data } catch (err) { console.error('查询失败:', err) // 按需处理错误逻辑 return [] } } // 执行查询 getSensorData()
如果不习惯async/await,也可以用.then()链式调用:
const myQuery = `from(bucket:"sensor_data") |> range(start: -24h) |> filter(fn: (r) => r._measurement == "temperature")` runFluxQuery(myQuery) .then(data => { console.log('查询结果:', data) }) .catch(err => { console.error('查询出错:', err) })
注意事项
- 不要试图在同步代码路径里直接拿异步查询的结果,所有依赖查询返回值的逻辑,必须写在
await之后或者.then()的回调里,否则只会拿到undefined - 如果数据量很大,可以按需在
next回调里做数据过滤、格式转换,不用等所有数据进数组再处理,减少内存占用
内容的提问来源于stack exchange,提问作者MPA
相关产品推荐
相关产品推荐

