(
dataset,
index_queue,
data_queue,
done_event,
transform,
collate,
batch_size,
seed,
worker_id,
num_workers,
datakind,
parallel_stream,
)
| 836 | |
| 837 | |
| 838 | def _worker_loop( |
| 839 | dataset, |
| 840 | index_queue, |
| 841 | data_queue, |
| 842 | done_event, |
| 843 | transform, |
| 844 | collate, |
| 845 | batch_size, |
| 846 | seed, |
| 847 | worker_id, |
| 848 | num_workers, |
| 849 | datakind, |
| 850 | parallel_stream, |
| 851 | ): |
| 852 | _set_worker_signal_handlers() |
| 853 | random.seed(seed) |
| 854 | np.random.seed(seed) |
| 855 | watchdog = ManagerWatchdog() |
| 856 | iteration_end = False |
| 857 | fetcher = map_fetcher |
| 858 | if datakind == "stream": |
| 859 | global _worker_info |
| 860 | _worker_info = WorkerInfo(idx=worker_id, worker=num_workers, seed=seed) |
| 861 | dataset = iter(dataset) |
| 862 | fetcher = stream_fetcher |
| 863 | |
| 864 | while watchdog.is_alive(): |
| 865 | try: |
| 866 | r = index_queue.get(timeout=GLOBAL_TIMEOUT) |
| 867 | except queue.Empty: |
| 868 | continue |
| 869 | if r is None: |
| 870 | assert done_event.is_set() or iteration_end |
| 871 | break |
| 872 | elif done_event.is_set() or iteration_end: |
| 873 | continue |
| 874 | |
| 875 | idx, place_holder = r |
| 876 | try: |
| 877 | if data_monitor: |
| 878 | data, pid, dataset_time, transform_time, collate_time = fetcher( |
| 879 | dataset, |
| 880 | place_holder, |
| 881 | transform, |
| 882 | collate, |
| 883 | parallel_stream, |
| 884 | batch_size, |
| 885 | ) |
| 886 | |
| 887 | if idx < monitor_num_workers.value: |
| 888 | if monitor_workers[idx] == 0: |
| 889 | monitor_workers[idx] = pid |
| 890 | worker_idx = idx % monitor_num_workers.value |
| 891 | if ( |
| 892 | pid == monitor_workers[worker_idx] |
| 893 | and idx % data_monitor_frequency == 0 |
| 894 | ): |
| 895 | print( |
nothing calls this directly
no test coverage detected