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

Node.js中SQL函数顺序执行异常求助:DuckDB/Sqlite3数据读写问题

问题:Node.js中SQL函数无法按预期顺序执行

我是JavaScript新手,尝试用Node.js结合DuckDB和Sqlite3时,遇到SQL函数顺序执行异常的问题。以下是我的代码data2.js:

// data2.js
const fs = require('fs');
// const sql = require('sqlite3');
const sql = require('duckdb');
const csvParser = require('csv-parser');
const stream = require('stream');


// CSV data as a string
const csvData = `Wind speed (m/s),Output power (kW)
0,0
1,0
2,0
3,0
4,80
5,140
6,360
7,610
8,1000
9,1470
10,1900
11,2320
12,2690
13,2850
14,2950
15,3000
16,3000
17,3000
18,3000
19,3000
20,3000
21,3000
22,3000
23,3000
24,3000
25,3000`;


// Function to create a database, table, and insert data
function createAndStoreData() {
  console.log('createAndStoreData()');
  const db = new sql.Database('power_curve.db');

  db.serialize(() => {
    db.run('DROP TABLE IF EXISTS power_curve');
    db.run('CREATE TABLE IF NOT EXISTS power_curve (ws DECIMAL PRIMARY KEY, power DECIMAL)');

    const stmt = db.prepare('INSERT INTO power_curve VALUES (?, ?)');

    const dataStream = new stream.Readable();
    dataStream.push(csvData);
    dataStream.push(null);

    dataStream
      .pipe(csvParser())
      .on('data', (row) => {
        stmt.run(row['Wind speed (m/s)'], row['Output power (kW)']);
      })
      .on('end', () => {
        stmt.finalize();
        db.close((err) => {
          if (err) {
            return console.error('Error closing database:', err.message);
          }
          console.log('Data stored in power_curve.db successfully.');
        });
      })
      .on('error', (error) => {
        console.error('Error reading CSV data:', error.message);
      });
  });
}


// Function to read data from the database
function readData() {
  console.log('readData()');
  const db = new sql.Database('power_curve.db');

  db.serialize(() => {
    db.all('SELECT * FROM power_curve', (err, rows) => {
      if (err) {
        console.error('Error querying data:', err.message);
      } else {
        console.table(rows);
        console.log('Data read from power_curve.db successfully.');
      }

      db.close((err) => {
        if (err) {
          console.error('Error closing database:', err.message);
        }
      });
    });
  });
}


if (require.main === module) {
  createAndStoreData();
  readData();
}

首次运行代码node data2.js,输出如下:

createAndStoreData()
readData()
Error querying data: Connection Error: Connection was never established or has been closed already
Error closing database: Database was already closed
Data stored in power_curve.db successfully.

第二次运行时,输出如下:

createAndStoreData()
readData()
Error closing database: Database was already closed
┌─────────┬────┬───────┐
│ (index) │ ws │ power │
├─────────┼────┼───────┤
│    0    │ 0  │   0   │
│    1    │ 1  │   0   │
│    2    │ 2  │   0   │
│    3    │ 3  │   0   │
│    4    │ 4  │  80   │
│    5    │ 5  │  140  │
│    6    │ 6  │  360  │
│    7    │ 7  │  610  │
│    8    │ 8  │ 1000  │
│    9    │ 9  │ 1470  │
│   10    │ 10 │ 1900  │
│   11    │ 11 │ 2320  │
│   12    │ 12 │ 2690  │
│   13    │ 13 │ 2850  │
│   14    │ 14 │ 2950  │
│   15    │ 15 │ 3000  │
│   16    │ 16 │ 3000  │
│   17    │ 17 │ 3000  │
│   18    │ 18 │ 3000  │
│   19    │ 19 │ 3000  │
│   20    │ 20 │ 3000  │
│   21    │ 21 │ 3000  │
│   22    │ 22 │ 3000  │
│   23    │ 23 │ 3000  │
│   24    │ 24 │ 3000  │
│   25    │ 25 │ 3000  │
└─────────┴────┴───────┘
Data read from power_curve.db successfully.

请求解决SQL函数顺序执行异常的问题。


解决方案

问题核心是Node.js的异步特性:createAndStoreData里的数据库操作、CSV流处理都是异步任务,你直接在调用它后立即执行readData,此时数据还未写入完成,数据库连接可能还未就绪,导致第一次运行读取失败;第二次能读到数据是因为之前已经完成了写入,但仍有连接关闭的报错。

要让函数按顺序执行,需要通过Promise或回调控制异步流程,以下是两种可行方案:

方案1:用Promise+async/await控制顺序

将异步操作包装为Promise,通过async/await确保前一步完成后再执行下一步:

// 修改createAndStoreData为Promise版本
function createAndStoreData() {
  return new Promise((resolve, reject) => {
    console.log('createAndStoreData()');
    const db = new sql.Database('power_curve.db');

    db.serialize(() => {
      db.run('DROP TABLE IF EXISTS power_curve');
      db.run('CREATE TABLE IF NOT EXISTS power_curve (ws DECIMAL PRIMARY KEY, power DECIMAL)');

      const stmt = db.prepare('INSERT INTO power_curve VALUES (?, ?)');

      const dataStream = new stream.Readable();
      dataStream.push(csvData);
      dataStream.push(null);

      dataStream
        .pipe(csvParser())
        .on('data', (row) => {
          stmt.run(row['Wind speed (m/s)'], row['Output power (kW)']);
        })
        .on('end', () => {
          stmt.finalize();
          db.close((err) => {
            if (err) {
              console.error('Error closing database:', err.message);
              reject(err);
              return;
            }
            console.log('Data stored in power_curve.db successfully.');
            resolve();
          });
        })
        .on('error', (error) => {
          console.error('Error reading CSV data:', error.message);
          reject(error);
        });
    });
  });
}

// 修改readData为Promise版本
function readData() {
  return new Promise((resolve, reject) => {
    console.log('readData()');
    const db = new sql.Database('power_curve.db');

    db.serialize(() => {
      db.all('SELECT * FROM power_curve', (err, rows) => {
        if (err) {
          console.error('Error querying data:', err.message);
          reject(err);
        } else {
          console.table(rows);
          console.log('Data read from power_curve.db successfully.');
          resolve(rows);
        }

        db.close((err) => {
          if (err) {
            console.error('Error closing database:', err.message);
            reject(err);
          }
        });
      });
    });
  });
}

// 用async/await控制执行顺序
if (require.main === module) {
  async function run() {
    try {
      await createAndStoreData();
      await readData();
    } catch (err) {
      console.error('执行出错:', err);
    }
  }
  run();
}

方案2:用回调函数传递执行信号

给createAndStoreData添加回调参数,在数据写入完成后触发readData:

// 修改createAndStoreData,添加回调参数
function createAndStoreData(callback) {
  console.log('createAndStoreData()');
  const db = new sql.Database('power_curve.db');

  db.serialize(() => {
    db.run('DROP TABLE IF EXISTS power_curve');
    db.run('CREATE TABLE IF NOT EXISTS power_curve (ws DECIMAL PRIMARY KEY, power DECIMAL)');

    const stmt = db.prepare('INSERT INTO power_curve VALUES (?, ?)');

    const dataStream = new stream.Readable();
    dataStream.push(csvData);
    dataStream.push(null);

    dataStream
      .pipe(csvParser())
      .on('data', (row) => {
        stmt.run(row['Wind speed (m/s)'], row['Output power (kW)']);
      })
      .on('end', () => {
        stmt.finalize();
        db.close((err) => {
          if (err) {
            return console.error('Error closing database:', err.message);
          }
          console.log('Data stored in power_curve.db successfully.');
          // 写入完成后调用回调执行readData
          callback();
        });
      })
      .on('error', (error) => {
        console.error('Error reading CSV data:', error.message);
        callback(error);
      });
  });
}

// 修改主执行逻辑,通过回调触发readData
if (require.main === module) {
  createAndStoreData(() => {
    readData();
  });
}

额外说明

  • 关于Error closing database: Database was already closed报错:DuckDB的serialize方法可能会自动管理连接关闭,可尝试去掉手动的db.close调用,或确保只执行一次关闭操作。
  • 两种方案的核心都是等待异步操作完成后再执行后续步骤,这是Node.js异步编程的基础逻辑。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 09:40:59