【发布时间】:2020-10-26 05:03:53
【问题描述】:
我正在尝试以并行方式运行作业,但问题是某些作业依赖于其他作业。如果我同时触发它们,那么相关作业将失败。所以我在想依赖的工作将等到父母的工作完全完成。有没有办法做到这一点?
例子:
Job Job_Type
A independent
B independent
C A
D B
在上面的示例中,您可以看到作业 A 和 B 是独立的,因此它们将同时运行,但 C 和 D 分别依赖于 A 和 B。所以C在A完成后立即启动,D在B完成后立即启动。我们可以做这样的过程吗?
在我的代码中,我是基于序列号的类别-
Job Job_Type seq_num
A independent 1
B independent 1
C A 2
D B 2
第一个 seq_num = 1 将触发,seq_num 1 中的所有作业将并行运行。 1 seq_num 2 完成后将启动。所以假设工作 A 需要 10 分钟。 B 需要 15 分钟,所以 seq_num 2 将在 15 分钟后开始但我希望进程 Job C 将立即开始 Job A 的完成意味着 10 分钟后 Job C 将开始,15 分钟后 Job D 将开始。请帮我实现这个逻辑。
我正在使用的代码:
def parallel_Execution():
logging.info("parallel_Execution..................[started]")
par_temp_loc = '/medaff/Temp/'
'''Reading the metadata file and creating as a dataframe'''
df = pd.read_csv(par_temp_loc+'metadata_file_imedical.txt', delimiter='|',error_bad_lines=False)
uni_master_job = df['Master Job Name'].unique().tolist()
print(uni_master_job)
'''getting unique execution sequence'''
logging.info("Getting the unique Execution Sequence Number!")
unique_exec_seq = df['Execution Sequence'].unique().tolist()
unique_exec_seq.sort()
print(unique_exec_seq)
num_unique_seq = len(unique_exec_seq)
logging.info("Total Number of unique sequence Number : %2d" %(num_unique_seq))
p2 = ThreadWithReturnValue(target = partial(parallel_temp2, unique_exec_seq, df ))
p2.start()
r2 = p2.join()
print(r2)
#r1 = r1.append(r2)
mail_df(r2)
'''Parallel Processing Function'''
def parallel_temp2(unique_exec_seq, df):
list_df = []
df_main4 = pd.DataFrame()
for exec_seq in unique_exec_seq:
seq_num = exec_seq
temp_df = df[df['Execution Sequence'] == exec_seq].copy()
unique_master_job = temp_df['Master Job Name'].unique().tolist()
print(unique_master_job)
#logging.info("%s Master Job Started." %(unique_master_job))
if(len(unique_master_job)>0):
num_processes = len(unique_master_job)
pool = ThreadPool(processes=num_processes)
result1 = pool.map(partial(parallel_view_creation, exec_seq, temp_df), unique_master_job)
pool.close()
pool.join()
df_main = pd.DataFrame(result1)
#print("printing df_main")
#print(df_main)
for m_job in df_main.master_job.unique():
temp_df1 = df_main[df_main['master_job'] == m_job]
status = temp_df1.status.unique()[0]
if(status == 0):
unique_master_job.remove(m_job)
pool = ThreadPool(processes=num_processes)
result2 = pool.map(partial(parallel_build_query, exec_seq, temp_df), unique_master_job)
pool.close()
pool.join()
if(result2):
df_main2 = pd.DataFrame(result2)
df_main3 = pd.concat([df_main,df_main2], sort = False)
status_df_list = df_main3['status'].unique().tolist()
print(status_df_list)
if(0 in status_df_list):
break
if(0 in status_df_list):
break
else:
df_main4 = df_main4.append(df_main3)
if(0 in status_df_list):
df_main4 = df_main4.append(df_main3)
return df_main4
如果您需要任何信息,请告诉我。
【问题讨论】:
标签: python multithreading pyspark multiprocessing threadpool