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

异步SQLite语句执行封装函数问题:建表后查询失败

问题:创建并填充SQLite表后立即查询失败(clean=true时必现)

在创建并填充SQLite表后立即执行查询,会出现表不存在的错误,推测是建表和插入操作未完成就执行了后续查询。

现象

  • 当createDatabaseFromSql的clean参数为false时,多次尝试偶尔能成功:
$ node ./setup.js 
CREATE database
Executing: DROP TABLE IF EXISTS `accounts`
Executing: CREATE TABLE IF NOT EXISTS `accounts` ( `id` INTEGER NOT NULL PRIMARY KEY AUTOINCREMENT, `username` VARCHAR(50) NOT NULL, `password` VARCHAR(255) NOT NULL, `email` VARCHAR(100) NOT NULL )
Executing: INSERT INTO `accounts` (`username`, `password`, `email`) VALUES ('admin', 'admin', 'admin@admin.com'), ('test', 'test', 'test@test.com')
FETCH test account
Connected to the 'auth' database.
Connected to the 'auth' database.
true
Close the 'auth' database connection.
Close the 'auth' database connection.
  • 当clean参数为true(强制删除现有数据库重建)时,查询必失败,报错SQLITE_ERROR: no such table: accounts:
$ node ./setup.js 
CREATE database
Executing: DROP TABLE IF EXISTS `accounts`
Executing: CREATE TABLE IF NOT EXISTS `accounts` ( `id` INTEGER NOT NULL PRIMARY KEY AUTOINCREMENT, `username` VARCHAR(50) NOT NULL, `password` VARCHAR(255) NOT NULL, `email` VARCHAR(100) NOT NULL )
Executing: INSERT INTO `accounts` (`username`, `password`, `email`) VALUES ('admin', 'admin', 'admin@admin.com'), ('test', 'test', 'test@test.com')
FETCH test account
Connected to the 'auth' database.
Connected to the 'auth' database.
node:internal/process/esm_loader:97
    internalBinding('errors').triggerUncaughtException(
                              ^

[Error: SQLITE_ERROR: no such table: accounts] {
  errno: 1,
  code: 'SQLITE_ERROR'
}

Node.js v18.14.0

相关代码文件

setup.js

import { createDatabaseFromSql, executeQuery } from "./query-utils.js";

console.log("CREATE database");
await createDatabaseFromSql("auth", "./setup.sql", true);

console.log("FETCH test account");
const row = await executeQuery(
  "auth",
  "SELECT * FROM `accounts` where `username` = ?",
  ["test"]
);
console.log(row?.id === 2);

query-utils.js

import fs from "fs";
import sqlite3 from "sqlite3";

const SQL_DEBUG = true;

const loadSql = (filename, delimiter = ";") =>
  fs
    .readFileSync(filename)
    .toString()
    .replace(/(\r\n|\n|\r)/gm, " ")
    .replace(/\s+/g, " ")
    .split(delimiter)
    .map((statement) => statement.trim())
    .filter((statement) => statement.length);

const executeSerializedQueries = async (databaseName, callback) => {
  let db;
  try {
    db = new sqlite3.Database(`./${databaseName}.db`, (err) => {
      if (err) console.error(err.message);
      console.log(`Connected to the '${databaseName}' database.`);
    });
    db.serialize(() => {
      callback(db);
    });
  } catch (e) {
    throw Error(e);
  } finally {
    if (db) {
      db.close((err) => {
        if (err) console.error(err.message);
        console.log(`Close the '${databaseName}' database connection.`);
      });
    }
  }
};

const createDatabaseFromSql = async (databaseName, sqlFilename, clean) =>
  new Promise((resolve, reject) => {
    if (clean) {
      fs.rmSync(`./${databaseName}.db`, { force: true }); // Remove existing
    }
    try {
      executeSerializedQueries(databaseName, (db) => {
        loadSql(sqlFilename).forEach((statement) => {
          if (SQL_DEBUG) {
            console.log("Executing:", statement);
          }
          db.run(statement);
        });
        resolve();
      });
    } catch (e) {
      reject(e);
    }
  });

const executeQuery = async (databaseName, query, params = []) =>
  new Promise((resolve, reject) => {
    try {
      executeSerializedQueries(databaseName, (db) => {
        db.get(query, params, (error, row) => {
          if (error) reject(error);
          else resolve(row);
        });
      });
    } catch (e) {
      reject(e);
    }
  });

const executeQueryAll = async (databaseName, query, params = []) =>
  new Promise((resolve, reject) => {
    try {
      executeSerializedQueries(databaseName, (db) => {
        db.all(query, params, (error, rows) => {
          if (error) reject(error);
          else resolve(rows);
        });
      });
    } catch (e) {
      reject(e);
    }
  });

export {
  createDatabaseFromSql,
  executeSerializedQueries,
  executeQuery,
  executeQueryAll,
  loadSql,
};

setup.sql

DROP TABLE IF EXISTS `accounts`;

CREATE TABLE IF NOT EXISTS `accounts` (
  `id` INTEGER NOT NULL PRIMARY KEY AUTOINCREMENT,
  `username` VARCHAR(50) NOT NULL,
  `password` VARCHAR(255) NOT NULL,
  `email` VARCHAR(100) NOT NULL
);

INSERT INTO `accounts` (`username`, `password`, `email`)
  VALUES ('admin', 'admin', 'admin@admin.com'),
    ('test', 'test', 'test@test.com');

package.json

{
  "dependencies": {
    "sqlite3": "^5.1.4"
  }
}

问题原因

核心问题是没有正确等待SQLite异步操作完成:

  1. createDatabaseFromSql中调用db.run(statement)后立即resolve(),但db.run是异步执行的,此时SQL语句可能还没真正执行完毕。
  2. executeSerializedQueries中,db.serialize只是保证回调内的操作按顺序执行,但不会等待所有操作完成后再继续后续代码。当clean=true时,数据库文件刚被删除重建,异步操作的延迟更明显,导致后续查询时表还未创建。

修复方案

需要修改executeSerializedQueries和createDatabaseFromSql,确保所有SQL操作完成后再resolve Promise:

修改后的query-utils.js

import fs from "fs";
import sqlite3 from "sqlite3";

const SQL_DEBUG = true;

const loadSql = (filename, delimiter = ";") =>
  fs
    .readFileSync(filename)
    .toString()
    .replace(/(\r\n|\n|\r)/gm, " ")
    .replace(/\s+/g, " ")
    .split(delimiter)
    .map((statement) => statement.trim())
    .filter((statement) => statement.length);

// 重构:让executeSerializedQueries返回Promise,等待所有操作完成
const executeSerializedQueries = async (databaseName, callback) => {
  return new Promise((resolve, reject) => {
    let db;
    try {
      db = new sqlite3.Database(`./${databaseName}.db`, (err) => {
        if (err) {
          console.error(err.message);
          reject(err);
          return;
        }
        console.log(`Connected to the '${databaseName}' database.`);
        
        db.serialize(() => {
          // 执行回调,回调需要返回Promise或处理完所有异步操作后通知
          const done = callback(db);
          if (done instanceof Promise) {
            done.then(() => {
              closeDb(db).then(resolve).catch(reject);
            }).catch(reject);
          } else {
            // 如果回调是同步的,直接关闭数据库
            closeDb(db).then(resolve).catch(reject);
          }
        });
      });
    } catch (e) {
      reject(e);
      if (db) closeDb(db).catch(console.error);
    }
  });
};

// 封装关闭数据库的Promise
const closeDb = (db) => {
  return new Promise((resolve, reject) => {
    db.close((err) => {
      if (err) {
        console.error(err.message);
        reject(err);
      } else {
        const dbName = db.filename.split('/').pop().replace('.db','');
        console.log(`Close the '${dbName}' database connection.`);
        resolve();
      }
    });
  });
};

const createDatabaseFromSql = async (databaseName, sqlFilename, clean) => {
  if (clean) {
    fs.rmSync(`./${databaseName}.db`, { force: true }); // Remove existing
  }
  
  return executeSerializedQueries(databaseName, async (db) => {
    // 遍历SQL语句,逐个等待执行完成
    for (const statement of loadSql(sqlFilename)) {
      if (SQL_DEBUG) {
        console.log("Executing:", statement);
      }
      await new Promise((resolve, reject) => {
        db.run(statement, (err) => {
          if (err) reject(err);
          else resolve();
        });
      });
    }
  });
};

const executeQuery = async (databaseName, query, params = []) => {
  return executeSerializedQueries(databaseName, (db) => {
    return new Promise((resolve, reject) => {
      db.get(query, params, (error, row) => {
        if (error) reject(error);
        else resolve(row);
      });
    });
  });
};

const executeQueryAll = async (databaseName, query, params = []) => {
  return executeSerializedQueries(databaseName, (db) => {
    return new Promise((resolve, reject) => {
      db.all(query, params, (error, rows) => {
        if (error) reject(error);
        else resolve(rows);
      });
    });
  });
};

export {
  createDatabaseFromSql,
  executeSerializedQueries,
  executeQuery,
  executeQueryAll,
  loadSql,
};

修复要点

  1. executeSerializedQueries改为返回Promise,确保数据库连接、操作执行、关闭的流程都被正确等待。
  2. createDatabaseFromSql中逐个等待db.run完成:使用for...of循环配合await,确保每条SQL语句执行完成后再执行下一条,而不是批量触发后直接resolve。
  3. 封装closeDb为Promise:避免关闭数据库的异步操作被忽略。
  4. db.serialize结合Promise:保证操作顺序的同时,等待所有异步操作完成。

修改后,无论clean参数是true还是false,都能确保建表和插入操作完成后再执行后续查询,不会出现表不存在的错误。


内容的提问来源于stack exchange,提问作者Mr. Polywhirl

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 09:45:14