【发布时间】:2015-06-14 22:03:32
【问题描述】:
我正在编写一个 Python 2.7 脚本,该脚本从 MySQL 表中检索行,遍历数据以对其进行处理,然后应该按此顺序执行以下操作:
更新我们之前获得的表行以设置锁定值 在每一行中为 TRUE
在 UPDATE 查询通过 MySQLdb 执行并提交之后,线程池应该在原始循环中的数据上运行。
实际发生的是 UPDATE 查询似乎在 ThreadPool 完成后以某种方式提交。我尝试将其重构为 try/finally 语句以确保,但现在它要么之后仍然这样做,要么只是不提交 UPDATE 并运行 ThreadPool。
可以肯定的是,这是一个令人头疼的问题。我认为我只是在做一些非常错误且明显的事情,但在看了这么久之后并没有抓住它。非常感谢任何输入!
这里是要点:
from multiprocessing.pool import ThreadPool, IMapIterator
import MySQLdb as mdb
import os, sys, time
import re
from boto.s3.connection import S3Connection
from boto.s3.bucket import Bucket
...
con = mdb.connect('localhost', 'user', 'pass', 'db')
with con:
cur = con.cursor()
cur.execute("SELECT preview_queue.filename, preview_queue.product_id, preview_queue.track, products.name, preview_queue.id FROM preview_queue join `catalog_module-products` AS products on products.id = preview_queue.product_id where locked != 1")
rows = cur.fetchall()
mp3s_to_download = []
lock_ids = []
last_directory = ""
if len(rows) > 0:
for row in rows:
base_dir = str(get_base_dir(row[1], row[3]))
mp3s_to_download.append([base_dir, str(row[0])])
if last_directory != "preview_temp/"+base_dir:
if not os.path.exists("preview_temp/"+base_dir):
try:
os.makedirs("preview_temp/"+base_dir)
except OSError, e:
pass
last_directory = "preview_temp/"+base_dir
lock_ids.append(str(row[4]))
if len(lock_ids) > 0:
action_ids = ','.join(lock_ids)
try:
cur.execute("UPDATE preview_queue SET locked = 1 WHERE id IN ({})".format(action_ids))
con.commit()
finally:
pool = ThreadPool(processes=20)
pool.map(download_file, mp3s_to_download)
cur.close()
【问题讨论】:
标签: python python-2.7 threadpool mysql-python