-
Notifications
You must be signed in to change notification settings - Fork 1
/
Copy pathreprocess.py
executable file
·160 lines (133 loc) · 4.74 KB
/
reprocess.py
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
#!/usr/bin/env python3
# Runs the pipeline with the provided arguments across all deliveries. Useful
# when there's a new stage, or a stage needs to be re-run.
#
# Usage: ./reprocess.py \
# --max-jobs <N> --log-prefix <LP> [--sample-level] -- arguments-for-run
#
# Examples:
# ./reprocess.py \
# --max-jobs 12 --log-prefix rl -- --stages readlengths
#
# You can pass --deliveries A,B,C to run on only a subset of deliveries.
#
# ./reprocess.py \
# --deliveries PRJNA729801 --max-jobs 12 --log-prefix rl \
# -- --stages readlengths
#
# And if you run --sample-level it will invoke run.py for each sample instead
# of for each delivery, allowing more parallelism but also more overhead.
#
# ./reprocess.py \
# --deliveries PRJNA729801 --sample-level --max-jobs 12 \
# --log-prefix rl -- --stages readlengths
import os
import re
import sys
import random
import datetime
import argparse
import subprocess
from concurrent.futures import ThreadPoolExecutor
log_date = datetime.datetime.now().date().isoformat()
log_dir = "log"
if not os.path.exists(log_dir):
os.mkdir(log_dir)
regular_deliveries = os.listdir("deliveries")
restricted_deliveries = []
restricted_dir = os.path.join("..", "mgs-restricted")
if os.path.exists(restricted_dir):
restricted_deliveries = os.listdir(
os.path.join(restricted_dir, "deliveries")
)
def prepare_job(delivery, log_prefix, sample, run_args):
logfile = "%s/%s.%s.%s" % (log_dir, log_date, log_prefix, delivery)
if sample:
logfile = "%s.%s" % (logfile, sample)
run_args = list(run_args) + ["--sample", sample]
return logfile, ["./run.py", "--delivery", delivery, *run_args]
def run_job(job):
logfile, cmd = job
with open(logfile, "w") as outf:
result = subprocess.run(cmd, stdout=outf, stderr=subprocess.STDOUT)
if result.returncode != 0:
outf.write("ERROR: %s\n" % (result.returncode))
def get_sample_priority(sample):
m = re.findall(r"L00\d$", sample)
if not m:
return "A"
m, = m
return m
def parallelize(config, deliveries, run_args):
job_queue = []
for delivery in deliveries:
args = run_args[:]
if delivery in restricted_deliveries:
args.append("--restricted")
root_dir = restricted_dir
elif delivery in regular_deliveries:
root_dir = "."
else:
raise Exception("Unknown delivery %r" % delivery)
if config.sample_level:
prioritized_samples = []
with open(os.path.join(root_dir, "deliveries", delivery,
"metadata", "metadata.tsv")) as inf:
for line in inf:
sample = line.strip().split()[0]
prioritized_samples.append(
(get_sample_priority(sample), sample))
prioritized_samples.sort()
for priority, sample in prioritized_samples:
job_queue.append(prepare_job(
delivery, config.log_prefix, sample, args))
else:
job_queue.append(prepare_job(
delivery, config.log_prefix, None, args))
if config.shuffle:
random.shuffle(job_queue)
with ThreadPoolExecutor(max_workers=config.max_jobs) as executor:
for job in job_queue:
executor.submit(run_job, job)
def start():
argv = sys.argv[1:]
if "--" not in argv:
raise Exception("Use -- to separate arguments to ./run.py.")
our_args = argv[: argv.index("--")]
run_args = argv[argv.index("--") + 1 :]
parser = argparse.ArgumentParser()
parser.add_argument(
"--max-jobs",
metavar="N",
type=int,
required=True,
help="maximum number of jobs to run at once",
)
parser.add_argument(
"--log-prefix",
required=True,
help="Log prefix, for storing this run under log/",
)
parser.add_argument(
"--deliveries",
help="The IDs of the delivery to process, comma separated",
)
parser.add_argument(
"--sample-level",
action="store_true",
help="Parallelize at the sample level instead of the delivery level")
parser.add_argument(
"--shuffle",
action="store_true",
help="Run jobs in random order. Allows greater parallelism if "
"inputs vary dramatically in size")
config = parser.parse_args(our_args)
if config.deliveries:
deliveries = config.deliveries.split(",")
else:
deliveries = regular_deliveries + restricted_deliveries
subprocess.check_call(["./prepare-shm-kraken.sh"])
subprocess.check_call(["./prepare-shm-bowtie.sh"])
parallelize(config, deliveries, run_args)
if __name__ == "__main__":
start()