| FazBrowse GitHub Viewer | Trending | | Home |
| Tools: [Download Repo ZIP] [Original HTTPS Page] |
1 parent d3d6840 commit 2096991
5 files changed
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -17,5 +17,5 @@ | |||
| 17 | 17 | @pytest.fixture | |
| 18 | 18 | def workload_params(request): | |
| 19 | 19 | params = request.param | |
| 20 | - files_names = [f"fio-go_storage_fio.0.{i}" for i in range(0, params.num_processes)] | ||
| 20 | + files_names = [f"fio-go_storage_fio.0.{i}" for i in range(0, params.num_files)] | ||
| 21 | 21 | return params, files_names | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -80,10 +80,10 @@ def _get_params() -> Dict[str, List[TimeBasedReadParameters]]: | |||
| 80 | 80 | chunk_size_bytes = chunk_size_kib * 1024 | |
| 81 | 81 | bucket_name = bucket_map[bucket_type] | |
| 82 | 82 | ||
| 83 | - num_files = num_processes * num_coros | ||
| 83 | + num_files = num_processes | ||
| 84 | 84 | ||
| 85 | 85 | # Create a descriptive name for the parameter set | |
| 86 | - name = f"{pattern}_{bucket_type}_{num_processes}p_{file_size_mib}MiB_{chunk_size_kib}KiB_{num_ranges_val}ranges" | ||
| 86 | + name = f"{pattern}_{bucket_type}_{num_processes}p_{num_coros}c_{file_size_mib}MiB_{chunk_size_kib}KiB_{num_ranges_val}ranges" | ||
| 87 | 87 | ||
| 88 | 88 | params[workload_name].append( | |
| 89 | 89 | TimeBasedReadParameters( | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -20,9 +20,10 @@ workload: | |||
| 20 | 20 | ||
| 21 | 21 | - name: "read_rand_multi_process" | |
| 22 | 22 | pattern: "rand" | |
| 23 | - coros: [1] | ||
| 23 | + coros: [1, 16] | ||
| 24 | 24 | processes: [1] | |
| 25 | 25 | ||
| 26 | + | ||
| 26 | 27 | defaults: | |
| 27 | 28 | DEFAULT_RAPID_ZONAL_BUCKET: "chandrasiri-benchmarks-zb" | |
| 28 | 29 | DEFAULT_STANDARD_BUCKET: "chandrasiri-benchmarks-rb" | |
| Original file line number | Diff line number | Diff line change | |
|---|---|---|---|
@@ -115,47 +115,51 @@ def _download_time_based_json(client, filename, params): | |||
| 115 | 115 | ||
| 116 | 116 | ||
| 117 | 117 | async def _download_time_based_async(client, filename, params): | |
| 118 | - total_bytes_downloaded = 0 | ||
| 119 | - | ||
| 120 | 118 | mrd = AsyncMultiRangeDownloader(client, params.bucket_name, filename) | |
| 121 | 119 | await mrd.open() | |
| 122 | 120 | ||
| 123 | - offset = 0 | ||
| 124 | - is_warming_up = True | ||
| 125 | - start_time = time.monotonic() | ||
| 126 | - warmup_end_time = start_time + params.warmup_duration | ||
| 127 | - test_end_time = warmup_end_time + params.duration | ||
| 128 | - | ||
| 129 | - while time.monotonic() < test_end_time: | ||
| 130 | - current_time = time.monotonic() | ||
| 131 | - if is_warming_up and current_time >= warmup_end_time: | ||
| 132 | - is_warming_up = False | ||
| 133 | - total_bytes_downloaded = 0 # Reset counter after warmup | ||
| 134 | - | ||
| 135 | - ranges = [] | ||
| 136 | - if params.pattern == "rand": | ||
| 137 | - for _ in range(params.num_ranges): | ||
| 138 | - offset = random.randint( | ||
| 139 | - 0, params.file_size_bytes - params.chunk_size_bytes | ||
| 140 | - ) | ||
| 141 | - ranges.append((offset, params.chunk_size_bytes, BytesIO())) | ||
| 142 | - else: # seq | ||
| 143 | - for _ in range(params.num_ranges): | ||
| 144 | - ranges.append((offset, params.chunk_size_bytes, BytesIO())) | ||
| 145 | - offset += params.chunk_size_bytes | ||
| 146 | - if offset + params.chunk_size_bytes > params.file_size_bytes: | ||
| 147 | - offset = 0 # Reset offset if end of file is reached | ||
| 148 | - | ||
| 149 | - await mrd.download_ranges(ranges) | ||
| 150 | - | ||
| 151 | - bytes_in_buffers = sum(r[2].getbuffer().nbytes for r in ranges) | ||
| 152 | - assert bytes_in_buffers == params.chunk_size_bytes * params.num_ranges | ||
| 153 | - | ||
| 154 | - if not is_warming_up: | ||
| 155 | - total_bytes_downloaded += params.chunk_size_bytes * params.num_ranges | ||
| 121 | + async def _worker_coro(): | ||
| 122 | + total_bytes_downloaded = 0 | ||
| 123 | + offset = 0 | ||
| 124 | + is_warming_up = True | ||
| 125 | + start_time = time.monotonic() | ||
| 126 | + warmup_end_time = start_time + params.warmup_duration | ||
| 127 | + test_end_time = warmup_end_time + params.duration | ||
| 128 | + | ||
| 129 | + while time.monotonic() < test_end_time: | ||
| 130 | + current_time = time.monotonic() | ||
| 131 | + if is_warming_up and current_time >= warmup_end_time: | ||
| 132 | + is_warming_up = False | ||
| 133 | + total_bytes_downloaded = 0 # Reset counter after warmup | ||
| 134 | + | ||
| 135 | + ranges = [] | ||
| 136 | + if params.pattern == "rand": | ||
| 137 | + for _ in range(params.num_ranges): | ||
| 138 | + offset = random.randint( | ||
| 139 | + 0, params.file_size_bytes - params.chunk_size_bytes | ||
| 140 | + ) | ||
| 141 | + ranges.append((offset, params.chunk_size_bytes, BytesIO())) | ||
| 142 | + else: # seq | ||
| 143 | + for _ in range(params.num_ranges): | ||
| 144 | + ranges.append((offset, params.chunk_size_bytes, BytesIO())) | ||
| 145 | + offset += params.chunk_size_bytes | ||
| 146 | + if offset + params.chunk_size_bytes > params.file_size_bytes: | ||
| 147 | + offset = 0 # Reset offset if end of file is reached | ||
| 148 | + | ||
| 149 | + await mrd.download_ranges(ranges) | ||
| 150 | + | ||
| 151 | + bytes_in_buffers = sum(r[2].getbuffer().nbytes for r in ranges) | ||
| 152 | + assert bytes_in_buffers == params.chunk_size_bytes * params.num_ranges | ||
| 153 | + | ||
| 154 | + if not is_warming_up: | ||
| 155 | + total_bytes_downloaded += params.chunk_size_bytes * params.num_ranges | ||
| 156 | + return total_bytes_downloaded | ||
| 157 | + | ||
| 158 | + tasks = [asyncio.create_task(_worker_coro()) for _ in range(params.num_coros)] | ||
| 159 | + results = await asyncio.gather(*tasks) | ||
| 156 | 160 | ||
| 157 | 161 | await mrd.close() | |
| 158 | - return total_bytes_downloaded | ||
| 162 | + return sum(results) | ||
| 159 | 163 | ||
| 160 | 164 | ||
| 161 | 165 | def _download_files_worker(process_idx, filename, params, bucket_type): | |
| Back | FazBrowse Home | New Git URL |
0 commit comments