Files
2026-10-07 13:06:47 +05:00

111 lines
5.0 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# -*- 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()