Init
This commit is contained in:
commit
c394677d07
59 files changed
+5149
No files matched your search
@@ -0,0 +1,41 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""Демо: spawn — каждый процесс заново импортирует модуль.
|
||||
|
||||
Запуск из корня курса:
|
||||
python lab04_multiprocessing/examples/01_spawn_demo.py
|
||||
|
||||
Показывает, что на Windows/spawn дочерний процесс заново выполняет
|
||||
модуль верхнего уровня. Поэтому защита __main__ обязательна.
|
||||
Здесь же видно время старта процесса.
|
||||
"""
|
||||
import multiprocessing
|
||||
import os
|
||||
import sys
|
||||
import time
|
||||
|
||||
print(f"[импорт модуля] pid={os.getpid()}, parent={os.getppid()}, "
|
||||
f"метод={multiprocessing.get_start_method()}")
|
||||
|
||||
|
||||
def worker(x):
|
||||
return x * x
|
||||
|
||||
|
||||
def main():
|
||||
print(f"\n[главный процесс] pid={os.getpid()}")
|
||||
ctx = multiprocessing.get_context() # spawn на Windows и macOS
|
||||
|
||||
t0 = time.perf_counter()
|
||||
with ctx.Pool(2) as pool:
|
||||
results = pool.map(worker, range(8))
|
||||
dt = time.perf_counter() - t0
|
||||
|
||||
print(f"pool.map -> {results}")
|
||||
print(f"время (включая старт 2 процессов): {dt * 1000:.1f} мс")
|
||||
print("\nОбратите внимание: строка [импорт модуля] могла напечататься")
|
||||
print("дважды-трижды (для каждого воркера при spawn). На fork — один раз.")
|
||||
print("Вывод: старт воркеров не бесплатный — это накладные расходы лабы 4.")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -0,0 +1,62 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""Демо: что можно, а что нельзя передать в Pool (pickle).
|
||||
|
||||
Запуск из корня курса:
|
||||
python lab04_multiprocessing/examples/02_pickle_demo.py
|
||||
|
||||
Вложенные функции и лямбды не пиккелятся — PicklingError.
|
||||
Модульные функции — можно. Отсюда правило: kernel для Pool живёт
|
||||
на уровне модуля (как dispatch_kernel в common.tasks).
|
||||
"""
|
||||
import multiprocessing
|
||||
import os
|
||||
import sys
|
||||
|
||||
sys.path.insert(0, os.path.abspath(os.path.join(os.path.dirname(__file__), "..", "..")))
|
||||
|
||||
from common.tasks import dispatch_kernel, smoke_params
|
||||
|
||||
|
||||
def module_level_double(x): # МОЖНО: функция уровня модуля
|
||||
return x * 2
|
||||
|
||||
|
||||
def make_nested():
|
||||
def nested_double(x): # НЕЛЬЗЯ: вложенная функция
|
||||
return x * 2
|
||||
return nested_double
|
||||
|
||||
|
||||
def main():
|
||||
ctx = multiprocessing.get_context()
|
||||
|
||||
print("--- модульная функция: работает ---")
|
||||
with ctx.Pool(2) as pool:
|
||||
print(f" pool.map(module_level_double, [1,2,3]) = "
|
||||
f"{pool.map(module_level_double, [1, 2, 3])}")
|
||||
|
||||
print("\n--- lambda: PicklingError ---")
|
||||
try:
|
||||
with ctx.Pool(2) as pool:
|
||||
pool.map(lambda x: x * 2, [1, 2, 3])
|
||||
except Exception as e:
|
||||
print(f" {type(e).__name__}: {str(e)[:90]}...")
|
||||
|
||||
print("\n--- kernel задачи через dispatch_kernel: работает ---")
|
||||
task_name = "primes"
|
||||
chunks = None
|
||||
task = __import__("common.tasks", fromlist=["TASKS"]).TASKS[task_name]
|
||||
data = task["build"](smoke_params(task_name))
|
||||
chunks = task["split"](data, 2)
|
||||
payload = [(task_name, ch) for ch in chunks]
|
||||
with ctx.Pool(2) as pool:
|
||||
partials = pool.map(dispatch_kernel, payload)
|
||||
print(f" подсчёт простых ≤ 2000: {task['combine'](partials)} "
|
||||
f"(ожидалось 303)")
|
||||
|
||||
print("\nВывод: всё, что уходит воркерам, должно пиккелиться —")
|
||||
print("поэтому в этом курсе kernel диспетчеризуется по имени задачи.")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -0,0 +1,66 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""Демо: Process + Queue — ручная организация воркеров.
|
||||
|
||||
Запуск из корня курса:
|
||||
python lab04_multiprocessing/examples/03_process_queue_demo.py
|
||||
|
||||
Pool — удобный «конвейер». Process+Queue — ручная схема: воркеры сами берут
|
||||
задачи из очереди. Это основа для producer-consumer в лабе 5.
|
||||
"""
|
||||
import multiprocessing
|
||||
import os
|
||||
import sys
|
||||
import time
|
||||
|
||||
sys.path.insert(0, os.path.abspath(os.path.join(os.path.dirname(__file__), "..", "..")))
|
||||
|
||||
|
||||
def worker(task_q, result_q):
|
||||
"""Воркер: берёт задачи, пока не встретит None — «маркер конца»."""
|
||||
name = multiprocessing.current_process().name
|
||||
while True:
|
||||
task = task_q.get()
|
||||
if task is None:
|
||||
result_q.put((name, "done"))
|
||||
break
|
||||
idx, value = task
|
||||
time.sleep(0.05) # «полезная работа»
|
||||
result_q.put((name, f"задача {idx}: {value}² = {value ** 2}"))
|
||||
|
||||
|
||||
def main():
|
||||
ctx = multiprocessing.get_context()
|
||||
task_q = ctx.Queue()
|
||||
result_q = ctx.Queue()
|
||||
|
||||
n_workers = 3
|
||||
workers = [ctx.Process(target=worker, args=(task_q, result_q),
|
||||
name=f"воркер-{i}")
|
||||
for i in range(n_workers)]
|
||||
for w in workers:
|
||||
w.start()
|
||||
|
||||
# кладём задачи
|
||||
for i in range(9):
|
||||
task_q.put((i, i + 1))
|
||||
# маркеры конца — по одному на воркер
|
||||
for _ in range(n_workers):
|
||||
task_q.put(None)
|
||||
|
||||
# собираем результаты
|
||||
done = 0
|
||||
t0 = time.perf_counter()
|
||||
while done < n_workers:
|
||||
name, msg = result_q.get()
|
||||
if msg == "done":
|
||||
done += 1
|
||||
print(f" [{name}] завершился")
|
||||
else:
|
||||
print(f" [{name}] {msg}")
|
||||
for w in workers:
|
||||
w.join()
|
||||
print(f"время: {time.perf_counter() - t0:.2f} c (9 задач по 0.05 c на 3 воркерах)")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -0,0 +1,111 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""Лабораторная работа 4: мультипроцессорный код.
|
||||
|
||||
Заполните функции ниже. ИНТЕРФЕЙСЫ МЕНЯТЬ НЕЛЬЗЯ — по ним работают
|
||||
автотесты (tests/test_lab04.py, запускаются через подпроцесс).
|
||||
|
||||
Запуск из корня курса:
|
||||
python lab04_multiprocessing/solution.py
|
||||
"""
|
||||
import multiprocessing
|
||||
import json
|
||||
import os
|
||||
import sys
|
||||
import time
|
||||
|
||||
sys.path.insert(0, os.path.abspath(os.path.join(os.path.dirname(__file__), "..")))
|
||||
|
||||
from common.benchmark import make_table, save_speedup_plot, time_call
|
||||
from common.tasks import TASKS, dispatch_kernel, get_variant
|
||||
|
||||
# TODO: впишите свой номер в журнале (1..20)
|
||||
VARIANT_NUMBER = 0
|
||||
|
||||
|
||||
def _sum_list(data):
|
||||
"""Функция уровня модуля — суммирует список (для measure_overhead)."""
|
||||
return sum(data)
|
||||
|
||||
|
||||
def run_cpu_processes(task_name: str, params: dict, p: int):
|
||||
"""CPU-задача на p процессах через Pool и dispatch_kernel.
|
||||
|
||||
Вернуть (checksum, elapsed_seconds). Подсказки:
|
||||
- payload = [(task_name, chunk) для каждого куска];
|
||||
- dispatch_kernel пиккелится, лямбды и вложенные функции — НЕТ;
|
||||
- ctx = multiprocessing.get_context(); with ctx.Pool(p) as pool: ...
|
||||
"""
|
||||
task = TASKS[task_name]
|
||||
# TODO: build -> split -> pool.map(dispatch_kernel, payload) ->
|
||||
# combine -> checksum (замерить время)
|
||||
raise NotImplementedError
|
||||
|
||||
|
||||
def measure_overhead(items_count_list):
|
||||
"""Стоимость передачи данных процессу.
|
||||
|
||||
Для каждого n из items_count_list: отправить в воркер список из n чисел
|
||||
(list(range(n))), воркер возвращает их сумму (_sum_list). Пул создаётся
|
||||
ОДИН раз, первый вызов — прогрев. Вернуть {n: лучшее_время из 3}.
|
||||
"""
|
||||
# TODO: ctx.Pool(1) -> для каждого n: замер pool.map(_sum_list, [data])
|
||||
raise NotImplementedError
|
||||
|
||||
|
||||
def main():
|
||||
# --- smoke-режим для автотестов: быстрый прогон и JSON-отчёт -----------
|
||||
# (не удаляйте: тесты запускают solution.py подпроцессом с LAB_SMOKE=1)
|
||||
if os.environ.get("LAB_SMOKE") == "1":
|
||||
from common.tasks import run_sequential, smoke_params
|
||||
name = "integrate" # smoke не зависит от варианта: быстрая задача
|
||||
sp = smoke_params(name)
|
||||
expected, _ = run_sequential(name, sp, parts=1)
|
||||
got, _t = run_cpu_processes(name, sp, p=2)
|
||||
ok = (got == expected) or (isinstance(got, float) and
|
||||
abs(got - expected) < 1e-6 * max(1, abs(expected)))
|
||||
ov = measure_overhead([1_000, 50_000])
|
||||
print("LAB_SMOKE_JSON:" + json.dumps({
|
||||
"ok": True, "result_ok": bool(ok), "p": 2,
|
||||
"overhead": {str(k): t for k, t in ov.items()}}))
|
||||
return
|
||||
|
||||
v = get_variant(VARIANT_NUMBER)
|
||||
print(f"Вариант {v['variant']}: {v['describe']}")
|
||||
print(f"метод запуска: {multiprocessing.get_start_method()}\n")
|
||||
|
||||
# --- T(p) на процессах ---------------------------------------------------
|
||||
rows, times = [], {}
|
||||
expected = None
|
||||
for p in v["p_list"]:
|
||||
t, (cs, _) = time_call(run_cpu_processes,
|
||||
(v["task_name"], v["cpu_params"], p), repeats=2)
|
||||
times[p] = t
|
||||
rows.append((f"процессы, p={p}", t))
|
||||
if expected is None:
|
||||
expected = cs
|
||||
assert cs == expected or (isinstance(cs, float) and
|
||||
abs(cs - expected) < 1e-6 * max(1, abs(expected))), \
|
||||
f"checksum при p={p} не совпал!"
|
||||
print(make_table(rows, title="CPU-задача на процессах"))
|
||||
|
||||
base = times[1]
|
||||
for p, t in times.items():
|
||||
print(f" S(p={p}) = {base / t:.2f}")
|
||||
print()
|
||||
|
||||
# --- Накладные расходы передачи -------------------------------------------
|
||||
overhead = measure_overhead([1_000, 100_000, 1_000_000])
|
||||
print("Стоимость передачи данных процессу (туда+обратно):")
|
||||
for n, t in sorted(overhead.items()):
|
||||
print(f" {n:>9,} чисел: {t * 1000:8.1f} мс".replace(",", " "))
|
||||
|
||||
save_speedup_plot({"процессы": sorted(times.items())},
|
||||
f"Ускорение CPU-задачи процессами ({v['describe']})",
|
||||
os.path.join(os.path.dirname(__file__), "speedup.png"))
|
||||
|
||||
# TODO: в отчёт — S(p) и число физических ядер вашей машины,
|
||||
# сравните время kernel варианта со временем передачи из measure_overhead.
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -0,0 +1,118 @@
|
||||
# Лабораторная работа 4. Мультипроцессорный код
|
||||
|
||||
**Время на работу:** 2 пары. **Зависимости:** лабы 1–3.
|
||||
|
||||
## Цель работы
|
||||
|
||||
Освоить настоящую параллельность на CPU-задачах: процессы обходят GIL.
|
||||
Измерить накладные расходы multiprocessing и понять, когда они съедают
|
||||
выигрыш.
|
||||
|
||||
## Теория
|
||||
|
||||
### Процесс против потока
|
||||
|
||||
Процесс — отдельный экземпляр интерпретатора со **своей** памятью. У каждого
|
||||
процесса свой GIL, поэтому CPU-задачи на p процессах исполняются по-настоящему
|
||||
параллельно (на p ядрах). Цена: создание процесса ~10–100 мс (против ~0.1 мс
|
||||
у потока), а любые данные передаются между процессами через **pickle** —
|
||||
сериализацию. Большие данные (матрицы, списки миллионов чисел) передаются
|
||||
долго, и это становится главным ограничением.
|
||||
|
||||
### Методы запуска: fork и spawn
|
||||
|
||||
- **fork** (Linux/macOS): дочерний процесс — копия родителя «снимком» памяти.
|
||||
Быстро, но опасно с потоками и блокировками.
|
||||
- **spawn** (Windows, по умолчанию на macOS): запускается новый интерпретатор,
|
||||
который заново импортирует ваш модуль. Медленнее и строже: **весь код
|
||||
создания процессов должен быть под `if __name__ == "__main__":`**, иначе
|
||||
бесконечное размножение процессов.
|
||||
|
||||
Проверить и выбрать: `multiprocessing.get_start_method()` /
|
||||
`multiprocessing.get_context("spawn")`.
|
||||
|
||||
### Pool и передаваемые функции
|
||||
|
||||
`multiprocessing.Pool(p)` — пул из p воркеров-процессов:
|
||||
|
||||
```python
|
||||
with multiprocessing.Pool(p) as pool:
|
||||
results = pool.map(kernel, chunks) # как map, но по процессам
|
||||
```
|
||||
|
||||
ВАЖНО: `kernel` и `chunks` должны пиккелиться. Вложенные функции (def внутри
|
||||
def) и лямбды — **не** пиккелятся. Поэтому в `common.tasks` ядро задачи
|
||||
вызывается через модульный `dispatch_kernel((name, chunk))`.
|
||||
|
||||
### Что именно передаётся
|
||||
|
||||
Схема split → kernel → combine в мире процессов означает:
|
||||
- `chunks` сериализуются и уходят воркерам (расход на передачу);
|
||||
- частичные результаты сериализуются и возвращаются (расход на возврат);
|
||||
- если данные больше результата — передача может занять дольше, чем счёт
|
||||
(см. задачу «сортировка» в банке задач).
|
||||
|
||||
### Разные способы организации
|
||||
|
||||
- `Pool.map` — распределить список по воркерам;
|
||||
- `Process` + `Queue` — ручная схема: воркеры берут задачи из очереди и кладут
|
||||
ответы в другую (задел для лабы 5);
|
||||
- `ProcessPoolExecutor` —ThreadPoolExecutor-подобный интерфейс для процессов.
|
||||
|
||||
## Задание
|
||||
|
||||
CPU-задача — из вашего варианта. Заполните `solution.py`:
|
||||
|
||||
1. **`run_cpu_processes(task_name, params, p)`** — CPU-задача на `p`
|
||||
процессах через `Pool` и `dispatch_kernel`. Вернуть `(checksum, сек)`.
|
||||
Вся логика — под `if __name__ == "__main__"` (в `main()`).
|
||||
2. **`measure_overhead(items_count)`** — измерить стоимость передачи:
|
||||
отправить в воркер список из `items_count` чисел, воркер возвращает его
|
||||
длину. Время / количество = цена байта. Вернуть `{items_count: сек}`.
|
||||
3. **`main()`** — T(p) для p из варианта; таблица; график S(p) с «идеальной»
|
||||
прямой; вывод о накладных расходах.
|
||||
|
||||
### Ожидаемые результаты
|
||||
|
||||
- CPU-задача: S(p) растёт почти линейно до числа физических ядер;
|
||||
на машинах с гиперторингом S(p_max) может упираться в ~число физ. ядер.
|
||||
- `measure_overhead`: передача маленьких данных ~миллисекунды, больших —
|
||||
десятки-сотни мс. Сравните с временем kernel вашего варианта.
|
||||
|
||||
### Важно (Windows)
|
||||
|
||||
- Всё создание пулов — строго внутри `main()`/под `__main__`.
|
||||
- Не используйте лямбды и вложенные функции как аргументы `pool.map`.
|
||||
- На macOS/Linux можно сравнить `fork` и `spawn` (`get_context`) — добавьте
|
||||
в отчёт, если работаете не на Windows.
|
||||
|
||||
## Контрольные вопросы
|
||||
|
||||
1. Почему процессы ускоряют CPU-код, а потоки — нет?
|
||||
2. Что такое pickle и почему лямбда не может быть аргументом `pool.map`?
|
||||
3. Замерили T(p=8) на 4-ядерной машине с гиперторингом — ускорение ~4, а не 8.
|
||||
Почему?
|
||||
4. Когда передача данных между процессами дороже самих вычислений?
|
||||
5. Почему на Windows нужен `if __name__ == "__main__":`?
|
||||
|
||||
## Что сдаётся
|
||||
|
||||
`solution.py` + отчёт: таблица T(p), график S(p), замеры overhead передачи,
|
||||
анализ (до какого p растёт, где упирается), ответы на вопросы.
|
||||
|
||||
## Критерии оценки
|
||||
|
||||
| Пункт | Баллы |
|
||||
|---|---|
|
||||
| Корректный запуск на процессах, checksum совпадает | 2 |
|
||||
| S(p) > 1.5 при p=4 на CPU-задаче | 2 |
|
||||
| Замер overhead передачи данных | 2 |
|
||||
| График + анализ накладных расходов | 2 |
|
||||
| Ответы на контрольные вопросы | 2 |
|
||||
|
||||
## Типичные ошибки
|
||||
|
||||
- Нет `__main__` — на Windows бесконечный spawn (см. типичные_ошибки п.1).
|
||||
- Лямбда/вложенная функция в `pool.map` — `PicklingError`.
|
||||
- Слишком мелкие куски: создание процессов и передача дороже счёта,
|
||||
S(p) < 1. Увеличьте размер задачи.
|
||||
Reference in new issue
Block a user