(input_prefix, output_prefix, num_workers)
| 184 | ) |
| 185 | |
| 186 | def make_binary_alignment_dataset(input_prefix, output_prefix, num_workers): |
| 187 | nseq = [0] |
| 188 | |
| 189 | def merge_result(worker_result): |
| 190 | nseq[0] += worker_result["nseq"] |
| 191 | |
| 192 | input_file = input_prefix |
| 193 | offsets = Binarizer.find_offsets(input_file, num_workers) |
| 194 | pool = None |
| 195 | if num_workers > 1: |
| 196 | pool = Pool(processes=num_workers - 1) |
| 197 | for worker_id in range(1, num_workers): |
| 198 | prefix = "{}{}".format(output_prefix, worker_id) |
| 199 | pool.apply_async( |
| 200 | binarize_alignments, |
| 201 | ( |
| 202 | args, |
| 203 | input_file, |
| 204 | utils.parse_alignment, |
| 205 | prefix, |
| 206 | offsets[worker_id], |
| 207 | offsets[worker_id + 1], |
| 208 | ), |
| 209 | callback=merge_result, |
| 210 | ) |
| 211 | pool.close() |
| 212 | |
| 213 | ds = indexed_dataset.make_builder( |
| 214 | dataset_dest_file(args, output_prefix, None, "bin"), impl=args.dataset_impl |
| 215 | ) |
| 216 | |
| 217 | merge_result( |
| 218 | Binarizer.binarize_alignments( |
| 219 | input_file, |
| 220 | utils.parse_alignment, |
| 221 | lambda t: ds.add_item(t), |
| 222 | offset=0, |
| 223 | end=offsets[1], |
| 224 | ) |
| 225 | ) |
| 226 | if num_workers > 1: |
| 227 | pool.join() |
| 228 | for worker_id in range(1, num_workers): |
| 229 | prefix = "{}{}".format(output_prefix, worker_id) |
| 230 | temp_file_path = dataset_dest_prefix(args, prefix, None) |
| 231 | ds.merge_file_(temp_file_path) |
| 232 | os.remove(indexed_dataset.data_file_path(temp_file_path)) |
| 233 | os.remove(indexed_dataset.index_file_path(temp_file_path)) |
| 234 | |
| 235 | ds.finalize(dataset_dest_file(args, output_prefix, None, "idx")) |
| 236 | |
| 237 | logger.info("[alignments] {}: parsed {} alignments".format(input_file, nseq[0])) |
| 238 | |
| 239 | def make_dataset(vocab, input_prefix, output_prefix, lang, num_workers=1): |
| 240 | if args.dataset_impl == "raw": |
no test coverage detected