Mundarija (21)
- 1. Kirish va motivatsiya
- 2. Nazariya — chuqur tushuntirish
- 2.1. Ikki API
- 2.2. Executor API
- 2.3. multiprocessing.Pool
- 2.4. chunksize va initializer
- 2.5. Xatolar va muddatlar
- 2.6. To'xtatish
- 2.7. Qaysi API ni tanlash
- 3. Tez ma'lumotnoma
- 4. Batafsil misollar
- Misol 1 — Executor asoslari
- Misol 2 — Xatolar, muddatlar va to'xtatish
- Misol 3 — multiprocessing.Pool
- Misol 4 — Amaliy: hujjatlarni parallel tahlil qilish
- 5. To'g'ri va noto'g'ri tushunishlar
- 6. Keng tarqalgan xatolar va yechimlari
- 7. Integratsiya — bu bilim qayerda kerak bo'ladi
- 8. Eng yaxshi amaliyotlar
- 9. Amaliy topshiriq
- Xulosa
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_asyncchunksize,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
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:
cf.as_completed(kelajaklar, timeout=30) # tugash tartibida
tugagan, qolgan = cf.wait(kelajaklar, timeout=5, return_when=cf.FIRST_COMPLETED)2.3. multiprocessing.Pool
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.
def boshlash(yol: str) -> None:
global MODEL
MODEL = ogir_model_yukla(yol) # har vazifada emas, bir marta2.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
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
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'xtatmaydi4. Batafsil misollar
Misol 1 — Executor asoslari
"""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:
=== 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 almashadiNima ko'rsatdi: 2.1, 2.2-bo'limlar.
Misol 2 — Xatolar, muddatlar va to'xtatish
"""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:
=== 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 qiladiNima ko'rsatdi: 2.5, 2.6-bo'limlar.
Misol 3 — multiprocessing.Pool
"""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:
=== 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.
"""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:
=== 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
def ishlov(sorov):
with ProcessPoolExecutor() as ex: # ❌ har safar jarayonlar yaratiladi
...
POOL = ProcessPoolExecutor() # ✅ bir marta2. map bilan xatolarni yo'qotish
natijalar = list(ex.map(ishla, malumot)) # ⚠️ birinchi xatoda to'xtaydi
for k in as_completed(kelajaklar): ... # ✅ hammasi hisobga olinadi3. chunksize ni unutish
ex.map(mayda_funksiya, 1_000_000 * [...]) # ❌
ex.map(mayda_funksiya, malumot, chunksize=10_000) # ✅4. Og'ir tayyorgarlikni har vazifada bajarish
def vazifa(x):
model = model_yukla() # ❌ har vazifada
ProcessPoolExecutor(initializer=model_yukla) # ✅ har ishchida bir marta5. Natijani o'qimaslik
ex.submit(ishla, x) # ❌ xato jim qoladi6. BrokenProcessPool ni hisobga olmaslik
# ⚠️ ishchi yiqilsa (segfault, OOM) — butun pul ishlamay qoladi
except cf.process.BrokenProcessPool: # ✅ qayta yaratish yoki qayd qilish7. Pul ichidan pulga vazifa berish
def vazifa():
return POOL.submit(boshqa).result() # ⚠️ deadlock xavfi8. Katta ma'lumotni har vazifaga uzatish
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):
ThreadPoolExecutorasoslari - 14.4-dars (o'tilgan): jarayonlar va ma'lumot almashish
- 14.7-dars:
asynciobilan 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
concurrent.futures— sukut tanlov.Pulni bir marta yarating va qayta ishlating.
as_completedbilan har natijani tekshiring.chunksizeni ish hajmiga qarab tanlang.Og'ir tayyorgarlik —
initializerda.maxtasksperchild— uzoq ishlaydigan pullar uchun.Muddatlar bering va ular vazifani to'xtatmasligini yodda tuting.
BrokenProcessPoolni qayta ishga tushirish bilan boshqaring.
9. Amaliy topshiriq
Vazifa 1: Natijani bashorat qiling
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, 2, 3]3 True FalseValueErrorTrue—cancel()TrueyokiFalseqaytaradiFIRST_COMPLETED ALL_COMPLETED[1, 2, 4]1 0True— 3.14[1, 2][8]5 True[1, 2, 3]
Vazifa 2: Xatolarni tuzating
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
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:
rejim:"iplar","jarayonlar","interpretatorlar","ketma-ket","auto""auto"da: namunada o'lchab, mos rejimni tanlasin (12.6-dars)chunksizeni avtomatik hisoblasin- Xatolarni yig'ib,
(natijalar, xatolar)qaytarsin - Uch xil yukda sinab, natijalarni jadvalda ko'rsating
Vazifa 4: chunksize egri chizig'i
- 1 000 000 mayda vazifa uchun
chunksizeni 1 dan 100 000 gacha o'zgartiring - Har qiymatda vaqtni o'lchang
- Optimal oraliqni toping va sababini tushuntiring
- Vazifa og'irligini oshirib (10 ms), egri chiziq qanday o'zgarishini ko'rsating
Vazifa 5: Ishonchli pul
IshonchliPul klassini yozing:
ProcessPoolExecutorustida:BrokenProcessPoolda avtomatik qayta yaratish- Har vazifa uchun muddat va qayta urinish
- Bekor qilish (
cancel_futures) bilan xushmuomala to'xtatish - Statistika: bajarilgan, xato, qayta urinish, o'rtacha vaqt
- Ishchini ataylab yiqitib (
os._exit), qayta tiklanishini ko'rsating
Vazifa 6: Rasm quvuri
- 500 ta rasm (yoki sun'iy massiv) uchun: o'qish (I/O) → o'zgartirish (CPU) → saqlash (I/O)
- Har bosqich uchun mos pul turini tanlang
- Bosqichlarni
Queueyokias_completedbilan ulang - Umumiy o'tkazuvchanlikni (rasm/soniya) o'lchang
- 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
- Jarayonlarda umumiy navbat kerak: ishchilar alohida xotira fazasida — vazifani "o'g'irlash" IPC talab qiladi, ya'ni arzon emas
- Iplarda foyda kam: GIL tufayli Python kodi baribir navbatma-navbat bajariladi
- Vazifalar odatda yirik: Python'da bo'lak hajmi millisekundlarda o'lchanadi (12.6-dars) — umumiy navbat yetarli
- Soddalik:
ExecutorAPI 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
Executor— JavaExecutorServicening Python varianti: sodda va yetarli- Ish o'g'irlash Python'da GIL va jarayonlar modeli tufayli foyda bermaydi
- Notekis yukda
imap_unordered/as_completed+ kichikchunksizeyetarli muvozanat beradi - 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:
concurrent.futures— bitta API, uch xil ishchi.ThreadPoolExecutor,ProcessPoolExecutorva 3.14 dagiInterpreterPoolExecutorbir xil interfeysga ega:submit→Future,map→ tartibda natijalar,as_completed→ tugash tartibida. Model almashtirish uchun faqat klass nomini o'zgartirasiz — qolgan kod o'zgarmaydi.Xatolar va muddatlar aniq qoidalarga bo'ysunadi. Istisno
Futureichida saqlanadi varesult()da qayta ko'tariladi;mapesa birinchi istisnoda to'xtaydi — hamma vazifani hisobga olish uchunas_completedishlating.timeoutfaqat kutishni to'xtatadi, vazifani emas;cancel()esa faqat hali boshlanmagan vazifani bekor qiladi (shutdown(cancel_futures=True)— navbatdagilarning hammasini).Poolning 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.chunksizeesa 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.
Izohlar (0)
Izoh yozish uchun kiring.
- Hozircha izoh yo'q. Birinchi bo'ling!