异步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异步操作完成:
createDatabaseFromSql中调用db.run(statement)后立即resolve(),但db.run是异步执行的,此时SQL语句可能还没真正执行完毕。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, };
修复要点
executeSerializedQueries改为返回Promise,确保数据库连接、操作执行、关闭的流程都被正确等待。createDatabaseFromSql中逐个等待db.run完成:使用for...of循环配合await,确保每条SQL语句执行完成后再执行下一条,而不是批量触发后直接resolve。- 封装
closeDb为Promise:避免关闭数据库的异步操作被忽略。 db.serialize结合Promise:保证操作顺序的同时,等待所有异步操作完成。
修改后,无论clean参数是true还是false,都能确保建表和插入操作完成后再执行后续查询,不会出现表不存在的错误。
内容的提问来源于stack exchange,提问作者Mr. Polywhirl
相关产品推荐
相关产品推荐

