多线程同步访问SQLite查询的问题排查与技术求助
多线程股票数据处理与SQLite连接问题排查
问题进展
- 初始状态:主进程维护股票列表,创建4个线程分发数据处理,出现线程未同步和重复数据库连接问题。已使用QMutex锁和工厂模式避免重复连接名,但问题未解决。
- 第一次修改:采纳@user4581301建议,将Mutex改为引用传递,线程同步问题修复,但重复连接名问题仍存在,推测与工厂模式实现有关。
- 重构后新问题:简化重构代码后,成功创建4个WAL模式数据库连接,但出现
database is locked Unable to fetch row错误,寻求技术帮助。
相关代码
初始版本
main.cpp
QThreadPool threadPool; threadPool.setMaxThreadCount(4); QMutex mutex; while (stockListQuery.next()) { QString name_ = stockListQuery.value("name").toString(); mutex.lock(); testRunable* query = new testRunable(db, name_); threadPool.start(query); mutex.unlock(); } threadPool.waitForDone();
.h
class testRunable : public QRunnable { public: testRunable(QSqlDatabase db, const QString &name); void run() override; private: QSqlDatabase db_; QMutex m_mutex; QString name_; };
.cc
testRunable::testRunable(QSqlDatabase db, const QString &name): db_(db), name_(name) { } void testRunable::run() { QString ids = QString::number((int)QThread::currentThread()); QSqlDatabase db = DBFactory(ids).getDatabase(); int iterations = 0; while (1) { m_mutex.lock(); qDebug() << name_ << " " << db.connectionName(); m_mutex.unlock(); if (iterations >= 1) { break; // stop the loop and exit the thread } iterations++; QThread::sleep(1); // seconds } }
修改Mutex后版本
main.cpp
QThreadPool threadPool; threadPool.setMaxThreadCount(4); QMutex mutex; // Submit the query task to the thread pool while (stockListQuery.next()) { QString name_ = stockListQuery.value("name").toString(); testRunable* query = new testRunable(db, name_, mutex); // ref of mutex should be passed in the function threadPool.start(query); } threadPool.waitForDone();
.h
class DBFactory { public: QSqlDatabase getDatabase(); void setConnetName(const QString& connectName); private: QString conName_; QMap<QString, QSqlDatabase> m_; }; class testRunable : public QRunnable { public: testRunable(QSqlDatabase db, const QString &name, QMutex& mutex); void run() override; private: DBFactory dbFac_; QSqlDatabase db_; QMutex& m_mutex; QString name_; };
.cpp
QSqlDatabase DBFactory::getDatabase() { if (m_.contains(conName_)) { return m_[conName_]; } // Create a new database connection object QSqlDatabase db = QSqlDatabase::addDatabase("QSQLITE", conName_); db.setDatabaseName("STOCKLIST.db"); if (!db.isOpen()) { db.open(); } QSqlQuery query(db); query.exec("PRAGMA journal_mode=WAL;"); if (query.next()) { QString mode = query.value(0).toString(); //qDebug() << "Current SQLite mode is:" << mode; } // save to map m_[conName_] = db; return db; } void DBFactory::setConnetName(const QString& connectName) { conName_ = connectName; } testRunable::testRunable(QSqlDatabase db, const QString &name, QMutex& mutex): db_(db), name_(name), m_mutex(mutex) { } void testRunable::run() { QString ids = QString::number((int)QThread::currentThread()); dbFac_.setConnetName(ids); QSqlDatabase db = dbFac_.getDatabase(); m_mutex.lock(); qDebug() << name_ << " " << db.connectionName(); stockCrawler(name_); saveToTable("result_table", db); m_mutex.unlock(); }
重构后版本
main.cpp
... QSharedPointer<sqliteConnectionFactory> factory(new sqliteConnectionFactory("STOCKLIST.db")); QThreadPool threadPool; threadPool.setMaxThreadCount(4); auto start = std::chrono::steady_clock::now(); // Submit the query task to the thread pool while (stockListQuery.next()) { QString name = stockListQuery.value("name").toString(); QString symbol = stockListQuery.value("symbol").toString(); testRunable* query = new testRunable(name, symbol, factory); threadPool.start(query); } threadPool.waitForDone();
.h
class sqliteConnectionFactory { public: sqliteConnectionFactory(const QString &dbName); QSqlDatabase createConnection(); private: QString dbName_; static QMap<QString, QSqlDatabase> connMap_; }; class testRunable : public QRunnable { public: testRunable(const QString& name, const QString& symbol, const QSharedPointer<sqliteConnectionFactory>& factory); void run() override; private: static QMutex mutex_; QString name_; QString symbol_; QSharedPointer<sqliteConnectionFactory> factory_; };
.cpp
sqliteConnectionFactory::sqliteConnectionFactory(const QString& dbName): dbName_(dbName) { } QSqlDatabase sqliteConnectionFactory::createConnection() { QString connectname = QString::number((int)(QThread::currentThread())); if (!connMap_.contains(connectname)) { QSqlDatabase db = QSqlDatabase::addDatabase("QSQLITE", connectname); db.setDatabaseName(dbName_); if (!db.open()) { qDebug() << db.lastError().text(); assert(0); } QSqlQuery setting(db); setting.exec("PRAGMA journal_mode=WAL;"); if (setting.next()) { QString mode = setting.value(0).toString(); qDebug() << "Create a new connection: " << connectname; qDebug() << "Current SQLite mode is: " << mode; } connMap_.insert(connectname, db); } return connMap_.value(connectname); } QMap<QString, QSqlDatabase> sqliteConnectionFactory::connMap_; QMutex testRunable::mutex_; testRunable::testRunable(const QString &name, const QString &symbol, const QSharedPointer<sqliteConnectionFactory>& factory): name_(name), symbol_(symbol), factory_(factory) { } void testRunable::run() { QSqlDatabase db = factory_->createConnection(); top_crawler* stockCrawler_ = new top_crawler(); stockCrawler_->InitCrawler(symbol_.toStdString()); std::vector<Spot> spots_ = stockCrawler_->GetSpots(); std::vector<double> ma5_ = stockCrawler_->GetMA5(); std::vector<double> ma10_ = stockCrawler_->GetMA10(); std::vector<QDateTime> series_ = stockCrawler_->GetDateTime(); std::vector<double> close_ = stockCrawler_->GetClosef(); auto latestDayPrice = stockCrawler_->GetLatestClosePrice(); int status = -1; int cross = 0; CrossList CL; QSqlQuery query3(db); for (int i = 0; i < close_.size(); ++i) { QVector<QString> queryBuffer; if (ma5_[i] >= ma10_[i] && status != 1) { if (status == 0) { cross++; CL.addSpot(close_[i], series_[i], i); queryBuffer.push_back(symbol_); queryBuffer.push_back(series_[i].toString("yyyy-MM-dd")); queryBuffer.push_back(QString::number(close_[i])); } status = 1; } else if (ma5_[i] < ma10_[i] && status != 0) { if (status == 1) { cross++; CL.addSpot(close_[i], series_[i], i); queryBuffer.push_back(symbol_); queryBuffer.push_back(series_[i].toString("yyyy-MM-dd")); queryBuffer.push_back(QString::number(close_[i])); } status = 0; } else { // error } if (!queryBuffer.empty()) { QMutexLocker locker(&mutex_); query3.prepare("INSERT INTO cross_point (stock_id, date, price) VALUES (?, ?, ?)"); for (int i = 0; i < 3; ++i) { query3.bindValue(i, queryBuffer[i]); } if (!query3.exec()) { qWarning() << "Failed to insert data into table"; qWarning() << query3.lastError().text(); } } } qWarning() << "cross:" << cross; CL.printList(); }
内容的提问来源于stack exchange,提问作者Weimin Chan
相关产品推荐
相关产品推荐

