【发布时间】:2017-02-06 19:36:22
【问题描述】:
我有一个脚本,它从 MySQL 表中检索活动“作业”列表,然后使用多处理库为每个活动作业实例化一次我的主脚本。我的多处理脚本具有检查给定作业是否已被另一个线程占用的功能。它通过检查 DB 表中的特定列是否为/不是 NULL 来做到这一点。 DB 查询返回单个项目元组:
def check_if_job_claimed():
#...
job_claimed = cursor.fetchone() #Returns (claim_id,) for claimed jobs, and (None,) for unclaimed jobs
if job_claimed:
print "This job has already been claimed by another thread."
return
else:
do_stuff_to_claim_the_job
当我在没有多处理部分的情况下运行此功能时,声明检查工作得很好。但是当我尝试并行运行作业时,声明检查将所有 (None,) 元组读取为具有价值并因此具有真实性,因此该函数假定该工作已被声明。
我已尝试调整多处理器使用的并发进程数,但声明检查仍然不起作用...即使我将进程数设置为 1。我也尝试过使用 if 语句看看我能不能让它这样工作:
if job_claimed == True
if job_claimed == (None,)
# etc.
不过运气不好。
是否有人知道多处理库中的某些内容会阻止我的声明检查函数正确解释 job_claimed 元组?也许我的代码有问题?
编辑
我在调试模式下对 job_claimed 变量运行了一些真实性测试。以下是这些测试的结果:
(pdb) job_claimed
(None,)
(pdb) len(job_claimed)
1
(pdb) job_claimed == True
False
(pdb) job_claimed == False
False
(pdb) job_claimed[0]
None
(pdb) job_claimed[0] == True
False
(pdb) job_claimed[0] == False
False
(pdb) any(job_claimed)
False
(pdb) all(job_claimed)
False
(pdb) job_claimed is not True
True
(pdb) job_claimed is not False
True
编辑
根据要求:
with open('Resource_File.txt', 'r') as f:
creds = eval(f.read())
connection = mysql.connector.connect(user=creds["mysql_user"],password=creds["mysql_pw"],host=creds["mysql_host"],database=creds["mysql_db"],use_pure=False,buffered=True)
def check_if_job_claimed(job_id):
cursor = connection.cursor()
thread_id_query = "SELECT Thread_Id FROM jobs WHERE Job_ID=\'{}\';".format(job_id)
cursor.execute(thread_id_query)
job_claimed = cursor.fetchone()
job_claimed = job_claimed[0]
if job_claimed:
print "This job has already been claimed by another thread. Moving on to next job..."
cursor.close()
return False
else:
thread_id = socket.gethostname()+':'+str(random.randint(0,1000))
claim_job = "UPDATE jobs SET Thread_Id = \'{}\' WHERE Job_ID = \'{}\';".format(job_id)
cursor.execute(claim_job)
connection.commit()
print "Job is now claimed"
cursor.close()
return True
def call_the_queen(dict_of_job_attributes):
if check_if_job_claimed(dict_of_job_attributes['job_id']):
instance = OM(dict_of_job_attributes) #<-- Create instance of my target class
instance.queen_bee()
#multiprocessing code
import multiprocessing as mp
if __name__ == '__main__':
active_jobs = get_active_jobs()
pool = mp.Pool(processes = 4)
pool.map(call_the_queen,active_jobs)
pool.close()
pool.join()
【问题讨论】:
-
这个怎么样 - 不要做这种复杂的笨拙的事情,而是将所有作业 ID 放入一个队列(例如 Redis 中的列表),然后只需简单地
pop()一次一个作业 ID。这是一个原子操作,因此当 worker 检索到作业 ID 时,没有其他进程可以窃取它。 -
能否包含多处理代码,以及创建光标的代码。我想您在进程中重用光标对象,并且只有 1 个项目
-
是的,那些真实性测试没有用,这是每个 python 程序的预期结果。
标签: python mysql boolean multiprocessing