IlmHamroh
Python kursi/Konkurentlik va parallellik5/8-dars22 daqiqa
Mundarija (21)

14.5-dars: Pool va concurrent.futures

14-QISM — KONKURENTLIK VA PARALLELLIK · 5-dars


1. Kirish va motivatsiya

14.4-darsda jarayonlarni qo'lda boshqardik: Process yaratish, Queue orqali natija olish, join, exitcode, xatolarni qo'lda uzatish. Bu ish har safar takrorlanadi — va har safar xato qilish imkoni bor.

Real vaziyat. Ma'lumotlar jamoasi 200 000 ta hujjatni tahlil qiladigan skript yozdi: Process, Queue, sentinel lar, qayta urinish — 180 qator kod. Uni ProcessPoolExecutor ga ko'chirganda 22 qator qoldi va uchta yashirin xato o'z-o'zidan yo'qoldi (yo'qolgan natijalar, jim xatolar, to'xtatishda osilib qolish).

multiprocessing.Pool va concurrent.futures — bir xil g'oya: ishchilar puli va vazifalar. Siz faqat funksiya va ma'lumotni berasiz; yaratish, taqsimlash, natijalarni yig'ish va xatolarni uzatish esa kutubxona zimmasida.

Bu darsda:

  • concurrent.futures: ThreadPoolExecutor, ProcessPoolExecutor, InterpreterPoolExecutor
  • Future: result, exception, cancel, as_completed, wait
  • multiprocessing.Pool: map, imap, imap_unordered, starmap, apply_async
  • chunksize, initializer, maxtasksperchild
  • Xatolar, muddatlar va to'xtatish (shutdown(cancel_futures=True))
  • Qaysi API ni qachon tanlash
  • Amaliy: hujjatlarni parallel tahlil qilish quvuri

2. Nazariya — chuqur tushuntirish

2.1. Ikki API

concurrent.futures multiprocessing.Pool
Iplar uchun ThreadPoolExecutor (ThreadPool — hujjatsiz)
Jarayonlar uchun ProcessPoolExecutor
Subinterpretatorlar (3.14) InterpreterPoolExecutor
Natija obyekti Future (umumiy interfeys) AsyncResult
Dangasa iteratsiya map (natijalar yig'iladi) imap, imap_unordered
Ishchini boshlash funksiyasi initializer initializer
Ishchini yangilash maxtasksperchild

Yangi kodda concurrent.futures ni sukut tanlov qiling: bitta API bilan iplar, jarayonlar va subinterpretatorlar orasida almashish mumkin.

2.2. Executor API

python
with cf.ProcessPoolExecutor(max_workers=4) as ex:
    kelajak = ex.submit(ishla, x)             # bitta vazifa → Future
    natijalar = list(ex.map(ishla, malumot, chunksize=100))
Metod Izoh
submit(f, *args) Future
map(f, iterable, timeout=None, chunksize=1) Tartibda natijalar (generator)
shutdown(wait=True, cancel_futures=False) Tugatish; with avtomatik chaqiradi

Future:

Metod Izoh
result(timeout) Natija yoki istisno (qayta ko'tariladi)
exception(timeout) Istisno obyekti yoki None
cancel() Faqat hali boshlanmagan vazifani bekor qiladi (True/False)
done(), running(), cancelled() Holat
add_done_callback(f) Tugaganda chaqiriladi

Kutish:

python
cf.as_completed(kelajaklar, timeout=30)                  # tugash tartibida
tugagan, qolgan = cf.wait(kelajaklar, timeout=5, return_when=cf.FIRST_COMPLETED)

2.3. multiprocessing.Pool

python
with mp.Pool(processes=4, initializer=boshlash, initargs=(sozlama,), maxtasksperchild=100) as pul:
    natijalar = pul.map(ishla, malumot, chunksize=50)
    tartibsiz = pul.imap_unordered(ishla, malumot)        # tayyor bo'lgani chiqadi
    bitta = pul.apply_async(ishla, (x,), error_callback=qayd)
Metod Xususiyati
map Hammasini yig'adi, tartibda qaytaradi
imap Dangasa, tartibda
imap_unordered Dangasa, tayyor bo'lish tartibida (eng tez)
starmap Ko'p argumentli funksiyalar uchun
apply_async Bitta vazifa, AsyncResult

maxtasksperchild — ishchi belgilangan vazifadan keyin qayta ishga tushadi: xotira sizishiga qarshi klassik himoya.

2.4. chunksize va initializer

chunksize (12.6-dars): vazifalar ishchilarga paketlab yuboriladi.

Vazifa hajmi chunksize
Juda mayda (mikrosoniyalar) Katta: 1 000 – 10 000
O'rtacha (millisekundlar) 10 – 100
Og'ir (soniyalar) 1

initializer — har ishchi jarayonda bir marta bajariladi: ulanish ochish, model yuklash, sozlamani tayyorlash.

python
def boshlash(yol: str) -> None:
    global MODEL
    MODEL = ogir_model_yukla(yol)          # har vazifada emas, bir marta

2.5. Xatolar va muddatlar

Holat Xulq
Vazifada istisno Future.result() da qayta ko'tariladi; Pool.map da birinchi istisno
result(timeout=5) TimeoutError — vazifa bekor qilinmaydi, ishlayveradi
as_completed(..., timeout=) Muddat tugasa — TimeoutError
Ishchi jarayon yiqildi BrokenProcessPool — pul ishlatib bo'lmaydi
cancel() Faqat navbatda turgan vazifa uchun

Muddat vazifani to'xtatmaydi: jarayonlarda uni to'xtatish uchun terminate yoki o'z bekor qilish mexanizmingiz kerak.

2.6. To'xtatish

python
ex.shutdown(wait=True)                       # hammasi tugaguncha kutadi (with ning sukuti)
ex.shutdown(wait=False, cancel_futures=True) # navbatdagilarni bekor qiladi
Vaziyat Nima qilish
Normal tugatish with bloki
Tezkor to'xtatish cancel_futures=True — boshlanmaganlar bekor qilinadi
Ishlayotgan vazifalar Baribir tugaydi (majburan to'xtatish yo'q)
Jarayonni majburan Pool.terminate() yoki Process.terminate()

2.7. Qaysi API ni tanlash

Ehtiyoj Tanlov
Oddiy parallel map Executor.map
Natijalar tayyor bo'lgani zahoti kerak as_completed yoki imap_unordered
Iplar va jarayonlar orasida almashish concurrent.futures (bitta API)
Ishchilarni yangilash (maxtasksperchild) multiprocessing.Pool
asyncio bilan integratsiya loop.run_in_executor (14.7-dars)
Subinterpretatorlar (3.14) InterpreterPoolExecutor

3. Tez ma'lumotnoma

python
import concurrent.futures as cf

with cf.ProcessPoolExecutor(max_workers=4) as ex:
    for kelajak in cf.as_completed({ex.submit(ishla, x): x for x in malumot}):
        natija = kelajak.result()            # istisno shu yerda

with mp.Pool(4, initializer=boshlash, maxtasksperchild=100) as pul:
    for natija in pul.imap_unordered(ishla, malumot, chunksize=50):
        ...

Qoidalar

har Future natijasini o'qing
chunksize: mayda vazifa → katta paket
initializer: og'ir tayyorgarlik bir marta
shutdown(cancel_futures=True): navbatdagilarni bekor qiladi
timeout vazifani to'xtatmaydi

4. Batafsil misollar

Misol 1 — Executor asoslari

python
"""submit va map; Future holatlari; as_completed va wait; cancel; iplar va jarayonlar bitta API bilan."""

import concurrent.futures as cf
import time


def sekin_kvadrat(x: int) -> int:
    time.sleep(0.05)
    return x * x


def uzoq_kvadrat(x: int) -> int:
    time.sleep(0.3)
    return x * x


def main() -> None:
    print("=== 1. map va submit ===")
    with cf.ThreadPoolExecutor(max_workers=4) as ex:
        natijalar = list(ex.map(sekin_kvadrat, range(8)))
        kelajak = ex.submit(sekin_kvadrat, 10)
        print(f"  map natijalari (tartibda): {natijalar}")
        print(f"  submit → {type(kelajak).__name__}, natija: {kelajak.result(timeout=10)}")

    print("\n=== 2. Future holatlari ===")
    with cf.ThreadPoolExecutor(max_workers=1) as ex:
        birinchi = ex.submit(uzoq_kvadrat, 1)
        ikkinchi = ex.submit(uzoq_kvadrat, 2)           # navbatda kutadi
        time.sleep(0.05)
        print(f"  birinchi: running={birinchi.running()}, done={birinchi.done()}")
        print(f"  ikkinchi navbatda: running={ikkinchi.running()}, cancel() → {ikkinchi.cancel()}")
        print(f"  bekor qilingandan keyin: cancelled={ikkinchi.cancelled()}")
        print(f"  birinchi natijasi: {birinchi.result(timeout=10)}")
        print(f"  ishlayotgan vazifani bekor qilib bo'lmaydi: {birinchi.cancel()}")

    print("\n=== 3. as_completed: tugash tartibida ===")

    def turli_vaqt(x: int) -> tuple[int, float]:
        kutish = 0.05 * (5 - x)
        time.sleep(kutish)
        return x, kutish

    with cf.ThreadPoolExecutor(max_workers=5) as ex:
        kelajaklar = [ex.submit(turli_vaqt, i) for i in range(5)]
        tartib = [kelajak.result()[0] for kelajak in cf.as_completed(kelajaklar, timeout=10)]
    print(f"  yuborilish tartibi: {list(range(5))}")
    print(f"  tugash tartibi:     {tartib}  ← eng qisqa kutgani birinchi")

    print("\n=== 4. wait ===")
    with cf.ThreadPoolExecutor(max_workers=4) as ex:
        kelajaklar = [ex.submit(sekin_kvadrat, i) for i in range(4)]
        tugagan, qolgan = cf.wait(kelajaklar, timeout=0.01, return_when=cf.FIRST_COMPLETED)
        print(f"  10 ms dan keyin: tugagan {len(tugagan)}, qolgan {len(qolgan)}")
        tugagan, qolgan = cf.wait(kelajaklar, timeout=10, return_when=cf.ALL_COMPLETED)
        print(f"  kutgandan keyin: tugagan {len(tugagan)}, qolgan {len(qolgan)}")

    print("\n=== 5. Bitta API, uch xil ishchi ===")
    pullar = {
        "iplar": cf.ThreadPoolExecutor(max_workers=4),
        "jarayonlar": cf.ProcessPoolExecutor(max_workers=4),
        "subinterpretatorlar": cf.InterpreterPoolExecutor(max_workers=4),
    }
    for nom, ex in pullar.items():
        with ex:
            natija = list(ex.map(sekin_kvadrat, range(4)))
        print(f"  {nom:22} {natija}")
    print("  ⭐ kod bir xil — faqat Executor klassi almashadi")


if __name__ == "__main__":
    main()

Natijaning muhim qismi:

text
=== 1. map va submit ===
  map natijalari (tartibda): [0, 1, 4, 9, 16, 25, 36, 49]
  submit → Future, natija: 100

=== 2. Future holatlari ===
  birinchi: running=True, done=False
  ikkinchi navbatda: running=False, cancel() → True
  bekor qilingandan keyin: cancelled=True
  birinchi natijasi: 1
  ishlayotgan vazifani bekor qilib bo'lmaydi: False

=== 3. as_completed: tugash tartibida ===
  yuborilish tartibi: [0, 1, 2, 3, 4]
  tugash tartibi:     [4, 3, 2, 1, 0]  ← eng qisqa kutgani birinchi

=== 4. wait ===
  10 ms dan keyin: tugagan 0, qolgan 4
  kutgandan keyin: tugagan 4, qolgan 0

=== 5. Bitta API, uch xil ishchi ===
  iplar                  [0, 1, 4, 9]
  jarayonlar             [0, 1, 4, 9]
  subinterpretatorlar    [0, 1, 4, 9]
  ⭐ kod bir xil — faqat Executor klassi almashadi

Nima ko'rsatdi: 2.1, 2.2-bo'limlar.

Misol 2 — Xatolar, muddatlar va to'xtatish

python
"""Istisnolar Future ichida; map va as_completed farqi; timeout vazifani to'xtatmaydi; shutdown(cancel_futures)."""

import concurrent.futures as cf
import time

BAJARILGANLAR: list[int] = []


def beqaror(x: int) -> int:
    if x % 4 == 3:
        raise ValueError(f"x={x} ishlov bermaydi")
    time.sleep(0.02)
    return x * 10


def juda_sekin(x: int) -> int:
    time.sleep(0.5)
    BAJARILGANLAR.append(x)
    return x


def main() -> None:
    print("=== 1. Istisno Future ichida ===")
    with cf.ThreadPoolExecutor(max_workers=4) as ex:
        yaxshi, yomon = ex.submit(beqaror, 1), ex.submit(beqaror, 3)
        print(f"  yaxshi natija: {yaxshi.result(timeout=5)}")
        print(f"  yomon exception(): {type(yomon.exception(timeout=5)).__name__}")
        try:
            yomon.result()
        except ValueError as xato:
            print(f"  yomon result(): ValueError: {xato}")

    print("\n=== 2. map: birinchi xatoda to'xtaydi ===")
    with cf.ThreadPoolExecutor(max_workers=4) as ex:
        olingan = []
        try:
            for natija in ex.map(beqaror, range(8)):
                olingan.append(natija)
        except ValueError as xato:
            print(f"  {len(olingan)} natija olindi, keyin: {xato}")

    print("\n=== 3. as_completed: hammasi ko'riladi ===")
    with cf.ThreadPoolExecutor(max_workers=4) as ex:
        kelajaklar = {ex.submit(beqaror, i): i for i in range(8)}
        natijalar, xatolar = [], []
        for kelajak in cf.as_completed(kelajaklar, timeout=10):
            try:
                natijalar.append(kelajak.result())
            except ValueError:
                xatolar.append(kelajaklar[kelajak])
    print(f"  muvaffaqiyatli: {len(natijalar)}, xato bergan qiymatlar: {sorted(xatolar)}")
    print("  ⭐ hamma vazifa bajarildi va har biri hisobga olindi")

    print("\n=== 4. ⚠️ timeout vazifani to'xtatmaydi ===")
    BAJARILGANLAR.clear()
    with cf.ThreadPoolExecutor(max_workers=2) as ex:
        kelajak = ex.submit(juda_sekin, 1)
        try:
            kelajak.result(timeout=0.05)
        except TimeoutError:
            print("  result(timeout=0.05) → TimeoutError")
        print(f"  lekin vazifa hali ishlayapti: {not kelajak.done()}")
    print(f"  with blokidan chiqishda kutildi va bajarildi: {BAJARILGANLAR == [1]}")

    print("\n=== 5. shutdown(cancel_futures=True) ===")
    ex = cf.ThreadPoolExecutor(max_workers=2)
    kelajaklar = [ex.submit(beqaror, i) for i in (0, 1, 2, 4, 5, 6, 8, 9, 10, 12)]
    ex.shutdown(wait=True, cancel_futures=True)
    bekor = sum(k.cancelled() for k in kelajaklar)
    bajarilgan = sum(k.done() and not k.cancelled() for k in kelajaklar)
    print(f"  har vazifa yo bajarildi, yo bekor qilindi: {bajarilgan + bekor == len(kelajaklar)}")
    print(f"  ko'pchiligi bekor qilindi: {bekor >= 5}")
    print(f"  ishlab ulgurganlari bekor qilinmadi: {bajarilgan >= 1}")
    print("  ⭐ cancel_futures faqat navbatda turganlarga ta'sir qiladi")


if __name__ == "__main__":
    main()

Natijaning muhim qismi:

text
=== 1. Istisno Future ichida ===
  yaxshi natija: 10
  yomon exception(): ValueError
  yomon result(): ValueError: x=3 ishlov bermaydi

=== 2. map: birinchi xatoda to'xtaydi ===
  3 natija olindi, keyin: x=3 ishlov bermaydi

=== 3. as_completed: hammasi ko'riladi ===
  muvaffaqiyatli: 6, xato bergan qiymatlar: [3, 7]
  ⭐ hamma vazifa bajarildi va har biri hisobga olindi

=== 4. ⚠️ timeout vazifani to'xtatmaydi ===
  result(timeout=0.05) → TimeoutError
  lekin vazifa hali ishlayapti: True
  with blokidan chiqishda kutildi va bajarildi: True

=== 5. shutdown(cancel_futures=True) ===
  har vazifa yo bajarildi, yo bekor qilindi: True
  ko'pchiligi bekor qilindi: True
  ishlab ulgurganlari bekor qilinmadi: True
  ⭐ cancel_futures faqat navbatda turganlarga ta'sir qiladi

Nima ko'rsatdi: 2.5, 2.6-bo'limlar.

Misol 3 — multiprocessing.Pool

python
"""map, imap, imap_unordered, starmap, apply_async; initializer; maxtasksperchild; chunksize ta'siri."""

import multiprocessing as mp
import os
import time

ISHCHI_HOLATI: dict[str, object] = {}


def boshlash(nom: str) -> None:
    """Har ishchi jarayonda bir marta bajariladi."""
    ISHCHI_HOLATI["nom"] = f"{nom}-{os.getpid()}"
    ISHCHI_HOLATI["chaqirilgan"] = int(ISHCHI_HOLATI.get("chaqirilgan", 0)) + 1


def holat_ol(_: int) -> tuple[str, int]:
    return (
        str(ISHCHI_HOLATI.get("nom", "tayyorlanmagan")),
        int(ISHCHI_HOLATI.get("chaqirilgan", 0)),
    )


def kvadrat(x: int) -> int:
    return x * x


def juft_kopaytir(a: int, b: int) -> int:
    return a * b


def mayda(x: int) -> int:
    return x + 1


def pid_ol(_: int) -> int:
    return os.getpid()


def xatoli(x: int) -> int:
    if x == 3:
        raise ValueError("uchinchi element buzilgan")
    return x


def main() -> None:
    print("=== 1. map, imap, imap_unordered ===")
    with mp.Pool(processes=3) as pul:
        print(f"  map:             {pul.map(kvadrat, range(6))}")
        print(f"  imap (tartibda): {list(pul.imap(kvadrat, range(6)))}")
        tartibsiz = list(pul.imap_unordered(kvadrat, range(6)))
        print(f"  imap_unordered:  to'plam bir xil: {sorted(tartibsiz) == [x * x for x in range(6)]}")
        print(f"  starmap:         {pul.starmap(juft_kopaytir, [(2, 3), (4, 5), (6, 7)])}")

    print("\n=== 2. initializer ===")
    with mp.Pool(processes=3, initializer=boshlash, initargs=("ishchi",)) as pul:
        holatlar = pul.map(holat_ol, range(12))
    nomlar = {nom for nom, _ in holatlar}
    chaqiruvlar = {son for _, son in holatlar}
    print(f"  12 vazifa bajarildi, hammasi tayyorlangan ishchida: {'tayyorlanmagan' not in nomlar}")
    print(f"  har ishchida initializer aynan bir marta chaqirilgan: {chaqiruvlar == {1}}")
    print(f"  ishchi nomlari soni 3 tadan oshmaydi: {len(nomlar) <= 3}")
    print(f"  asosiy jarayonda holat bo'sh: {ISHCHI_HOLATI == {}}")

    print("\n=== 3. apply_async va xatolar ===")
    with mp.Pool(processes=2) as pul:
        natija = pul.apply_async(kvadrat, (9,))
        print(f"  apply_async: {natija.get(timeout=20)}, tayyor: {natija.ready()}")
        xato_natija = pul.apply_async(xatoli, (3,))
        try:
            xato_natija.get(timeout=20)
        except ValueError as xato:
            print(f"  xatoli vazifa: ValueError: {xato}")
        try:
            pul.map(xatoli, range(6))
        except ValueError as xato:
            print(f"  map birinchi xatoda to'xtaydi: {xato}")

    print("\n=== 4. maxtasksperchild ===")
    with mp.Pool(processes=2, maxtasksperchild=2) as pul:
        pidlar = pul.map(pid_ol, range(8), chunksize=1)
    print(f"  8 vazifa, 2 ishchi, maxtasksperchild=2 → PID lar soni 2 dan ko'p: {len(set(pidlar)) > 2}")
    print(f"  har PID ko'pi bilan 2 vazifa bajargan: {max(pidlar.count(p) for p in set(pidlar)) <= 2}")
    print("  ⭐ xotira sizishiga qarshi himoya")

    print("\n=== 5. chunksize ta'siri ===")
    MAYDA = list(range(200_000))
    with mp.Pool(processes=4) as pul:
        bosh = time.perf_counter()
        pul.map(mayda, MAYDA, chunksize=1)
        kichik = time.perf_counter() - bosh
        bosh = time.perf_counter()
        pul.map(mayda, MAYDA, chunksize=10_000)
        katta = time.perf_counter() - bosh
    print(f"  chunksize=10 000 chunksize=1 dan kamida 3 barobar tez: {kichik / katta >= 3}")
    print("  ⭐ mayda vazifalarda uzatish narxi hal qiluvchi (12.6-dars)")


if __name__ == "__main__":
    main()

Natijaning muhim qismi:

text
=== 1. map, imap, imap_unordered ===
  map:             [0, 1, 4, 9, 16, 25]
  imap (tartibda): [0, 1, 4, 9, 16, 25]
  imap_unordered:  to'plam bir xil: True
  starmap:         [6, 20, 42]

=== 2. initializer ===
  12 vazifa bajarildi, hammasi tayyorlangan ishchida: True
  har ishchida initializer aynan bir marta chaqirilgan: True
  ishchi nomlari soni 3 tadan oshmaydi: True
  asosiy jarayonda holat bo'sh: True

=== 3. apply_async va xatolar ===
  apply_async: 81, tayyor: True
  xatoli vazifa: ValueError: uchinchi element buzilgan
  map birinchi xatoda to'xtaydi: uchinchi element buzilgan

=== 4. maxtasksperchild ===
  8 vazifa, 2 ishchi, maxtasksperchild=2 → PID lar soni 2 dan ko'p: True
  har PID ko'pi bilan 2 vazifa bajargan: True
  ⭐ xotira sizishiga qarshi himoya

=== 5. chunksize ta'siri ===
  chunksize=10 000 chunksize=1 dan kamida 3 barobar tez: True
  ⭐ mayda vazifalarda uzatish narxi hal qiluvchi (12.6-dars)

Nima ko'rsatdi: 2.3, 2.4-bo'limlar.

Misol 4 — Amaliy: hujjatlarni parallel tahlil qilish

Kirishdagi vaziyat: ko'p hujjatni tahlil qilish. Pul ishchilari og'ir "modelni" bir marta yuklaydi (initializer), vazifalar paketlanadi (chunksize), natijalar tayyor bo'lgani zahoti yig'iladi, xatolar esa hisobotga tushadi. Oxirida iplar bilan solishtiramiz.

python
"""initializer bilan tayyorgarlik; as_completed bilan yig'ish; xatolar hisoboti; iplar va jarayonlar solishtiruvi."""

import concurrent.futures as cf
import hashlib
import os
import time
from dataclasses import dataclass

HUJJATLAR = [f"hujjat-{i:04d}" for i in range(240)]
ISHCHILAR = 4
LUGAT: dict[str, int] = {}


@dataclass(frozen=True, slots=True)
class Tahlil:
    hujjat: str
    ball: int
    ishchi: int


def model_yukla() -> None:
    """Og'ir tayyorgarlik: har ishchida bir marta."""
    global LUGAT
    LUGAT = {f"soz{i}": i for i in range(20_000)}


def tahlil_qil(hujjat: str) -> Tahlil:
    if not LUGAT:
        raise RuntimeError("model yuklanmagan")
    if hujjat.endswith("13"):
        raise ValueError("buzilgan format")
    xesh = hashlib.sha256(hujjat.encode())
    for _ in range(300):
        xesh.update(b"x")
    ball = 0
    for _ in range(6):                       # sof Python hisob — GIL ni band qiladi
        ball = sum(LUGAT.get(f"soz{i}", 0) for i in range(0, 20_000, 3)) % 1000
    return Tahlil(hujjat, ball + xesh.digest()[0], os.getpid())


def quvur(Ex, ishchilar: int = ISHCHILAR) -> tuple[list[Tahlil], list[str], float]:
    bosh = time.perf_counter()
    natijalar: list[Tahlil] = []
    xatolar: list[str] = []
    with Ex(max_workers=ishchilar, initializer=model_yukla) as ex:
        kelajaklar = {ex.submit(tahlil_qil, h): h for h in HUJJATLAR}
        for kelajak in cf.as_completed(kelajaklar, timeout=300):
            hujjat = kelajaklar[kelajak]
            try:
                natijalar.append(kelajak.result())
            except ValueError as xato:
                xatolar.append(f"{hujjat}: {xato}")
    return natijalar, xatolar, time.perf_counter() - bosh


def main() -> None:
    print("=== 1. Jarayonlar puli ===")
    natijalar, xatolar, jarayon_vaqti = quvur(cf.ProcessPoolExecutor)
    print(f"  {len(HUJJATLAR)} hujjat: {len(natijalar)} tahlil, {len(xatolar)} xato")
    ishchilar = {t.ishchi for t in natijalar}
    print(f"  hammasi hisobga olindi: {len(natijalar) + len(xatolar) == len(HUJJATLAR)}")
    print(f"  bir nechta ishchi jarayon ishladi: {1 < len(ishchilar) <= 4}")
    print(f"  asosiy jarayon ishlatilmadi: {os.getpid() not in ishchilar}")
    print(f"  asosiy jarayonda LUGAT bo'sh qoldi: {len(LUGAT) == 0}")

    print("\n=== 2. Xatolar hisoboti ===")
    for xato in sorted(xatolar)[:3]:
        print(f"  {xato}")
    print(f"  jami: {len(xatolar)}")

    print("\n=== 3. Iplar bilan solishtirish ===")
    ip_natijalari, ip_xatolari, ip_vaqti = quvur(cf.ThreadPoolExecutor)
    print(f"  iplar: {len(ip_natijalari)} tahlil, {len(ip_xatolari)} xato")
    print(f"  natijalar mos: {sorted(t.hujjat for t in natijalar) == sorted(t.hujjat for t in ip_natijalari)}")
    print(f"  iplarda bitta jarayon: {len({t.ishchi for t in ip_natijalari}) == 1}")
    print(f"  jarayonlar iplardan kamida 2 barobar tez: {ip_vaqti / jarayon_vaqti >= 2}")
    print("  ⭐ sof Python hisob GIL ni band qiladi — iplar yordam bermaydi")

    print("\n=== 4. initializer nima berdi ===")
    print(f"  model har ishchida bir marta yuklandi (eng ko'pi bilan {ISHCHILAR} marta),")
    print(f"  har vazifada emas ({len(HUJJATLAR)} marta bo'lardi)")
    print(f"  iplar esa asosiy jarayon xotirasini to'ldirdi: {len(LUGAT) > 0}")
    print("  ⚠️ iplar bitta xotirani bo'lishadi — global holat ular uchun umumiy")

    print("\n=== 5. Qoidalar ===")
    print("  ✅ og'ir tayyorgarlik — initializer da")
    print("  ✅ natijalar as_completed bilan, har biri tekshiriladi")
    print("  ✅ xatolar hisobotga tushadi, quvur to'xtamaydi")
    print("  ✅ sof Python hisob — jarayonlar; I/O bo'lsa — iplar (14.8-dars)")


if __name__ == "__main__":
    main()

Natijaning muhim qismi:

text
=== 1. Jarayonlar puli ===
  240 hujjat: 237 tahlil, 3 xato
  hammasi hisobga olindi: True
  bir nechta ishchi jarayon ishladi: True
  asosiy jarayon ishlatilmadi: True
  asosiy jarayonda LUGAT bo'sh qoldi: True

=== 2. Xatolar hisoboti ===
  hujjat-0013: buzilgan format
  hujjat-0113: buzilgan format
  hujjat-0213: buzilgan format
  jami: 3

=== 3. Iplar bilan solishtirish ===
  iplar: 237 tahlil, 3 xato
  natijalar mos: True
  iplarda bitta jarayon: True
  jarayonlar iplardan kamida 2 barobar tez: True
  ⭐ sof Python hisob GIL ni band qiladi — iplar yordam bermaydi

=== 4. initializer nima berdi ===
  model har ishchida bir marta yuklandi (eng ko'pi bilan 4 marta),
  har vazifada emas (240 marta bo'lardi)
  iplar esa asosiy jarayon xotirasini to'ldirdi: True
  ⚠️ iplar bitta xotirani bo'lishadi — global holat ular uchun umumiy

=== 5. Qoidalar ===
  ✅ og'ir tayyorgarlik — initializer da
  ✅ natijalar as_completed bilan, har biri tekshiriladi
  ✅ xatolar hisobotga tushadi, quvur to'xtamaydi
  ✅ sof Python hisob — jarayonlar; I/O bo'lsa — iplar (14.8-dars)

Nima ko'rsatdi: 2.2, 2.4, 2.7-bo'limlar.


5. To'g'ri va noto'g'ri tushunishlar

Noto'g'ri fikr To'g'risi
"Pool va Executor — bir xil" O'xshash, lekin API va imkoniyatlar farq qiladi
"map hamma xatoni qaytaradi" Birinchi istisnoda to'xtaydi
"timeout vazifani to'xtatadi" Faqat kutishni to'xtatadi
"cancel() ishlayotgan vazifani to'xtatadi" Faqat navbatdagini
"chunksize — mayda sozlama" Mayda vazifalarda hal qiluvchi
"initializer har vazifada ishlaydi" Har ishchida bir marta
"Pul yaratish arzon" Jarayonlar puli — millisekundlar; qayta ishlating
"InterpreterPoolExecutor — jarayonlar bilan bir xil" Yengilroq, lekin uzatish cheklovlari bor

6. Keng tarqalgan xatolar va yechimlari

1. Har so'rovda yangi pul

python
def ishlov(sorov):
    with ProcessPoolExecutor() as ex:     # ❌ har safar jarayonlar yaratiladi
        ...
POOL = ProcessPoolExecutor()              # ✅ bir marta

2. map bilan xatolarni yo'qotish

python
natijalar = list(ex.map(ishla, malumot))  # ⚠️ birinchi xatoda to'xtaydi
for k in as_completed(kelajaklar): ...    # ✅ hammasi hisobga olinadi

3. chunksize ni unutish

python
ex.map(mayda_funksiya, 1_000_000 * [...])            # ❌
ex.map(mayda_funksiya, malumot, chunksize=10_000)    # ✅

4. Og'ir tayyorgarlikni har vazifada bajarish

python
def vazifa(x):
    model = model_yukla()                 # ❌ har vazifada
ProcessPoolExecutor(initializer=model_yukla)   # ✅ har ishchida bir marta

5. Natijani o'qimaslik

python
ex.submit(ishla, x)                       # ❌ xato jim qoladi

6. BrokenProcessPool ni hisobga olmaslik

python
# ⚠️ ishchi yiqilsa (segfault, OOM) — butun pul ishlamay qoladi
except cf.process.BrokenProcessPool:      # ✅ qayta yaratish yoki qayd qilish

7. Pul ichidan pulga vazifa berish

python
def vazifa():
    return POOL.submit(boshqa).result()   # ⚠️ deadlock xavfi

8. Katta ma'lumotni har vazifaga uzatish

python
ex.map(ishla, [(katta_jadval, i) for i in ...])   # ❌ (14.4-dars)

7. Integratsiya — bu bilim qayerda kerak bo'ladi

  • 12.6-dars (o'tilgan): bo'lak hajmi va uzatish narxi
  • 14.2-dars (o'tilgan): ThreadPoolExecutor asoslari
  • 14.4-dars (o'tilgan): jarayonlar va ma'lumot almashish
  • 14.7-dars: asyncio bilan pullarni birlashtirish
  • 17-qism: parallel testlar (pytest-xdist)
  • 24-qism: ma'lumot tahlili quvurlari
  • 26-qism: avtomatlashtirish va paketli ishlov
  • 29-qism: miqyoslash va ishchilar

8. Eng yaxshi amaliyotlar

  1. concurrent.futures — sukut tanlov.

  2. Pulni bir marta yarating va qayta ishlating.

  3. as_completed bilan har natijani tekshiring.

  4. chunksize ni ish hajmiga qarab tanlang.

  5. Og'ir tayyorgarlik — initializer da.

  6. maxtasksperchild — uzoq ishlaydigan pullar uchun.

  7. Muddatlar bering va ular vazifani to'xtatmasligini yodda tuting.

  8. BrokenProcessPool ni qayta ishga tushirish bilan boshqaring.


9. Amaliy topshiriq

Vazifa 1: Natijani bashorat qiling

python
import concurrent.futures as cf
1.  with cf.ThreadPoolExecutor(2) as ex:
        print(list(ex.map(abs, [-1, -2, 3])))
2.  with cf.ThreadPoolExecutor(2) as ex:
        f = ex.submit(sum, [1, 2])
    print(f.result(), f.done(), f.cancelled())
3.  with cf.ThreadPoolExecutor(2) as ex:
        f = ex.submit(int, "x")
    print(type(f.exception()).__name__)
4.  with cf.ThreadPoolExecutor(1) as ex:
        a = ex.submit(sum, [1])
        b = ex.submit(sum, [2])
        print(b.cancel() in (True, False))
5.  print(cf.FIRST_COMPLETED, cf.ALL_COMPLETED)
6.  with cf.ThreadPoolExecutor(2) as ex:
        fs = [ex.submit(pow, 2, i) for i in range(3)]
        print(sorted(f.result() for f in cf.as_completed(fs)))
7.  with cf.ThreadPoolExecutor(2) as ex:
        t, q = cf.wait([ex.submit(pow, 2, 3)], timeout=5)
    print(len(t), len(q))
8.  print(hasattr(cf, "InterpreterPoolExecutor"))
9.  import multiprocessing as mp
    with mp.Pool(2) as p:
        print(p.map(abs, [-1, -2]))
10. with mp.Pool(2) as p:
        print(p.starmap(pow, [(2, 3)]))
11. with mp.Pool(2) as p:
        r = p.apply_async(abs, (-5,))
        print(r.get(timeout=10), r.ready())
12. with mp.Pool(2) as p:
        print(sorted(p.imap_unordered(abs, [-3, -1, -2])))
Javoblar
  1. [1, 2, 3]
  2. 3 True False
  3. ValueError
  4. True — cancel() True yoki False qaytaradi
  5. FIRST_COMPLETED ALL_COMPLETED
  6. [1, 2, 4]
  7. 1 0
  8. True — 3.14
  9. [1, 2]
  10. [8]
  11. 5 True
  12. [1, 2, 3]

Vazifa 2: Xatolarni tuzating

python
1.  def ishlov(sorovlar):
        with ProcessPoolExecutor() as ex:
            return list(ex.map(og_ir, sorovlar))

2.  natijalar = list(ex.map(tahlil, hujjatlar))   # ba'zi hujjatlar buzilgan

3.  ex.map(kichik_funksiya, range(1_000_000))

4.  def vazifa(x):
        model = katta_model_yukla()
        return model.bashorat(x)
    ex.map(vazifa, malumot)

5.  for x in malumot:
        ex.submit(saqla, x)
    print("saqlandi")
Javoblar
python
1.  EX = ProcessPoolExecutor()                    # modul darajasida, bir marta
    def ishlov(sorovlar):
        return list(EX.map(og_ir, sorovlar))

2.  kelajaklar = {ex.submit(tahlil, h): h for h in hujjatlar}
    natijalar, xatolar = [], []
    for k in as_completed(kelajaklar):
        try:
            natijalar.append(k.result())
        except Exception as xato:
            xatolar.append((kelajaklar[k], xato))

3.  ex.map(kichik_funksiya, range(1_000_000), chunksize=10_000)

4.  def boshlash():
        global MODEL
        MODEL = katta_model_yukla()
    def vazifa(x):
        return MODEL.bashorat(x)
    with ProcessPoolExecutor(initializer=boshlash) as ex:
        ex.map(vazifa, malumot)

5.  kelajaklar = [ex.submit(saqla, x) for x in malumot]
    for k in as_completed(kelajaklar):
        k.result()                                # xatolar ko'rinadi
    print("saqlandi")

Vazifa 3: Universal xarita funksiyasi

xarita(funksiya, malumot, rejim="auto") yozing:

  1. rejim: "iplar", "jarayonlar", "interpretatorlar", "ketma-ket", "auto"
  2. "auto" da: namunada o'lchab, mos rejimni tanlasin (12.6-dars)
  3. chunksize ni avtomatik hisoblasin
  4. Xatolarni yig'ib, (natijalar, xatolar) qaytarsin
  5. Uch xil yukda sinab, natijalarni jadvalda ko'rsating

Vazifa 4: chunksize egri chizig'i

  1. 1 000 000 mayda vazifa uchun chunksize ni 1 dan 100 000 gacha o'zgartiring
  2. Har qiymatda vaqtni o'lchang
  3. Optimal oraliqni toping va sababini tushuntiring
  4. Vazifa og'irligini oshirib (10 ms), egri chiziq qanday o'zgarishini ko'rsating

Vazifa 5: Ishonchli pul

IshonchliPul klassini yozing:

  1. ProcessPoolExecutor ustida: BrokenProcessPool da avtomatik qayta yaratish
  2. Har vazifa uchun muddat va qayta urinish
  3. Bekor qilish (cancel_futures) bilan xushmuomala to'xtatish
  4. Statistika: bajarilgan, xato, qayta urinish, o'rtacha vaqt
  5. Ishchini ataylab yiqitib (os._exit), qayta tiklanishini ko'rsating

Vazifa 6: Rasm quvuri

  1. 500 ta rasm (yoki sun'iy massiv) uchun: o'qish (I/O) → o'zgartirish (CPU) → saqlash (I/O)
  2. Har bosqich uchun mos pul turini tanlang
  3. Bosqichlarni Queue yoki as_completed bilan ulang
  4. Umumiy o'tkazuvchanlikni (rasm/soniya) o'lchang
  5. Bitta pul bilan hammasini bajarish bilan solishtiring

Vazifa 7: O'ylash

Java'da ExecutorService va ForkJoinPool (ish o'g'irlash bilan), .NET'da Task Parallel Library, Go'da esa qo'lda yoziladigan ishchilar puli. Python'dagi Executor API si ular bilan qanday taqqoslanadi va nega Python'da "ish o'g'irlash" (work stealing) rejalashtiruvchisi yo'q?

Javob

Qisqa javob: Python'ning Executor API si Java'dagi ExecutorService dan ilhomlangan va unga juda o'xshash, lekin ForkJoinPool kabi ish o'g'irlash va rekursiv bo'linishni qo'llab-quvvatlamaydi. Sababi — GIL va jarayonlar modeli: Python'da ishchilar orasida ish o'g'irlash uchun umumiy xotira yo'q (jarayonlar) yoki foyda yo'q (iplar, GIL tufayli).

1. Taqqoslash

Platforma Pul Rejalashtirish
Java ExecutorService, ForkJoinPool Navbat; ForkJoinPool da ish o'g'irlash
.NET Task, ThreadPool Ish o'g'irlash
Go Qo'lda pul Runtime M:N, goroutine lar uchun ish o'g'irlash
Python Executor, Pool Oddiy umumiy navbat

2. Ish o'g'irlash nima

Har ishchining o'z navbati bo'ladi; bo'shagan ishchi boshqasining navbatidan vazifa "o'g'irlaydi". Bu:

  • Notekis vazifalarda yukni yaxshi taqsimlaydi
  • Rekursiv bo'linadigan algoritmlar (fork/join) uchun ideal
  • Umumiy navbatdagi raqobatni kamaytiradi

3. Nega Python'da yo'q

  1. Jarayonlarda umumiy navbat kerak: ishchilar alohida xotira fazasida — vazifani "o'g'irlash" IPC talab qiladi, ya'ni arzon emas
  2. Iplarda foyda kam: GIL tufayli Python kodi baribir navbatma-navbat bajariladi
  3. Vazifalar odatda yirik: Python'da bo'lak hajmi millisekundlarda o'lchanadi (12.6-dars) — umumiy navbat yetarli
  4. Soddalik: Executor API si kichik va tushunarli; murakkab rejalashtiruvchi qo'llab-quvvatlashni qiyinlashtirardi

4. Python'da nima qilish mumkin

Muammo Yechim
Notekis vazifalar chunksize=1 va imap_unordered — tugagan ishchi darhol yangisini oladi
Rekursiv bo'linish Vazifalarni o'zingiz bo'laklab, navbatga qo'ying
Dinamik yuk Navbat asosidagi ishchilar (14.3-dars) — tabiiy yuk muvozanati
Haqiqiy ish o'g'irlash kerak Dask, Ray kabi kutubxonalar (29-qism)

5. Xulosa

  1. Executor — Java ExecutorService ning Python varianti: sodda va yetarli
  2. Ish o'g'irlash Python'da GIL va jarayonlar modeli tufayli foyda bermaydi
  3. Notekis yukda imap_unordered/as_completed + kichik chunksize yetarli muvozanat beradi
  4. Katta miqyosda — taqsimlangan freymvorklar (Dask, Ray)

Nimani mustahkamlaydi: 2.1–2.7-bo'limlar.


Xulosa

Bu darsda pullar bilan ishlashni o'rgandik.

Eng muhim uch fikr:

  1. concurrent.futures — bitta API, uch xil ishchi. ThreadPoolExecutor, ProcessPoolExecutor va 3.14 dagi InterpreterPoolExecutor bir xil interfeysga ega: submit → Future, map → tartibda natijalar, as_completed → tugash tartibida. Model almashtirish uchun faqat klass nomini o'zgartirasiz — qolgan kod o'zgarmaydi.

  2. Xatolar va muddatlar aniq qoidalarga bo'ysunadi. Istisno Future ichida saqlanadi va result() da qayta ko'tariladi; map esa birinchi istisnoda to'xtaydi — hamma vazifani hisobga olish uchun as_completed ishlating. timeout faqat kutishni to'xtatadi, vazifani emas; cancel() esa faqat hali boshlanmagan vazifani bekor qiladi (shutdown(cancel_futures=True) — navbatdagilarning hammasini).

  3. Pool ning qo'shimcha imkoniyatlari. imap_unordered — natijalarni tayyor bo'lgani zahoti beradi, starmap — ko'p argumentli funksiyalar uchun, initializer — og'ir tayyorgarlikni har ishchida bir marta bajaradi, maxtasksperchild — ishchilarni vaqti-vaqti bilan yangilab, xotira sizishidan himoya qiladi. chunksize esa mayda vazifalarda hal qiluvchi ahamiyatga ega.

Keyingi darsda uchinchi modelga o'tamiz: asyncio hodisa sikli — u qanday ishlaydi, korutinalar qanday rejalashtiriladi va nega minglab ulanish uchun eng arzon yechim shu.

Ulashish:Telegram'da

Izohlar (0)

Izoh yozish uchun kiring.

  • Hozircha izoh yo'q. Birinchi bo'ling!
14.5-dars: Pool va concurrent.futures — IlmHamroh