# -*- coding: utf-8 -*- """ common.tasks — ядра CPU-задач и унифицированный доступ к вариантам. Все ядра «честно» последовательные на чистом Python (без NumPy), чтобы эффект GIL и накладные расходы параллелизма были видны честно. Единая схема каждой CPU-задачи (MapReduce-подобная): data = build(params) # подготовить данные chunks = split(data, parts) # разбить работу на части partial = kernel(chunk) # вычислить одну часть raw = combine(partials) # собрать итог value = checksum(raw) # свёртка результата для таблиц/сравнений Важно: checksum не зависит от способа разбиения — поэтому последовательный и любой параллельный запуск обязаны давать одинаковый checksum. Для multiprocessing функции передавать напрямую нельзя (вложенные функции не пиккелятся). Используйте dispatch_kernel модульного уровня: from common.tasks import dispatch_kernel with Pool(p) as pool: partials = pool.map(dispatch_kernel, [(name, ch) for ch in chunks]) На верхнем уровне модуля только определения — безопасно для spawn (Windows). """ from __future__ import annotations import hashlib import random import re import time from collections import Counter from typing import Any, Callable, Dict, List, Sequence, Tuple # --------------------------------------------------------------------------- # Реестр задач # --------------------------------------------------------------------------- TASKS: Dict[str, Dict[str, Callable]] = {} def register(name: str): def wrap(fn): TASKS[name] = fn() return fn return wrap # --------------------------------------------------------------------------- # Задача 1. Умножение матриц (наивное, чистый Python) # --------------------------------------------------------------------------- @register("matmul") def _task_matmul(): def build(params): n = params["n"] rng = random.Random(42) A = [[rng.random() for _ in range(n)] for _ in range(n)] B = [[rng.random() for _ in range(n)] for _ in range(n)] return (A, B) def split(data, parts): A, B = data n = len(A) bounds = [i * n // parts for i in range(parts + 1)] return [(A, B, bounds[i], bounds[i + 1]) for i in range(parts) if bounds[i] < bounds[i + 1]] def kernel(chunk): A, B, i0, i1 = chunk n = len(A) # транспонируем B для линейного доступа по памяти (эффект кэша) Bt = [[B[j][k] for j in range(n)] for k in range(n)] C = [] for i in range(i0, i1): Ai = A[i] row = [] for j in range(n): Bj = Bt[j] s = 0.0 for k in range(n): s += Ai[k] * Bj[k] row.append(s) C.append(row) return (i0, C) def combine(partials): C = [] for _, rows in sorted(partials, key=lambda r: r[0]): C.extend(rows) return C def checksum(raw): s = 0.0 for row in raw: for x in row: s += x return s def params(level): return {"n": {"S": 140, "M": 180, "L": 220}[level]} def describe(params): return f"умножение матриц {params['n']}x{params['n']}" return dict(build=build, split=split, kernel=kernel, combine=combine, checksum=checksum, params=params, describe=describe) # --------------------------------------------------------------------------- # Задача 2. Число π методом Монте-Карло # --------------------------------------------------------------------------- # --------------------------------------------------------------------------- # Задача 2. Число π методом Монте-Карло # # Важный приём: total бросков делится на ФИКСИРОВАННЫЕ серии по SERIES_SIZE # бросков, у каждой серии своё зерно по её номеру. Кусок = диапазон номеров # серий. Поэтому checksum не зависит от того, как серии сгруппированы в # куски — приём, реально применяемый в MC-расчётах для воспроизводимости. # --------------------------------------------------------------------------- @register("pi_mc") def _task_pi_mc(): SERIES_SIZE = 10_000 # бросков в серии; серия = единица воспроизводимости def build(params): return params["total"] def split(data, parts): total = data n_series = total // SERIES_SIZE bounds = [i * n_series // parts for i in range(parts + 1)] return [((i, bounds[i], bounds[i + 1]), total) for i in range(parts) if bounds[i] < bounds[i + 1]] def kernel(chunk): (idx, s0, s1), _total = chunk inside = 0 count = 0 for s in range(s0, s1): rng = random.Random(s) # зерно = номер серии for _ in range(SERIES_SIZE): x = rng.random() y = rng.random() if x * x + y * y <= 1.0: inside += 1 count += SERIES_SIZE return (inside, count) def combine(partials): inside = sum(p[0] for p in partials) total = sum(p[1] for p in partials) return 4.0 * inside / total def checksum(raw): return raw def params(level): return {"total": {"S": 2_000_000, "M": 4_000_000, "L": 8_000_000}[level]} def describe(params): return f"π методом Монте-Карло, {params['total']:,} бросков".replace(",", " ") return dict(build=build, split=split, kernel=kernel, combine=combine, checksum=checksum, params=params, describe=describe) # --------------------------------------------------------------------------- # Задача 3. Численное интегрирование (метод средних прямоугольников) # --------------------------------------------------------------------------- @register("integrate") def _task_integrate(): def f(x): return 4.0 / (1.0 + x * x) # на [0, 1] интеграл = π (проверяемый ответ) def build(params): return params def split(data, parts): a, b, n = data["a"], data["b"], data["n"] # кусок = диапазон ИНДЕКСОВ отрезков [i0, i1): kernel считает всегда # одни и те же отрезки, checksum не зависит от разбиения bounds = [i * n // parts for i in range(parts + 1)] return [(data, bounds[i], bounds[i + 1]) for i in range(parts) if bounds[i] < bounds[i + 1]] def kernel(chunk): data, i0, i1 = chunk a, b, n = data["a"], data["b"], data["n"] h = (b - a) / n s = 0.0 for i in range(i0, i1): s += f(a + (i + 0.5) * h) return s * h def combine(partials): return sum(partials) def checksum(raw): return raw def params(level): n = {"S": 3_000_000, "M": 6_000_000, "L": 10_000_000}[level] return {"a": 0.0, "b": 10.0, "n": n} def describe(params): n_str = f"{params['n']:,}".replace(",", " ") return (f"интегрирование [{params['a']:g}, {params['b']:g}] " f"{n_str} отрезков") return dict(build=build, split=split, kernel=kernel, combine=combine, checksum=checksum, params=params, describe=describe) # --------------------------------------------------------------------------- # Задача 4. Подсчёт простых чисел ≤ M (перебор делителей) # --------------------------------------------------------------------------- @register("primes") def _task_primes(): def _is_prime(k): if k < 2: return False if k % 2 == 0: return k == 2 d = 3 while d * d <= k: if k % d == 0: return False d += 2 return True def build(params): return params["M"] def split(data, parts): M = data lo = 2 span = (M + 1 - lo) // parts if span == 0: span = 1 bounds = [lo + i * span for i in range(parts)] + [M + 1] return [(bounds[i], min(bounds[i + 1], M + 1)) for i in range(parts) if bounds[i] < min(bounds[i + 1], M + 1)] def kernel(chunk): lo, hi = chunk c = 0 for k in range(lo, hi): if _is_prime(k): c += 1 return c def combine(partials): return sum(partials) def checksum(raw): return raw def params(level): return {"M": {"S": 1_000_000, "M": 2_000_000, "L": 4_000_000}[level]} def describe(params): return f"подсчёт простых чисел ≤ {params['M']:,}".replace(",", " ") return dict(build=build, split=split, kernel=kernel, combine=combine, checksum=checksum, params=params, describe=describe) # --------------------------------------------------------------------------- # Задача 5. Задача N ферзей (подсчёт расстановок) # --------------------------------------------------------------------------- @register("nqueens") def _task_nqueens(): def build(params): return params["N"] def split(data, parts): N = data # одна ветка = первый ферзь в колонке col return [(N, col) for col in range(N)] def kernel(chunk): N, col0 = chunk count = 0 cols = [-1] * N def place(row): nonlocal count if row == N: count += 1 return start = col0 if row == 0 else 0 for col in range(start, N): ok = True for r in range(row): c = cols[r] if c == col or abs(c - col) == row - r: ok = False break if ok: cols[row] = col place(row + 1) cols[row] = -1 place(0) return count def combine(partials): return sum(partials) def checksum(raw): return raw def params(level): return {"N": {"S": 10, "M": 11, "L": 12}[level]} def describe(params): return f"задача {params['N']} ферзей, подсчёт расстановок" return dict(build=build, split=split, kernel=kernel, combine=combine, checksum=checksum, params=params, describe=describe) # --------------------------------------------------------------------------- # Задача 6. Частотный анализ текста # --------------------------------------------------------------------------- _WORD_RE = re.compile(r"[а-яёa-z0-9]+") _SENTENCES = [ "параллельное программирование это просто", "потоки и процессы работают по разному", "гил ограничивает настоящую параллельность потоков", "асинхронный код ждёт ввод вывод эффективно", "метод монте карло считает число пи", "распределённые системы соединяют много машин", "синхронизация защищает общие данные", "очередь сообщений соединяет производителя и потребителя", "каждый воркер обрабатывает свой кусок работы", "ускорение зависит от доли последовательного кода", ] @register("wordcount") def _task_wordcount(): def build(params): lines = params["lines"] rng = random.Random(12345) words_pool = [w for s in _SENTENCES for w in s.split()] out = [] for _ in range(lines): k = rng.randint(5, 20) out.append(" ".join(rng.choice(words_pool) for _ in range(k))) return out def split(data, parts): lines = data n = len(lines) bounds = [i * n // parts for i in range(parts + 1)] return [lines[bounds[i]:bounds[i + 1]] for i in range(parts) if bounds[i] < bounds[i + 1]] def kernel(chunk): c = Counter() for line in chunk: for w in _WORD_RE.findall(line): c[w] += 1 return c def combine(partials): total = Counter() for c in partials: total.update(c) return total def checksum(raw): # полная сумма вхождений — не зависит от разбиения return sum(raw.values()) def params(level): return {"lines": {"S": 100_000, "M": 200_000, "L": 400_000}[level]} def describe(params): return f"частотный анализ, {params['lines']:,} строк".replace(",", " ") return dict(build=build, split=split, kernel=kernel, combine=combine, checksum=checksum, params=params, describe=describe) # --------------------------------------------------------------------------- # Задача 7. Box blur (размытие матрицы, чистый Python, кайма 1 строка) # --------------------------------------------------------------------------- @register("blur") def _task_blur(): def build(params): size = params["size"] rng = random.Random(777) return [[rng.random() for _ in range(size)] for _ in range(size)] def split(data, parts): img = data n = len(img) bounds = [i * n // parts for i in range(parts + 1)] chunks = [] for i in range(parts): i0, i1 = bounds[i], bounds[i + 1] if i0 >= i1: continue top = max(0, i0 - 1) bot = min(n, i1 + 1) sub = [row[:] for row in img[top:bot]] # копия с каймой (halo) chunks.append((sub, i0 - top, i1 - top)) return chunks def kernel(chunk): sub, lo, hi = chunk # вычисляем строки [lo, hi) внутри sub m = len(sub[0]) out = [] for r in range(lo, hi): new_row = [] for c in range(m): s = 0.0 cnt = 0 for dr in (-1, 0, 1): rr = r + dr if 0 <= rr < len(sub): row = sub[rr] for dc in (-1, 0, 1): cc = c + dc if 0 <= cc < m: s += row[cc] cnt += 1 new_row.append(s / cnt) out.append(new_row) return (lo, out) def combine(partials): out = [] for _, rows in sorted(partials, key=lambda r: r[0]): out.extend(rows) return out def checksum(raw): s = 0.0 for row in raw: for x in row: s += x return s def params(level): return {"size": {"S": 600, "M": 900, "L": 1200}[level]} def describe(params): return f"box blur {params['size']}x{params['size']}" return dict(build=build, split=split, kernel=kernel, combine=combine, checksum=checksum, params=params, describe=describe) # --------------------------------------------------------------------------- # Задача 8. Хеширование строк SHA-256 # --------------------------------------------------------------------------- @register("hashing") def _task_hashing(): def build(params): lines = params["lines"] rng = random.Random(555) pool = [w for s in _SENTENCES for w in s.split()] out = [] for _ in range(lines): k = rng.randint(30, 80) out.append(" ".join(rng.choice(pool) for _ in range(k)).encode("utf-8")) return out def split(data, parts): items = data n = len(items) bounds = [i * n // parts for i in range(parts + 1)] return [items[bounds[i]:bounds[i + 1]] for i in range(parts) if bounds[i] < bounds[i + 1]] def kernel(items): # свёртка по каждому элементу отдельно => результат не зависит # от границ кусков. hashlib отпускает GIL на время одного вызова: # на коротких строках эффект незаметен, на блоках ≥ 256 КБ потоки # реально параллелятся (демо в лабе 3) acc = 0 for it in items: d = hashlib.sha256(it).digest() acc = (acc + int.from_bytes(d[:8], "big")) & 0xFFFFFFFFFFFFFFFF return acc def combine(partials): acc = 0 for v in partials: acc = (acc + v) & 0xFFFFFFFFFFFFFFFF return acc def checksum(raw): return raw def params(level): return {"lines": {"S": 300_000, "M": 600_000, "L": 1_200_000}[level]} def describe(params): return f"хеширование SHA-256, {params['lines']:,} строк".replace(",", " ") return dict(build=build, split=split, kernel=kernel, combine=combine, checksum=checksum, params=params, describe=describe) # --------------------------------------------------------------------------- # Задача 9. Сортировка слиянием (куски + k-way merge) # --------------------------------------------------------------------------- @register("sortbig") def _task_sortbig(): def build(params): n = params["n"] rng = random.Random(999) return [rng.random() for _ in range(n)] def split(data, parts): arr = data n = len(arr) bounds = [i * n // parts for i in range(parts + 1)] return [arr[bounds[i]:bounds[i + 1]] for i in range(parts) if bounds[i] < bounds[i + 1]] def kernel(chunk): return sorted(chunk) def combine(partials): import heapq return list(heapq.merge(*partials)) def checksum(raw): s = 0.0 prev = -1.0 for x in raw: if x < prev: raise ValueError("массив не отсортирован") s += x prev = x return s def params(level): return {"n": {"S": 1_000_000, "M": 2_000_000, "L": 4_000_000}[level]} def describe(params): return f"сортировка слиянием {params['n']:,} чисел".replace(",", " ") return dict(build=build, split=split, kernel=kernel, combine=combine, checksum=checksum, params=params, describe=describe) # --------------------------------------------------------------------------- # Задача 10. Игра «Жизнь» (параллелятся строки одного поколения) # --------------------------------------------------------------------------- @register("life") def _task_life(): def build(params): size = params["size"] rng = random.Random(31337) grid = [[1 if rng.random() < 0.30 else 0 for _ in range(size)] for _ in range(size)] return (grid, params["steps"]) def split(data, parts): grid, _steps = data n = len(grid) bounds = [i * n // parts for i in range(parts + 1)] return [(grid, bounds[i], bounds[i + 1]) for i in range(parts) if bounds[i] < bounds[i + 1]] def kernel(chunk): # ОДНО поколение для строк [i0, i1). Соседние строки читаются, свои — пишутся. grid, i0, i1 = chunk m = len(grid[0]) band = [] for r in range(i0, i1): new_row = [0] * m for c in range(m): s = 0 for dr in (-1, 0, 1): rr = r + dr if 0 <= rr < len(grid): row = grid[rr] for dc in (-1, 0, 1): if dc == 0 and dr == 0: continue cc = c + dc if 0 <= cc < m: s += row[cc] alive = grid[r][c] == 1 if alive and (s == 2 or s == 3): new_row[c] = 1 elif not alive and s == 3: new_row[c] = 1 band.append(new_row) return (i0, band) def combine(partials): grid = [] for _, rows in sorted(partials, key=lambda r: r[0]): grid.extend(rows) return grid def checksum(raw): return sum(sum(row) for row in raw) def params(level): size = {"S": 500, "M": 600, "L": 700}[level] steps = {"S": 20, "M": 25, "L": 30}[level] return {"size": size, "steps": steps} def describe(params): return (f"игра «Жизнь» {params['size']}x{params['size']}, " f"{params['steps']} поколений") return dict(build=build, split=split, kernel=kernel, combine=combine, checksum=checksum, params=params, describe=describe) # --------------------------------------------------------------------------- # I/O-нагрузка (общая для всех вариантов, имитация сетевых запросов) # --------------------------------------------------------------------------- def io_fetch(item: int, delay: float = 0.05) -> Tuple[int, int]: """Имитация I/O-операции: блокирующая задержка + «полезная работа». Никаких настоящих сетевых вызовов — только time.sleep (блокирующий). Это позволяет сравнивать threading / multiprocessing / asyncio на равных. """ time.sleep(delay) return (item, item * item) def io_items(count: int) -> List[int]: return list(range(count)) # --------------------------------------------------------------------------- # Вызов kernel из процессов (модульного уровня — пиккелится) # --------------------------------------------------------------------------- def dispatch_kernel(payload: Tuple[str, Any]): """payload = (имя_задачи, chunk) -> частичный результат. Pool не умеет передавать вложенные функции, поэтому dispatch живёт на уровне модуля и диспетчеризует по имени задачи. """ name, chunk = payload return TASKS[name]["kernel"](chunk) # --------------------------------------------------------------------------- # Последовательный запуск (эталон для сравнения и проверки) # --------------------------------------------------------------------------- def run_sequential(task_name: str, params: dict, parts: int = 1): """Последовательное выполнение задачи. Возвращает (checksum, секунды). parts > 1 означает: те же куски, что пошли бы воркерам, но выполняются в главном процессе по очереди. Используется, чтобы проверить, что разбиение не меняет результат. """ task = TASKS[task_name] data = task["build"](params) t0 = time.perf_counter() if task_name == "life": grid, steps = data for _ in range(steps): chunks = task["split"]((grid, steps), parts) partials = [task["kernel"](ch) for ch in chunks] grid = task["combine"](partials) raw = grid else: chunks = task["split"](data, parts) partials = [task["kernel"](ch) for ch in chunks] raw = task["combine"](partials) elapsed = time.perf_counter() - t0 return task["checksum"](raw), elapsed # --------------------------------------------------------------------------- # Известные ответы (для автотестов) # --------------------------------------------------------------------------- _NQUEENS_KNOWN = {4: 2, 5: 10, 6: 4, 7: 40, 8: 92, 9: 352, 10: 724, 11: 2680, 12: 14200} def known_checksum(task_name: str, params: dict): """Известный правильный ответ, если он есть, иначе None. Для pi_mc и integrate ответ известен математически (π), но с ограниченной точностью — их проверяют автотесты отдельно. """ import math if task_name == "nqueens": return _NQUEENS_KNOWN.get(params["N"]) if task_name == "primes": M = params["M"] sieve = bytearray([1]) * (M + 1) sieve[0:2] = b"\x00\x00" for p in range(2, int(M ** 0.5) + 1): if sieve[p]: sieve[p * p:: p] = bytearray(len(sieve[p * p:: p])) return int(sum(sieve)) if task_name == "integrate" and params["a"] == 0.0 and params["b"] == 1.0: return math.pi if task_name == "pi_mc": return math.pi # приближённо, тест сравнит с допуском return None # --------------------------------------------------------------------------- # «Дымовые» параметры (маленькие — для автотестов и самопроверки) # --------------------------------------------------------------------------- _SMOKE = { "matmul": {"n": 12}, "pi_mc": {"total": 50_000}, # кратно SERIES_SIZE=10_000 "integrate": {"a": 0.0, "b": 1.0, "n": 20_000}, "primes": {"M": 2_000}, "nqueens": {"N": 6}, "wordcount": {"lines": 2_000}, "blur": {"size": 40}, "hashing": {"lines": 5_000}, "sortbig": {"n": 20_000}, "life": {"size": 20, "steps": 3}, } def smoke_params(task_name: str) -> dict: return dict(_SMOKE[task_name]) # --------------------------------------------------------------------------- # Варианты # --------------------------------------------------------------------------- _VARIANTS: Dict[int, Tuple[str, str]] = { 1: ("matmul", "S"), 2: ("pi_mc", "M"), 3: ("integrate", "L"), 4: ("primes", "S"), 5: ("nqueens", "M"), 6: ("wordcount", "L"), 7: ("blur", "S"), 8: ("hashing", "M"), 9: ("sortbig", "L"), 10: ("life", "S"), 11: ("matmul", "M"), 12: ("pi_mc", "L"), 13: ("integrate", "S"), 14: ("primes", "M"), 15: ("nqueens", "L"), 16: ("wordcount", "S"), 17: ("blur", "M"), 18: ("hashing", "L"), 19: ("sortbig", "S"), 20: ("life", "M"), } _IO_BY_LEVEL = {"S": (24, 0.05), "M": (48, 0.05), "L": (96, 0.05)} def get_variant(num: int) -> Dict[str, Any]: """Параметры варианта по номеру студента в журнале (1–20).""" if not 1 <= num <= 20: raise ValueError("номер варианта должен быть от 1 до 20") task_name, level = _VARIANTS[num] task = TASKS[task_name] cpu_params = task["params"](level) io_items_n, io_delay = _IO_BY_LEVEL[level] return { "variant": num, "task_name": task_name, "level": level, "cpu_params": cpu_params, "io": {"items": io_items_n, "delay": io_delay}, "p_list": [1, 2, 4] if level in ("S", "M") else [1, 2, 4, 8], "describe": task["describe"](cpu_params), }