使用Tornado时两个并发Worker访问MySQL读不到数据求助
问题:Tornado线程池内Worker无法读取到其他Worker插入的MySQL记录
我编写了一个包含两个Worker的简单程序:Worker1负责插入记录,Worker2读取这些记录。但程序执行时Worker2始终读取到0条记录,单独通过CLI运行两个Worker则正常,疑似问题出在Tornado上,求解决方案。
import time import munch from tornado import concurrent from tornado.ioloop import IOLoop import logging import pymysql config = munch.munchify({"mysql_host": "127.0.0.1", "mysql_port": 3306, "mysql_user": "", "mysql_password": "", "mysql_db": ""}) executor = concurrent.futures.ThreadPoolExecutor(max_workers=8) class Worker1(): CONST_TYPE_ID = 'worker1' def __init__(self): self.conn = pymysql.connect(host=config.mysql_host, port=config.mysql_port, user=config.mysql_user, passwd=config.mysql_password, db=config.mysql_db, charset='UTF8MB4', local_infile=True) self.cursor = self.conn.cursor(pymysql.cursors.DictCursor) def single_run(self): print("Single run of worker:", self.CONST_TYPE_ID) query = "insert into readwrite (timestamp) values (now())" self.conn.ping(True) self.cursor.execute(query) self.conn.commit() class Worker2(): CONST_TYPE_ID = 'worker2' def __init__(self): self.conn = pymysql.connect(host=config.mysql_host, port=config.mysql_port, user=config.mysql_user, passwd=config.mysql_password, db=config.mysql_db, charset='UTF8MB4', local_infile=True) self.cursor = self.conn.cursor(pymysql.cursors.DictCursor) def single_run(self): print("Single run of worker:", self.CONST_TYPE_ID) query = "select * from readwrite" self.conn.ping(True) self.cursor.execute(query) for row in self.cursor.fetchall(): print(row) def run_worker(worker): i = 0 instance = worker() while True: try: i += 1 print("Starting {name} - run {iter}".format(name=instance.CONST_TYPE_ID, iter=i)) instance.single_run() except Exception as e: logging.exception('Worker {} got error {!r}, errno is {}'.format(instance.CONST_TYPE_ID, e, e.args[0])) print("Waiting...") time.sleep(1) def main(): workers = [Worker1, Worker2] executor.map(run_worker, workers) IOLoop.current().start() if __name__ == '__main__': main()
问题原因
核心问题和Tornado本身无关,是MySQL事务隔离机制和pymysql连接的事务状态导致的:
- MySQL默认隔离级别是
REPEATABLE READ,在一个事务内,查询只能看到事务启动时的数据快照。 - pymysql会在第一次执行查询时隐式开启事务,Worker2初始化后一直复用同一个连接,这个连接的事务从未提交或回滚,所以它始终只能看到连接创建时的数据库状态,无法读取Worker1后续插入的新数据。
- 单独CLI运行时,每个Worker都是全新的连接,每次查询都在新事务中执行,因此能读到最新数据。
解决方案
方案1:Worker2每次查询后提交事务
修改Worker2的single_run方法,在查询完成后提交事务,刷新连接的事务视图:
def single_run(self): print("Single run of worker:", self.CONST_TYPE_ID) query = "select * from readwrite" self.conn.ping(True) self.cursor.execute(query) for row in self.cursor.fetchall(): print(row) # 提交事务,让连接的事务状态更新 self.conn.commit()
方案2:给Worker2的连接开启自动提交
在Worker2初始化数据库连接时,添加autocommit=True参数,让每次查询自动提交事务:
def __init__(self): self.conn = pymysql.connect(host=config.mysql_host, port=config.mysql_port, user=config.mysql_user, passwd=config.mysql_password, db=config.mysql_db, charset='UTF8MB4', local_infile=True, autocommit=True) # 开启自动提交 self.cursor = self.conn.cursor(pymysql.cursors.DictCursor)
方案3:每次查询前回滚旧事务
如果不想开启自动提交,可在Worker2每次查询前回滚旧事务,强制开启新的事务上下文:
def single_run(self): print("Single run of worker:", self.CONST_TYPE_ID) # 回滚旧事务,重置事务状态 self.conn.rollback() query = "select * from readwrite" self.conn.ping(True) self.cursor.execute(query) for row in self.cursor.fetchall(): print(row)
内容的提问来源于stack exchange,提问作者hyp3rv1p3r
相关产品推荐
相关产品推荐

