IlmHamroh
Data Science va sun'iy intellekt/MLOps va deploy3/14-dars39 daqiqa
Mundarija (25)

27.3-dars: Ma'lumot quvuri

27-QISM — MLOPS VA DEPLOY · 3-dars


1. Kirish va motivatsiya

27.1-darsda ML tizimida modeldan ko'ra uning atrofidagi quvur kattaroq ekanini ko'rdik, 27.2-darsda esa natijani qayta olish uchun ma'lumot xeshi kerakligini. Endi o'sha quvurning birinchi va eng ko'p buziladigan qismini quramiz: ma'lumotni manbadan olib, tekshirib, o'zgartirib, omborga yuklaydigan ma'lumot quvuri.

Ishlab chiqarishdagi ML nosozliklarining katta qismi aynan shu yerda tug'iladi. Quvur kechasi ishlaydi va ertalab hech kim unga qaramaydi. U tarmoq uzilganda yarim yo'lda to'xtaydi. Uni qo'lda qayta ishga tushirishadi — va jadvalda har buyurtma ikki marta bo'lib qoladi. Manba tizim bir ustunni summa dan miqdor ga qayta nomlaydi — va quvur yo xato beradi, yo, eng yomoni, bo'sh ustun bilan davom etadi. Kechikib kelgan yozuvlar hech qachon omborga tushmaydi. Model esa shu omborga ishonadi.

Yaxshi ma'lumot quvuri to'rt xususiyatga ega: bosqichlari va ularning bog'liqliklari aniq (DAG); qayta ishga tushirish xavfsiz (idempotent); faqat yangi ma'lumotni o'qiydi (inkremental) va buzilgan partiyani omborga kiritmaydi (sifat tekshiruvi va karantin). Bu darsda shularning har birini noldan, sqlite3, pandas va hashlib bilan quramiz. Airflow, Prefect yoki great_expectations kabi vositalar xuddi shu g'oyalarni katta miqyosda amalga oshiradi — ularni nazariya qismida ko'rsatamiz.

Real vaziyat. Chakana savdo kompaniyasida talab modeli bir hafta davomida sotuvni ikki baravar ko'p bashorat qildi va omborlar ortiqcha tovar bilan to'ldi. Sabab: tungi yuklash tarmoq xatosi bilan yiqilgan, operator uni qo'lda qayta ishga tushirgan va birinchi urinishda allaqachon yozilgan yarim partiya ikkinchi marta qo'shilgan. Jadvalda asosiy kalit yo'q edi, yuklash tranzaksiyasiz edi, dublikat tekshiruvi ham yo'q edi. Uchta arzon himoyadan birortasi ham bu haftani saqlab qolardi.

Bu darsda qayta ishga tushirilganda xavfsiz, buzilgan ma'lumotni o'tkazmaydigan va nima qilganini aniq aytadigan quvur quramiz.

Bu darsda:

  • ETL va ELT
  • Quvur — bosqichlar grafi (DAG) va topologik tartib
  • Idempotentlik: UPSERT, tranzaksiya, atomik yozish
  • Inkremental yuklash: watermark va kechikkan yozuvlar
  • Ma'lumot sifati tekshiruvlari va karantin
  • Keshlash va qayta urinish
  • Orkestratorlar: Airflow va Prefect (nazariyada)
  • Quvurdagi sizish (qisqa)
  • Amaliy: kichik freymvork, idempotent yuklash, sifat darvozasi, kesh

ℹ Misollar real pandas/sqlite3/sklearn bilan (Python 3.14). Kutish va vaqt — "virtual soat" bilan; fayllar faqat vaqtinchalik papkada.


2. Nazariya — chuqur tushuntirish

2.1. ETL va ELT

text
ETL (Extract -> Transform -> Load):
  manba --olish--> [quvur serverida o'zgartirish] --yuklash--> ombor
  + omborga faqat toza, tayyor ma'lumot tushadi
  - xom ma'lumot saqlanmaydi -> qayta o'zgartirish uchun manbaga qaytish kerak
  qachon: ombor qimmat yoki cheklangan, o'zgartirish murakkab (python, ML)

ELT (Extract -> Load -> Transform):
  manba --olish--> ombor (XOM qatlam) --SQL bilan o'zgartirish--> toza qatlam
  + xom ma'lumot saqlanadi -> transformatsiyani istalgan vaqtda qayta ishlatish
  + ombor kuchi (ustunli bazalar) ishlatiladi
  - xom qatlamda shaxsiy ma'lumot ham bor -> ruxsatlar muhim
  qachon: zamonaviy bulut omborlari, dbt kabi vositalar

ML UCHUN AMALIY SXEMA (qatlamlar):
  xom (bronze)    manbadan kelganidek, o'zgartirilmaydi, sana bo'yicha bo'lingan
  toza (silver)   tekshirilgan, turlari to'g'ri, dublikatsiz
  belgilar (gold) modelga tayyor jadval / feature store

Qaysi biri bo'lishidan qat'i nazar, bitta qoida umumiy: xom ma'lumotni o'zgartirmasdan saqlang. Bugun o'zgartirishda xato topilsa, ertaga uni qayta hisoblash uchun xom nusxa kerak bo'ladi. Bu 27.2 dagi reproduksiyaning ma'lumot tomonidagi davomi.

Xom qatlam — faqat qo'shiladi, hech qachon ustiga yozilmaydi. Qolgan hamma narsa undan qayta hisoblanishi mumkin.

2.2. Quvur — bosqichlar grafi (DAG)

text
extract ----> validate ----+
                           +--> transform --+--> load
valyuta_kursi -------------+                +--> hisobot

DAG = Directed Acyclic Graph (yo'naltirilgan, siklsiz graf)
  tugun  = bosqich (funksiya)
  qirra  = "B uchun A ning natijasi kerak"
  sikl   = xato (A -> B -> A: hech biri boshlanolmaydi)

TOPOLOGIK TARTIB (Kahn algoritmi):
  1. kirish qirrasi yo'q tugunlar -> "tayyor" navbat
  2. navbatdan olib bajarish; uning bolalaridan qirrani olib tashlash
  3. kirish qirrasi 0 ga tushgan bola -> navbatga
  4. tugunlar tugamasa -> SIKL bor
  determinizm uchun: navbat har doim sorted()

DAG NIMA BERADI:
  bajarish tartibi avtomatik
  mustaqil tarmoqlar parallel bajarilishi mumkin
  bitta bosqich o'zgarsa -> faqat uning PASTKI tarmog'i qayta hisoblanadi
  yiqilgan joydan davom ettirish

Bog'liqlikni kod tartibi emas, graf belgilaydi. Skriptdagi qatorlar tartibi yashirin bog'liqlik; DAG da u ochiq yozilgan.

2.3. Idempotentlik

Idempotent amal — bir marta bajarilgani bilan ko'p marta bajarilgani natijasi bir xil. Quvur uchun bu "qayta ishga tushirish xavfsiz" degani — va quvurlar doim qayta ishga tushiriladi: xatodan keyin, orkestrator qayta urinishida, qo'lda tuzatishdan keyin.

text
IDEMPOTENT EMAS:
  INSERT INTO t VALUES (...)          -> har yurishda dublikat
  df.to_csv(yol, mode="a")            -> fayl o'sib boradi
  hisoblagich += partiya_soni         -> ikki marta sanaladi

IDEMPOTENT:
  1. UPSERT (tabiiy kalit bo'yicha)
     INSERT ... ON CONFLICT(id) DO UPDATE SET ...
  2. BO'LIMNI ALMASHTIRISH
     BEGIN; DELETE WHERE kun = '2026-09-24'; INSERT ...; COMMIT;
  3. ATOMIK FAYL YOZISH
     vaqtinchalik faylga yozish -> os.replace(tmp, yakuniy)
     (o'quvchi yarim faylni hech qachon ko'rmaydi)

TRANZAKSIYA - "hammasi yoki hech narsa":
  yuklash o'rtasida uzilish -> ROLLBACK -> jadval oldingi holatda
  tranzaksiyasiz -> YARIM holat (qaysi qatorlar yozilgani noma'lum)

Idempotentlik va tranzaksiya birgalikda ishlaydi: tranzaksiya yiqilgan yurishdan keyin toza holatni kafolatlaydi, idempotentlik esa qayta yurish hech narsani ikki marta qo'shmasligini.

Har yuklash — tabiiy kalit bo'yicha UPSERT yoki bo'limni almashtirish, har doim tranzaksiya ichida.

2.4. Inkremental yuklash

text
TO'LIQ YUKLASH:  har yurishda butun manbani o'qish
  + oddiy, har doim to'g'ri
  - manba katta bo'lsa sekin va qimmat

INKREMENTAL (WATERMARK):
  holat jadvalida: watermark = oxirgi yuklangan "yangilangan" vaqti
  yurish: SELECT * FROM manba WHERE yangilangan > watermark
  yuklash (UPSERT) va watermark ni yangilash - BITTA tranzaksiyada

KECHIKKAN YOZUVLAR:
  yozuv vaqti 32, lekin manbaga watermark 34 bo'lgandan KEYIN keldi
  -> "yangilangan > 34" uni hech qachon olmaydi (JIM yo'qotish)
  yechim: orqaga qarash oynasi
    WHERE yangilangan > watermark - oyna
    + UPSERT (qayta o'qilganlar dublikat bermaydi - idempotentlik!)
  oyna kattaligi: kechikishlar taqsimotidan (masalan, 99.9-kvantil)

BOSHQA USULLAR:
  CDC (change data capture) - manba bazasi jurnalidan o'zgarishlar
  o'chirilgan qatorlar: "yumshoq o'chirish" ustuni yoki CDC

Watermark — manba vaqti bo'yicha, orqaga qarash oynasi bilan, UPSERT bilan birga. Oynasiz inkremental yuklash kechikkan ma'lumotni jim yo'qotadi.

2.5. Ma'lumot sifati tekshiruvlari

text
TEKSHIRUV TURLARI (great_expectations "expectation" g'oyasi):
  sxema        ustunlar nomi, tartibi, turlari
  to'liqlik    null ulushi <= chegara
  diapazon     summa >= 0, yosh 18..100
  noyoblik     id takrorlanmaydi
  to'plam      hudud ruxsat etilgan ro'yxatda
  taqsimot     kunlik o'rtacha tarixiy [m - 4s, m + 4s] ichida
  hajm         qatorlar soni kutilgan oraliqda
  yangilik     eng so'nggi yozuv 24 soatdan eski emas
  munosabat    har buyurtma_id mijozlar jadvalida bor

DARAJALAR:
  XATO          -> partiya yuklanmaydi, KARANTINGA
  OGOHLANTIRISH -> yuklanadi, signal (alert) yuboriladi

CHEGARALARNI QAYERDAN OLISH:
  qat'iy qoidalar - soha bilimidan (summa >= 0)
  statistik chegaralar - TARIXIY TOZA ma'lumotdan (mos yozuv)
  keyin: normal partiyalarda yolg'on signal sonini TEKSHIRISH

Muhim nozik jihat: tekshiruvlar ketma-ket bog'liq. Sxema buzilgan bo'lsa (masalan, summa matn bo'lib kelgan), diapazon tekshiruvini bajarib bo'lmaydi — u o'zi xato beradi. Shuning uchun sxema birinchi va u o'tmasa, qolganlari bajarilmaydi. Taqsimot tekshiruvlari esa sxema va diapazonga qaraganda "yumshoqroq" — o'rtachaning o'zgarishi haqiqiy biznes o'zgarishi ham bo'lishi mumkin, shuning uchun ular odatda ogohlantirish darajasida bo'ladi.

Tekshiruv chegaralari tarixiy toza ma'lumotdan olinadi va normal partiyalarda yolg'on signal soni o'lchanadi. Har kuni yolg'on signal beradigan tekshiruvni jamoa bir haftada o'chirib qo'yadi.

2.6. Karantin

text
partiya -> tekshiruv --o'tdi--> ombor
                    \
                     --XATO--> karantin/
                               <partiya>.csv   (xom nusxa, o'zgartirilmagan)
                               <partiya>.json  (qaysi tekshiruv, tafsilot)
                               -> alert, egasiga xabar

KARANTINDAN CHIQISH:
  sabab tuzatiladi (manba tizimda yoki tuzatish skriptida)
  partiya QAYTA tekshiriladi (xuddi shu tekshiruvlar bilan)
  o'tsa - odatiy yo'l bilan yuklanadi (idempotent!)
  karantin hech qachon "qo'lda omborga ko'chirish" bilan tugamaydi

Karantinning qiymati — quvur to'xtamaydi: bitta buzilgan partiya boshqa partiyalarni ushlab turmaydi, lekin omborni ham buzmaydi. Buzilgan partiya dalil sifatida saqlanadi — manba egasiga aniq "20% summa bo'sh keldi" deb aytish mumkin.

2.7. Keshlash

text
KESH KALITI = xesh(bosqich nomi + KOD VERSIYASI + KIRISH XESHLARI + parametrlar)

  kirish o'zgarmagan, parametr o'zgarmagan -> keshdan
  parametr o'zgardi (max_iter) -> faqat shu bosqich va pastkisi
  kirish o'zgardi (1 qiymat) -> hammasi qayta
  bosqich kodi o'zgardi -> VERSIYANI oshirish kerak (aks holda eski kesh!)

TUZOQLAR:
  kalitda kod versiyasi yo'q -> tuzatilgan kod eski natijani qaytaradi
  kalitda parametr yo'q -> boshqa sozlama bilan eski natija
  kalit fayl nomi/sanasidan (xeshdan emas) -> o'zgargan fayl sezilmaydi
  kesh ichidagi pickle - faqat ishonchli joydan (19.6)

Bu g'oya DVC (dvc repro), Prefect (cache_key_fn), sklearn ning Pipeline(memory=...) va Make ning asosida yotadi.

2.8. Qayta urinish

text
VAQTINCHALIK xatolar (timeout, 503, ulanish uzildi) -> qayta urinish
DOIMIY xatolar (sxema xato, 401 ruxsat yo'q, 404) -> darhol to'xtash

EKSPONENSIAL KUTISH:
  kutish_i = asos * 2^(i-1)       1, 2, 4, 8 ... soniya
  + jitter (tasodifiy qo'shimcha) -> ko'p mijoz bir vaqtda urinmaydi
  + maksimal urinishlar soni -> keyin xato KO'TARILADI (yashirilmaydi)

SHART: qayta urinadigan bosqich IDEMPOTENT bo'lishi kerak
  aks holda: 1-urinish yarim yozdi -> 2-urinish yana yozdi -> dublikat

25.6-darsda LLM API uchun xuddi shu retry/backoff ni ko'rgan edik; ma'lumot quvurida u manbadan o'qishda va omborga yozishda ishlatiladi.

2.9. Orkestratorlar: Airflow va Prefect

Bu darsdagi kichik freymvork bitta jarayonda ishlaydi. Ishlab chiqarishda orkestrator quvurni jadval bo'yicha ishga tushiradi, bosqichlarni turli mashinalarda bajaradi, qayta urinadi, holatni saqlaydi va interfeys beradi. Quyidagi kod faqat ma'lumotnoma — mashinada airflow va prefect o'rnatilmagan.

python
# Airflow: DAG - fayl, bosqichlar - operatorlar, jadval - cron
from datetime import datetime, timedelta

from airflow import DAG
from airflow.operators.python import PythonOperator

with DAG(
    dag_id="buyurtmalar_kunlik",
    start_date=datetime(2026, 9, 1),
    schedule="0 2 * * *",                 # har kuni 02:00 da
    catchup=False,
    default_args={"retries": 3, "retry_delay": timedelta(minutes=5),
                  "retry_exponential_backoff": True},
) as dag:
    olish = PythonOperator(task_id="extract", python_callable=extract)
    tekshirish = PythonOperator(task_id="validate", python_callable=validate)
    yuklash = PythonOperator(task_id="load", python_callable=load)
    olish >> tekshirish >> yuklash        # bog'liqliklar - DAG qirralari
python
# Prefect: oddiy python funksiyalari + dekoratorlar
from datetime import timedelta

from prefect import flow, task
from prefect.tasks import task_input_hash


@task(retries=3, retry_delay_seconds=[1, 2, 4])
def extract(sana):
    ...


@task(cache_key_fn=task_input_hash, cache_expiration=timedelta(days=1))
def transform(df):
    ...


@flow(name="buyurtmalar_kunlik")
def quvur(sana: str):
    df = extract(sana)
    toza = transform(df)
    load(toza)
Jihat Bu darsdagi freymvork Airflow Prefect
DAG Dekorator + Kahn Operatorlar, >> Funksiya chaqiruvlari
Jadval Yo'q Cron, sensorlar Deployment jadvali
Qayta urinish qayta_urin retries, retry_delay @task(retries=...)
Kesh Kirish xeshi Yo'q (idempotent vazifalar kutiladi) cache_key_fn
Holat va interfeys Yo'q Veb interfeys, loglar Veb interfeys, loglar

Orkestrator idempotentlik va sifat tekshiruvini o'rnini bosmaydi — u faqat bosqichlarni ishga tushiradi va qayta uradi. Bosqichning o'zi xavfsiz bo'lishi sizning ishingiz.

2.10. Quvurdagi sizish

18-qismda baholashdagi sizish (leakage) ni ko'rgan edik. Quvurda u ikki joyda yashirinadi:

text
1. VAQT BO'YICHA BO'LMASLIK
   tasodifiy bo'lish: model kelajak qatorlarining "qo'shnilari"ni ko'radi
   -> validatsiya kelajakni emas, INTERPOLYATSIYANI o'lchaydi
   to'g'ri: o'quv < validatsiya < test, vaqt bo'yicha 18.2-bob

2. XIZMAT VAQTIDA MAVJUD BO'LMAGAN BELGI
   rolling(7, center=True) -> bugungi belgi ertangi 3 kunni ko'radi
   o'qitishda bor, ishlab chiqarishda YO'Q -> training-serving skew 27.1-bob
   to'g'ri: faqat shift(1) dan keyingi (o'tmish) oynalar

QOIDA: har belgi uchun "bashorat qilinayotgan paytda bu qiymat
       ma'lummi?" savoli - quvur kodida, sharh bilan

2.11. Tuzoqlar

Asosiy tuzoqlar: asosiy kalitsiz jadvalga INSERT; tranzaksiyasiz yuklash; df.to_csv(mode="a"); watermark ni yuklashdan alohida yangilash; orqaga qarash oynasisiz inkremental yuklash; sxema tekshiruvisiz diapazon tekshiruvi; chegaralarni taxmin bilan qo'yish; buzilgan partiyani qo'lda omborga ko'chirish; kesh kalitida kod versiyasi yoki parametr yo'qligi; doimiy xatoga qayta urinish; idempotent bo'lmagan bosqichni qayta urish; with sqlite3.connect(...) ulanishni yopadi deb o'ylash (u faqat tranzaksiyani yakunlaydi); vaqt qatorida tasodifiy bo'lish.


3. Tez ma'lumotnoma

python
import hashlib
import json
import os
import sqlite3
from contextlib import closing

# 1. Idempotent yuklash: UPSERT + tranzaksiya; ulanishni YOPISH - closing
UPSERT = ("INSERT INTO ombor VALUES (?, ?, ?, ?) ON CONFLICT(id) DO UPDATE SET "
          "summa=excluded.summa, holat=excluded.holat, yangilangan=excluded.yangilangan")


def yukla(db_yol, qatorlar):
    with closing(sqlite3.connect(db_yol)) as db:   # closing - ulanishni yopadi
        with db:                                   # tranzaksiya: commit / rollback
            db.executemany(UPSERT, qatorlar)


# 2. Watermark + orqaga qarash oynasi (bitta tranzaksiyada)
def inkremental(db, manba, oyna):
    wm = db.execute("SELECT qiymat FROM holat WHERE kalit='wm'").fetchone()[0]
    yangi = [q for q in manba if q[3] > wm - oyna]
    with db:
        db.executemany(UPSERT, yangi)
        db.execute("UPDATE holat SET qiymat=? WHERE kalit='wm'",
                   (max([wm] + [q[3] for q in yangi]),))


# 3. Atomik fayl yozish
def atomik_yoz(yol, matn):
    tmp = f"{yol}.tmp"
    with open(tmp, "w", encoding="utf-8") as f:
        f.write(matn)
    os.replace(tmp, yol)                           # o'quvchi yarim faylni ko'rmaydi


# 4. Kesh kaliti
def kesh_kaliti(bosqich, versiya, kirish_xeshlari, param):
    s = json.dumps({"b": bosqich, "v": versiya, "k": kirish_xeshlari, "p": param},
                   sort_keys=True)
    return hashlib.sha256(s.encode()).hexdigest()[:16]


# 5. Eksponensial kutish bilan qayta urinish (faqat vaqtinchalik xatolar)
def qayta_urin(fn, kut, urinish=4, asos=1.0):
    for i in range(1, urinish + 1):
        try:
            return fn()
        except (ConnectionError, TimeoutError):
            if i == urinish:
                raise
            kut(asos * 2 ** (i - 1))

Tuzilma xulosasi

text
ETL/ELT: xom qatlam saqlanadi, qolgani undan qayta hisoblanadi
DAG: bog'liqliklar ochiq, topologik tartib, faqat pastki tarmoq qayta
idempotent: UPSERT / bo'limni almashtirish / atomik fayl + tranzaksiya
inkremental: watermark + orqaga oyna + UPSERT, bitta tranzaksiyada
sifat: sxema -> to'liqlik -> diapazon -> noyoblik -> to'plam -> taqsimot
karantin: xom nusxa + hisobot; tuzatib QAYTA tekshirish
kesh: kod versiyasi + kirish xeshi + parametr
retry: vaqtinchalik xatoga, eksponensial, idempotent bosqichga

4. Batafsil misollar

Misollar real pandas/sqlite3/sklearn bilan (Python 3.14). Kutishlar virtual soatda hisoblanadi — hech narsa haqiqatan kutmaydi.

Misol 1 — Kichik quvur freymvorki: bosqichlar, DAG, topologik tartib, qayta urinish

Bosqichlar — oddiy funksiyalar, bog'liqliklar — dekorator argumentlari. Freymvork Kahn algoritmi bilan tartibni hisoblaydi, siklni va noma'lum bog'liqlikni ushlaydi, bitta bosqich o'zgarganda qayta hisoblanadigan pastki tarmoqni topadi. Manba — ikki marta "timeout" beradigan beqaror API.

python
"""Kichik quvur freymvorki: bosqichlar, DAG, topologik tartib, qayta urinish."""

import sqlite3
import tempfile
from contextlib import closing
from pathlib import Path

import numpy as np
import pandas as pd


class VirtualSoat:
    """Haqiqiy kutish o'rniga vaqtni hisoblaydi (deterministik)."""

    def __init__(self):
        self.t = 0.0

    def kut(self, soniya):
        self.t += soniya


class Quvur:
    def __init__(self):
        self.bosqichlar = {}          # nom -> (funksiya, bog'liqliklar)

    def bosqich(self, *boglar):
        def bezak(fn):
            self.bosqichlar[fn.__name__] = (fn, list(boglar))
            return fn
        return bezak

    def tartib(self):
        """Kahn algoritmi; bir vaqtda tayyor bo'lganlar alifbo tartibida."""
        kirish = {n: len(b) for n, (_, b) in self.bosqichlar.items()}
        bolalar = {n: [] for n in self.bosqichlar}
        for n, (_, boglar) in self.bosqichlar.items():
            for b in boglar:
                if b not in self.bosqichlar:
                    raise KeyError(f"{n}: noma'lum bog'liqlik '{b}'")
                bolalar[b].append(n)
        tayyor = sorted(n for n, k in kirish.items() if k == 0)
        natija = []
        while tayyor:
            n = tayyor.pop(0)
            natija.append(n)
            for c in sorted(bolalar[n]):
                kirish[c] -= 1
                if kirish[c] == 0:
                    tayyor.append(c)
            tayyor.sort()
        if len(natija) != len(self.bosqichlar):
            qolgan = sorted(set(self.bosqichlar) - set(natija))
            raise ValueError(f"sikl bor: {qolgan}")
        return natija

    def pastki(self, nom):
        """nom va unga (bilvosita) bog'liq barcha bosqichlar."""
        natija = {nom}
        ozgardi = True
        while ozgardi:
            ozgardi = False
            for n, (_, boglar) in self.bosqichlar.items():
                if n not in natija and natija & set(boglar):
                    natija.add(n)
                    ozgardi = True
        return [n for n in self.tartib() if n in natija]

    def ishga_tushir(self):
        natijalar, jurnal = {}, []
        for n in self.tartib():
            fn, boglar = self.bosqichlar[n]
            natijalar[n] = fn(*[natijalar[b] for b in boglar])
            jurnal.append(n)
        return natijalar, jurnal


def qayta_urin(fn, soat, urinish=4, asos=1.0, jurnal=None):
    """Eksponensial kutish bilan qayta urinish (faqat vaqtinchalik xato)."""
    for i in range(1, urinish + 1):
        try:
            return fn()
        except ConnectionError as xato:
            if i == urinish:
                raise
            kutish = asos * 2 ** (i - 1)
            if jurnal is not None:
                jurnal.append(f"urinish {i}: {xato}; {kutish:.0f} s kutish")
            soat.kut(kutish)


def main() -> None:
    with tempfile.TemporaryDirectory() as t:
        papka = Path(t)
        soat = VirtualSoat()
        urinishlar = {"n": 0}
        qayta_jurnal = []

        def manba():                      # beqaror API: 2 marta xato, keyin ishlaydi
            urinishlar["n"] += 1
            if urinishlar["n"] <= 2:
                raise ConnectionError("timeout")
            rng = np.random.default_rng(0)
            return pd.DataFrame({"buyurtma_id": np.arange(1, 201),
                                 "summa": rng.gamma(2.0, 50_000, 200).round(-2),
                                 "hudud": rng.choice(["A", "B", "C"], 200)})

        q = Quvur()

        @q.bosqich()
        def extract():
            return qayta_urin(manba, soat, jurnal=qayta_jurnal)

        @q.bosqich()
        def valyuta_kursi():
            return {"USD": 12_650.0}      # faraziy kurs

        @q.bosqich("extract")
        def validate(df):
            assert df["buyurtma_id"].is_unique, "ID takrorlangan"
            assert (df["summa"] > 0).all(), "manfiy summa"
            return df

        @q.bosqich("validate", "valyuta_kursi")
        def transform(df, kurs):
            return df.assign(summa_usd=(df["summa"] / kurs["USD"]).round(2))

        @q.bosqich("transform")
        def load(df):
            with closing(sqlite3.connect(papka / "ombor.db")) as db:
                df.to_sql("buyurtma", db, if_exists="replace", index=False)
                return db.execute("SELECT COUNT(*) FROM buyurtma").fetchone()[0]

        @q.bosqich("transform")
        def hisobot(df):
            return df.groupby("hudud")["summa_usd"].sum().round(2).to_dict()

        print("=== 1. DAG va topologik tartib ===")
        for n, (_, b) in q.bosqichlar.items():
            print(f"  {n:<14} <- {b if b else '-'}")
        print(f"  tartib: {' -> '.join(q.tartib())}")

        print("\n=== 2. Ishga tushirish (qayta urinish bilan) ===")
        natija, jurnal = q.ishga_tushir()
        for qator in qayta_jurnal:
            print(f"  {qator}")
        print(f"  extract chaqiruvlari: {urinishlar['n']}, "
              f"virtual kutish: {soat.t:.0f} s")
        print(f"  bajarildi: {jurnal}")
        print(f"  load: {natija['load']} qator; hisobot: {natija['hisobot']}")

        print("\n=== 3. Faqat o'zgargan tarmoqni qayta hisoblash ===")
        pastki = q.pastki("valyuta_kursi")
        print(f"  valyuta_kursi o'zgardi -> qayta: {pastki}")
        print(f"  qayta hisoblanmaydi: {[n for n in q.tartib() if n not in pastki]}")

        print("\n=== 4. Noto'g'ri DAG ushlanadi ===")
        q2 = Quvur()
        q2.bosqichlar = {"a": (None, ["c"]), "b": (None, ["a"]), "c": (None, ["b"])}
        try:
            q2.tartib()
        except ValueError as x:
            print(f"  ValueError: {x}")
        q2.bosqichlar = {"a": (None, ["yoq"])}
        try:
            q2.tartib()
        except KeyError as x:
            print(f"  KeyError: {x}")

        print("\n=== 5. Qayta urinish ham tugaydi - xato yashirilmaydi ===")
        soat2 = VirtualSoat()
        j2 = []

        def doim_xato():
            raise ConnectionError("server javob bermadi")

        try:
            qayta_urin(doim_xato, soat2, urinish=4, jurnal=j2)
        except ConnectionError as x:
            print(f"  {len(j2)} kutish ({soat2.t:.0f} s), 4-urinishdan keyin: "
                  f"ConnectionError: {x}")


if __name__ == "__main__":
    main()

Natijaning muhim qismi:

text
=== 1. DAG va topologik tartib ===
  extract        <- -
  valyuta_kursi  <- -
  validate       <- ['extract']
  transform      <- ['validate', 'valyuta_kursi']
  load           <- ['transform']
  hisobot        <- ['transform']
  tartib: extract -> validate -> valyuta_kursi -> transform -> hisobot -> load

=== 2. Ishga tushirish (qayta urinish bilan) ===
  urinish 1: timeout; 1 s kutish
  urinish 2: timeout; 2 s kutish
  extract chaqiruvlari: 3, virtual kutish: 3 s
  bajarildi: ['extract', 'validate', 'valyuta_kursi', 'transform', 'hisobot', 'load']
  load: 200 qator; hisobot: {'A': 434.3, 'B': 649.81, 'C': 425.21}

=== 3. Faqat o'zgargan tarmoqni qayta hisoblash ===
  valyuta_kursi o'zgardi -> qayta: ['valyuta_kursi', 'transform', 'hisobot', 'load']
  qayta hisoblanmaydi: ['extract', 'validate']

=== 4. Noto'g'ri DAG ushlanadi ===
  ValueError: sikl bor: ['a', 'b', 'c']
  KeyError: "a: noma'lum bog'liqlik 'yoq'"

=== 5. Qayta urinish ham tugaydi - xato yashirilmaydi ===
  3 kutish (7 s), 4-urinishdan keyin: ConnectionError: server javob bermadi

Nima ko'rsatdi: 2.2, 2.8-bo'limlar. Bosqichlar kodda qaysi tartibda yozilganidan qat'i nazar, freymvork ularni bog'liqlik bo'yicha tartibladi; bir vaqtda tayyor bo'lganlar alifbo tartibida — shuning uchun tartib har ishga tushishda bir xil (27.2 dagi sorted qoidasi). Beqaror manba ikki marta yiqildi, quvur 1 va 2 soniya (virtual) kutib, uchinchi urinishda davom etdi — foydalanuvchi uchun bu shaffof, lekin jurnalda iz qoldi. Valyuta kursi o'zgarganda faqat to'rtta pastki bosqich qayta hisoblanishi kerak, extract va validate esa yo'q — bu 4-misoldagi keshning asosi. Sikl va noma'lum bog'liqlik ishga tushishdan oldin ushlandi. Va nihoyat, qayta urinishning chegarasi bor: uch kutishdan keyin xato yashirilmasdan ko'tarildi — orkestrator uni ko'radi va alert beradi.

load bosqichidagi closing(...) ga e'tibor bering: with sqlite3.connect(...) as db ulanishni yopmaydi, faqat tranzaksiyani yakunlaydi. Windows da ochiq qolgan ulanish vaqtinchalik papkani o'chirishga ham xalaqit beradi.

Misol 2 — Idempotent yuklash, tranzaksiya va watermark

Manba tizimdagi buyurtmalar (id, summa, holat, yangilangan — soat). Avval sodda INSERT ni uch marta ishga tushiramiz, keyin UPSERT ni. So'ng yuklash o'rtasida "tarmoq uzilishi" — tranzaksiya bilan va tranzaksiyasiz. Oxirida inkremental yuklash: watermark, yangilangan eski qator va kechikib kelgan yozuvlar.

python
"""Idempotent yuklash, tranzaksiya va watermark bilan inkremental yuklash."""

import hashlib
import sqlite3
import tempfile
from contextlib import closing
from pathlib import Path

import numpy as np

UPSERT = ("INSERT INTO ombor VALUES (?, ?, ?, ?) ON CONFLICT(id) DO UPDATE SET "
          "summa=excluded.summa, holat=excluded.holat, "
          "yangilangan=excluded.yangilangan")


def manba_yarat(n, rng, boshlanish=1, vaqt0=0):
    """Manba tizimdagi buyurtmalar: (id, summa, holat, yangilangan soat)."""
    return [(int(boshlanish + i), float(rng.gamma(2.0, 50.0) // 1),
             str(rng.choice(["yangi", "tolandi"])), int(vaqt0 + i // 10))
            for i in range(n)]


def jadval_xesh(db, jadval):
    qatorlar = db.execute(f"SELECT * FROM {jadval} ORDER BY id").fetchall()
    return hashlib.sha256(repr(qatorlar).encode()).hexdigest()[:12]


def yarat(db):
    db.execute("DROP TABLE IF EXISTS ombor")
    db.execute("CREATE TABLE ombor (id INTEGER PRIMARY KEY, summa REAL, "
               "holat TEXT, yangilangan INTEGER)")


def yukla_sodda(db, qatorlar):
    db.executemany("INSERT INTO sodda VALUES (?, ?, ?, ?)", qatorlar)
    db.commit()


def yukla_upsert(db, qatorlar):
    with db:                                  # tranzaksiya: hammasi yoki hech narsa
        db.executemany(UPSERT, qatorlar)


def yukla_uzilish_bilan(db, qatorlar, tranzaksiya):
    """Yuklash o'rtasida 'tarmoq uzildi'."""
    def yoz():
        for i, q in enumerate(qatorlar):
            if i == len(qatorlar) // 2:
                raise ConnectionError("uzildi")
            db.execute(UPSERT, q)
            if not tranzaksiya:
                db.commit()                   # har qatordan keyin - tipik xato
    try:
        if tranzaksiya:
            with db:
                yoz()
        else:
            yoz()
    except ConnectionError:
        pass


def watermark_ol(db):
    r = db.execute("SELECT qiymat FROM holat WHERE kalit='watermark'").fetchone()
    return r[0] if r else -1


def inkremental(db, manba, orqaga=0):
    wm = watermark_ol(db)
    yangi = [q for q in manba if q[3] > wm - orqaga]
    with db:                                  # yuklash va watermark - BITTA tranzaksiya
        db.executemany(UPSERT, yangi)
        if yangi:
            yangi_wm = max(wm, max(q[3] for q in yangi))
            db.execute("INSERT INTO holat VALUES ('watermark', ?) ON CONFLICT(kalit) "
                       "DO UPDATE SET qiymat=excluded.qiymat", (yangi_wm,))
    return len(yangi)


def main() -> None:
    manba = manba_yarat(500, np.random.default_rng(0))
    with tempfile.TemporaryDirectory() as t, \
            closing(sqlite3.connect(Path(t) / "ombor.db")) as db:
        print("=== 1. Sodda INSERT: qayta ishga tushirish dublikat beradi ===")
        db.execute("CREATE TABLE sodda (id INTEGER, summa REAL, holat TEXT, "
                   "yangilangan INTEGER)")
        for yurish in range(1, 4):
            yukla_sodda(db, manba)
            n = db.execute("SELECT COUNT(*) FROM sodda").fetchone()[0]
            s = db.execute("SELECT SUM(summa) FROM sodda").fetchone()[0]
            print(f"  {yurish}-yurish: {n} qator, summa jami {s:,.0f}")

        print("\n=== 2. Idempotent UPSERT (PRIMARY KEY) ===")
        yarat(db)
        xeshlar = []
        for yurish in range(1, 4):
            yukla_upsert(db, manba)
            n = db.execute("SELECT COUNT(*) FROM ombor").fetchone()[0]
            xeshlar.append(jadval_xesh(db, "ombor"))
            print(f"  {yurish}-yurish: {n} qator, jadval xeshi {xeshlar[-1]}")
        print(f"  uchala xesh bir xil: {len(set(xeshlar)) == 1}")

        print("\n=== 3. Yuklash o'rtasida uzilish ===")
        partiya = manba_yarat(100, np.random.default_rng(1), boshlanish=1001)
        for tranz in [False, True]:
            yarat(db)
            yukla_uzilish_bilan(db, partiya, tranzaksiya=tranz)
            n = db.execute("SELECT COUNT(*) FROM ombor").fetchone()[0]
            nom = "tranzaksiya bilan" if tranz else "tranzaksiyasiz"
            holat = "hammasi yoki hech narsa" if n in (0, 100) else "YARIM holat"
            print(f"  {nom:<17}: {n}/100 qator yozildi ({holat})")
        yukla_upsert(db, partiya)
        n = db.execute("SELECT COUNT(*) FROM ombor").fetchone()[0]
        print(f"  tranzaksiyadan keyin qayta yurish: {n}/100")

        print("\n=== 4. Inkremental yuklash: watermark ===")
        yarat(db)
        db.execute("CREATE TABLE holat (kalit TEXT PRIMARY KEY, qiymat INTEGER)")
        manba2 = manba_yarat(300, np.random.default_rng(2))       # soat 0..29
        print(f"  1-yurish: {inkremental(db, manba2)} qator, "
              f"watermark {watermark_ol(db)}")
        manba2 += manba_yarat(50, np.random.default_rng(3), boshlanish=301,
                              vaqt0=30)                           # yangi: soat 30..34
        manba2[10] = (manba2[10][0], 999.0, "tolandi", 33)        # eski qator yangilandi
        print(f"  2-yurish: {inkremental(db, manba2)} qator "
              f"(50 yangi + 1 yangilangan), watermark {watermark_ol(db)}")
        print(f"  3-yurish (o'zgarishsiz): {inkremental(db, manba2)} qator")
        s11 = db.execute("SELECT summa FROM ombor WHERE id=11").fetchone()[0]
        print(f"  id=11 summa omborda: {s11}")

        print("\n=== 5. Kechikib kelgan yozuvlar ===")
        manba2 += manba_yarat(5, np.random.default_rng(4), boshlanish=401,
                              vaqt0=32)                           # soat 32 < watermark 34
        manba2 += manba_yarat(10, np.random.default_rng(5), boshlanish=501,
                              vaqt0=35)
        n_oddiy = inkremental(db, manba2)
        kech = db.execute("SELECT COUNT(*) FROM ombor WHERE id BETWEEN 401 AND 405"
                          ).fetchone()[0]
        print(f"  oddiy watermark: {n_oddiy} qator olindi, kechikkan 5 tadan "
              f"omborda: {kech}")
        n_oyna = inkremental(db, manba2, orqaga=5)
        kech = db.execute("SELECT COUNT(*) FROM ombor WHERE id BETWEEN 401 AND 405"
                          ).fetchone()[0]
        print(f"  orqaga 5 soatlik oyna: {n_oyna} qator qayta o'qildi, "
              f"kechikkan: {kech}/5")
        n = db.execute("SELECT COUNT(*) FROM ombor").fetchone()[0]
        noyob = len({q[0] for q in manba2})
        print(f"  ombor: {n} qator, manba: {noyob} noyob id, "
              f"dublikat yo'q: {n == noyob}")


if __name__ == "__main__":
    main()

Natijaning muhim qismi:

text
=== 1. Sodda INSERT: qayta ishga tushirish dublikat beradi ===
  1-yurish: 500 qator, summa jami 47,776
  2-yurish: 1000 qator, summa jami 95,552
  3-yurish: 1500 qator, summa jami 143,328

=== 2. Idempotent UPSERT (PRIMARY KEY) ===
  1-yurish: 500 qator, jadval xeshi 01d9f88ebd5c
  2-yurish: 500 qator, jadval xeshi 01d9f88ebd5c
  3-yurish: 500 qator, jadval xeshi 01d9f88ebd5c
  uchala xesh bir xil: True

=== 3. Yuklash o'rtasida uzilish ===
  tranzaksiyasiz   : 50/100 qator yozildi (YARIM holat)
  tranzaksiya bilan: 0/100 qator yozildi (hammasi yoki hech narsa)
  tranzaksiyadan keyin qayta yurish: 100/100

=== 4. Inkremental yuklash: watermark ===
  1-yurish: 300 qator, watermark 29
  2-yurish: 51 qator (50 yangi + 1 yangilangan), watermark 34
  3-yurish (o'zgarishsiz): 0 qator
  id=11 summa omborda: 999.0

=== 5. Kechikib kelgan yozuvlar ===
  oddiy watermark: 10 qator olindi, kechikkan 5 tadan omborda: 0
  orqaga 5 soatlik oyna: 56 qator qayta o'qildi, kechikkan: 5/5
  ombor: 365 qator, manba: 365 noyob id, dublikat yo'q: True

Nima ko'rsatdi: 2.3, 2.4-bo'limlar. Sodda INSERT har yurishda jadvalni o'stirdi — uchinchi yurishdan keyin "jami summa" uch baravar. Bu kirishdagi real vaziyatning o'zi: raqamlar mantiqan to'g'ri ko'rinadi, faqat noto'g'ri. UPSERT bilan esa uch yurishning jadval xeshi aynan bir xil — idempotentlik o'lchandi, taxmin qilinmadi.

Uzilish tajribasi tranzaksiyaning qiymatini ko'rsatdi: har qatordan keyin commit qiladigan kod aynan yarmida to'xtab, jadvalni yarim holatda qoldirdi — qaysi qatorlar yozilganini bilish uchun qo'shimcha ish kerak. Tranzaksiya bilan jadval bo'sh qoldi, keyingi (idempotent) yurish esa hammasini toza yozdi.

Inkremental qismda watermark faqat yangi va yangilangan qatorlarni oldi (eski qatorning yangi summasi omborga tushdi), o'zgarishsiz yurish esa hech narsa o'qimadi. Lekin kechikib kelgan beshta yozuv oddiy watermark bilan jim yo'qoldi — hech qanday xato yo'q. Orqaga qarash oynasi ularni oldi va UPSERT tufayli qayta o'qilgan o'nlab qator dublikat bermadi: ombordagi qatorlar soni manbadagi noyob id lar soniga teng.

Misol 3 — Ma'lumot sifati tekshiruvlari noldan va karantin

great_expectations g'oyasida kichik "kutilmalar" to'plami: sxema, null ulushi, diapazon, noyoblik, ruxsat etilgan to'plam va kunlik o'rtacha chegarasi. Statistik chegaralar 30 kunlik toza tarixdan olinadi, keyin 30 ta yangi normal partiyada yolg'on signal soni tekshiriladi. So'ng yetti xil buzilgan partiya keladi: xato darajasidagilar karantinga, ogohlantirishlilar — omborga signal bilan.

python
"""Ma'lumot sifati tekshiruvlari noldan va karantin."""

import json
import sqlite3
import tempfile
from contextlib import closing
from dataclasses import asdict, dataclass
from pathlib import Path

import numpy as np
import pandas as pd


@dataclass
class Natija:
    nom: str
    otdi: bool
    daraja: str          # "xato" - karantin, "ogoh" - yuklanadi, signal beradi
    tafsilot: str


# --------------------- kutilmalar (expectations) ---------------------
def sxema(df, kutilgan):
    ustunlar = list(df.columns)
    if ustunlar != list(kutilgan):
        return Natija("sxema", False, "xato", f"ustunlar {ustunlar}")
    notogri = [c for c, tur in kutilgan.items() if not tur(df[c])]
    return Natija("sxema", not notogri, "xato",
                  f"tur mos emas: {notogri}" if notogri else "ok")


def null_ulushi(df, ustun, maks):
    u = float(df[ustun].isna().mean())
    return Natija(f"null({ustun})", u <= maks, "xato", f"{u:.1%} (maks {maks:.0%})")


def diapazon(df, ustun, past, yuqori):
    s = df[ustun].dropna()
    tashqi = int(((s < past) | (s > yuqori)).sum())
    return Natija(f"diapazon({ustun})", tashqi == 0, "xato",
                  f"{tashqi} ta [{past}, {yuqori}] dan tashqari")


def noyob(df, ustun):
    d = int(df[ustun].duplicated().sum())
    return Natija(f"noyob({ustun})", d == 0, "xato", f"{d} ta dublikat")


def toplamda(df, ustun, ruxsat):
    yangi = sorted(set(df[ustun].dropna()) - set(ruxsat))
    return Natija(f"toplam({ustun})", not yangi, "ogoh",
                  f"yangi: {yangi}" if yangi else "ok")


def ortacha_chegara(df, ustun, past, yuqori):
    m = float(df[ustun].mean())
    return Natija(f"ortacha({ustun})", past <= m <= yuqori, "ogoh",
                  f"{m:,.0f} (kutilgan [{past:,.0f}, {yuqori:,.0f}])")


def qoidalar_yasa(mos):
    """Chegaralar TARIXIY TOZA ma'lumotdan (mos yozuv) olinadi."""
    kunlik = [g["summa"].mean() for _, g in mos.groupby("kun")]
    m, s = float(np.mean(kunlik)), float(np.std(kunlik, ddof=1))
    son = pd.api.types.is_numeric_dtype
    satr = pd.api.types.is_string_dtype
    return [
        lambda d: sxema(d, {"id": son, "kun": son, "summa": son, "hudud": satr}),
        lambda d: null_ulushi(d, "summa", 0.02),
        lambda d: diapazon(d, "summa", 0, 5_000_000),
        lambda d: noyob(d, "id"),
        lambda d: toplamda(d, "hudud", sorted(mos["hudud"].unique())),
        lambda d: ortacha_chegara(d, "summa", m - 4 * s, m + 4 * s),
    ]


def tekshir(df, qoidalar):
    birinchi = qoidalar[0](df)          # sxema buzilsa - qolganini tekshirib bo'lmaydi
    if not birinchi.otdi:
        return [birinchi]
    return [birinchi] + [q(df) for q in qoidalar[1:]]


def partiya(kun, rng, n=400):
    return pd.DataFrame({
        "id": np.arange(kun * 10_000, kun * 10_000 + n),
        "kun": kun,
        "summa": rng.gamma(2.0, 60_000, n).round(-2),
        "hudud": rng.choice(["Andijon", "Buxoro", "Toshkent"], n, p=[0.3, 0.2, 0.5])})


def qayta_ishla(nom, df, qoidalar, db, karantin):
    natijalar = tekshir(df, qoidalar)
    xato = [r for r in natijalar if not r.otdi and r.daraja == "xato"]
    ogoh = [r for r in natijalar if not r.otdi and r.daraja == "ogoh"]
    if xato:
        df.to_csv(karantin / f"{nom}.csv", index=False)           # xom nusxa
        (karantin / f"{nom}.json").write_text(json.dumps(
            [asdict(r) for r in natijalar], ensure_ascii=False, indent=1),
            encoding="utf-8")                                    # sabab
        holat = "KARANTIN"
    else:
        with db:
            db.executemany("INSERT INTO buyurtma VALUES (?, ?, ?, ?)",
                           df.itertuples(index=False, name=None))
        holat = "yuklandi" + (" + OGOH" if ogoh else "")
    sabab = "; ".join(f"{r.nom}: {r.tafsilot}" for r in xato + ogoh)
    print(f"  {nom:<17} {holat:<15} {sabab or '-'}")


def main() -> None:
    rng = np.random.default_rng(0)
    mos = pd.concat([partiya(k, rng) for k in range(1, 31)], ignore_index=True)
    qoidalar = qoidalar_yasa(mos)

    print("=== 1. 30 ta yangi normal partiyada yolg'on signal ===")
    yolgon = {}
    for k in range(31, 61):
        for r in tekshir(partiya(k, rng), qoidalar):
            if not r.otdi:
                yolgon[r.nom] = yolgon.get(r.nom, 0) + 1
    print(f"  muvaffaqiyatsiz tekshiruvlar: {yolgon or 'yoq'}")

    print("\n=== 2. Normal va buzilgan partiyalar ===")
    partiyalar = [("normal_61", partiya(61, rng))]
    b = partiya(62, rng)
    b.loc[b.sample(frac=0.2, random_state=0).index, "summa"] = np.nan
    partiyalar.append(("null_portlash", b))
    b = partiya(63, rng)
    b.loc[:40, "id"] = b.loc[0, "id"]
    partiyalar.append(("dublikat_id", b))
    b = partiya(64, rng)
    b.loc[:9, "summa"] = -b.loc[:9, "summa"]
    partiyalar.append(("manfiy_summa", b))
    b = partiya(65, rng)
    b.loc[:59, "hudud"] = "Samarqand"
    partiyalar.append(("yangi_hudud", b))
    b = partiya(66, rng)
    b["summa"] = (b["summa"] * 1.6).round(-2)
    partiyalar.append(("siljigan_ortacha", b))
    b = partiya(67, rng)
    b["summa"] = b["summa"].astype(int).astype(str) + " so'm"
    partiyalar.append(("matn_summa", b))
    partiyalar.append(("ustun_nomi", partiya(68, rng).rename(
        columns={"summa": "miqdor"})))

    with tempfile.TemporaryDirectory() as t:
        karantin = Path(t) / "karantin"
        karantin.mkdir()
        with closing(sqlite3.connect(Path(t) / "ombor.db")) as db:
            db.execute("CREATE TABLE buyurtma (id INTEGER PRIMARY KEY, kun INTEGER, "
                       "summa REAL, hudud TEXT)")
            for nom, df in partiyalar:
                qayta_ishla(nom, df, qoidalar, db, karantin)
            n = db.execute("SELECT COUNT(*) FROM buyurtma").fetchone()[0]
            print(f"\n  omborda: {n} qator")
            print(f"  karantinda: {sorted(p.stem for p in karantin.glob('*.csv'))}")

            print("\n=== 3. Karantindan chiqish: tuzatish va QAYTA tekshirish ===")
            sabab = json.loads((karantin / "matn_summa.json").read_text(encoding="utf-8"))
            print(f"  saqlangan sabab: {sabab[0]['nom']} - {sabab[0]['tafsilot']}")
            df = pd.read_csv(karantin / "matn_summa.csv")
            df["summa"] = pd.to_numeric(df["summa"].str.replace(" so'm", "",
                                                                regex=False))
            qayta_ishla("matn_summa_tuz", df, qoidalar, db, karantin)
            n = db.execute("SELECT COUNT(*) FROM buyurtma").fetchone()[0]
            print(f"  omborda endi: {n} qator")


if __name__ == "__main__":
    main()

Natijaning muhim qismi:

text
=== 1. 30 ta yangi normal partiyada yolg'on signal ===
  muvaffaqiyatsiz tekshiruvlar: yoq

=== 2. Normal va buzilgan partiyalar ===
  normal_61         yuklandi        -
  null_portlash     KARANTIN        null(summa): 20.0% (maks 2%)
  dublikat_id       KARANTIN        noyob(id): 40 ta dublikat
  manfiy_summa      KARANTIN        diapazon(summa): 10 ta [0, 5000000] dan tashqari
  yangi_hudud       yuklandi + OGOH toplam(hudud): yangi: ['Samarqand']
  siljigan_ortacha  yuklandi + OGOH ortacha(summa): 198,321 (kutilgan [104,413, 135,966])
  matn_summa        KARANTIN        sxema: tur mos emas: ['summa']
  ustun_nomi        KARANTIN        sxema: ustunlar ['id', 'kun', 'miqdor', 'hudud']

  omborda: 1200 qator
  karantinda: ['dublikat_id', 'manfiy_summa', 'matn_summa', 'null_portlash', 'ustun_nomi']

=== 3. Karantindan chiqish: tuzatish va QAYTA tekshirish ===
  saqlangan sabab: sxema - tur mos emas: ['summa']
  matn_summa_tuz    yuklandi        -
  omborda endi: 1600 qator

Nima ko'rsatdi: 2.5, 2.6-bo'limlar. Avval chegaralar sinovdan o'tdi: 30 ta yangi normal partiyada birorta ham tekshiruv yiqilmadi — ya'ni signal kelsa, unga ishonish mumkin. Keyin har buzilish o'z tekshiruvi bilan ushlandi va sababi aniq aytildi: 20% bo'sh summa, 40 ta dublikat id, 10 ta manfiy summa, matnga aylangan summa, qayta nomlangan ustun. Bu beshtasi karantinga tushdi va omborga birorta qatori ham kirmadi.

Ikki holat ongli ravishda ogohlantirish darajasida: yangi hudud ("Samarqand") va o'rtachaning 1.6 baravar oshishi. Ular xato bo'lishi ham, haqiqiy biznes o'zgarishi (yangi filial, aksiya) bo'lishi ham mumkin — shuning uchun ma'lumot yuklanadi, lekin inson xabardor qilinadi. Ikkinchisi 27.1-misoldagi "so'm → ming so'm" ga o'xshash muammoni ham ushlaydi.

Oxirgi qism karantinning to'liq siklini ko'rsatdi: saqlangan sababdan nima buzilgani ma'lum, xom nusxa tuzatildi va xuddi shu tekshiruvlardan qayta o'tkazildi — shundan keyingina omborga yuklandi. Qo'lda "ko'chirib qo'yish" yo'q.

Misol 4 — Kirish xeshi bo'yicha keshlash va quvurdagi vaqt sizishi

Birinchi qism: har bosqich natijasi kalit bo'yicha keshlanadi — kalit = xesh(bosqich nomi + kod versiyasi + kirish xeshlari + parametrlar). Qaysi o'zgarish qaysi bosqichni qayta hisoblashini sanaymiz (vaqt emas, chaqiruvlar soni). Ikkinchi qism — qisqa: ikki yillik kunlik talab qatori; tasodifiy va vaqt bo'yicha bo'lish, hamda "kelajakni ko'radigan" belgi.

python
"""Kirish xeshi bo'yicha keshlash va quvurdagi vaqt sizishi."""

import hashlib
import json
import pickle
import tempfile
from pathlib import Path

import numpy as np
import pandas as pd
from sklearn.ensemble import HistGradientBoostingRegressor


def xesh_obyekt(obj):
    if isinstance(obj, pd.DataFrame):
        b = pd.util.hash_pandas_object(obj, index=True).to_numpy().tobytes()
        b += json.dumps(list(map(str, obj.columns))).encode()
    else:
        b = json.dumps(obj, sort_keys=True, default=str).encode()
    return hashlib.sha256(b).hexdigest()[:12]


class Kesh:
    def __init__(self, papka):
        self.papka = Path(papka)
        self.hisob = {"hisoblandi": [], "keshdan": []}

    def bosqich(self, versiya):
        """versiya - bosqich kodi o'zgarganda qo'lda oshiriladi."""
        def bezak(fn):
            def ichki(*kirishlar, **param):
                kalit = xesh_obyekt({"fn": fn.__name__, "v": versiya,
                                     "kirish": [xesh_obyekt(k) for k in kirishlar],
                                     "param": param})
                yol = self.papka / f"{fn.__name__}_{kalit}.pkl"
                if yol.exists():                   # o'z keshimiz - ishonchli fayl
                    self.hisob["keshdan"].append(fn.__name__)
                    return pickle.loads(yol.read_bytes())
                natija = fn(*kirishlar, **param)
                yol.write_bytes(pickle.dumps(natija))
                self.hisob["hisoblandi"].append(fn.__name__)
                return natija
            return ichki
        return bezak


def xom_malumot(kunlar=730, seed=0):
    rng = np.random.default_rng(seed)
    t = np.arange(kunlar)
    talab = (200 + 0.15 * t + 40 * np.sin(2 * np.pi * t / 7)
             + 60 * np.sin(2 * np.pi * t / 365) + rng.normal(0, 15, kunlar))
    harorat = 15 + 12 * np.sin(2 * np.pi * (t - 100) / 365) + rng.normal(0, 3, kunlar)
    return pd.DataFrame({"kun": t, "harorat": harorat.round(1),
                         "talab": talab.round(1)})


def main() -> None:
    with tempfile.TemporaryDirectory() as t:
        k = Kesh(t)

        @k.bosqich(versiya=1)
        def tozalash(df):
            return df.dropna().reset_index(drop=True)

        @k.bosqich(versiya=1)
        def belgilar(df, oyna=7):
            d = df.copy()
            d["hafta_kuni"] = d["kun"] % 7
            d["otgan_ortacha"] = d["talab"].shift(1).rolling(oyna).mean()
            return d.dropna().reset_index(drop=True)

        @k.bosqich(versiya=1)
        def orgat(df, max_iter=100):
            X, y = df[["kun", "harorat", "hafta_kuni", "otgan_ortacha"]], df["talab"]
            m = HistGradientBoostingRegressor(max_iter=max_iter, random_state=0)
            return m.fit(X.iloc[:-90], y.iloc[:-90])

        def quvur(df, oyna=7, max_iter=100):
            k.hisob = {"hisoblandi": [], "keshdan": []}
            orgat(belgilar(tozalash(df), oyna=oyna), max_iter=max_iter)
            return dict(k.hisob)

        print("=== 1. Kesh: kirish xeshi + parametr + kod versiyasi ===")
        df = xom_malumot()
        for izoh, kw in [("birinchi yurish", {}), ("aynan takror", {}),
                         ("max_iter=200", {"max_iter": 200}),
                         ("oyna=14", {"oyna": 14})]:
            h = quvur(df, **kw)
            print(f"  {izoh:<17} hisoblandi {h['hisoblandi']}, keshdan {h['keshdan']}")
        df2 = df.copy()
        df2.loc[5, "talab"] += 1.0
        h = quvur(df2)
        print(f"  {'1 qiymat ozgardi':<17} hisoblandi {h['hisoblandi']}, "
              f"keshdan {h['keshdan']}")
        print(f"  kesh fayllari: {len(list(Path(t).glob('*.pkl')))} ta")

    print("\n=== 2. Quvurdagi sizish: bo'lish usuli va kelajakni ko'radigan belgi ===")
    d = xom_malumot()
    d["hafta_kuni"] = d["kun"] % 7
    d["otgan"] = d["talab"].shift(1).rolling(7).mean()               # faqat o'tmish
    d["markaz"] = d["talab"].rolling(7, center=True).mean()          # 3 kun kelajak!
    d["markaz_xizmatda"] = d["talab"].shift(1).rolling(4).mean()     # xizmatda mavjudi
    d = d.dropna().reset_index(drop=True)
    kelajak = np.asarray(d.index >= len(d) - 90)                 # oxirgi 90 kun
    tarix = ~kelajak
    tasodif_val = tarix & (np.random.default_rng(0).random(len(d)) < 0.2)
    vaqt_val = tarix & np.asarray(d.index >= len(d) - 180)
    asos = ["kun", "harorat", "hafta_kuni"]
    print(f"  {'belgilar':<18} {'bolish':<10} {'val MAE':>8} {'kelajak MAE':>12}")
    xatolar = {}
    for nom, ustunlar in [("asos", asos), ("asos + otgan", asos + ["otgan"]),
                          ("asos + markaz", asos + ["markaz"])]:
        for bol, val in [("tasodifiy", tasodif_val), ("vaqt", vaqt_val)]:
            tr = tarix & ~val
            m = HistGradientBoostingRegressor(max_iter=100, random_state=0).fit(
                d.loc[tr, ustunlar], d.loc[tr, "talab"])
            v = np.abs(d.loc[val, "talab"] - m.predict(d.loc[val, ustunlar])).mean()
            Xk = d.loc[kelajak, ustunlar].copy()
            if "markaz" in ustunlar:                                 # xizmatda kelajak yo'q
                Xk["markaz"] = d.loc[kelajak, "markaz_xizmatda"]
            xato = np.abs(d.loc[kelajak, "talab"] - m.predict(Xk)).to_numpy()
            xatolar[(nom, bol)] = xato
            print(f"  {nom:<18} {bol:<10} {v:>8.2f} {xato.mean():>12.2f}")
    dd = xatolar[("asos + markaz", "vaqt")] - xatolar[("asos + otgan", "vaqt")]
    se = dd.std(ddof=1) / np.sqrt(len(dd))
    baho = "SEZILARLI yomon" if dd.mean() > 2 * se else "sezilarli emas"
    print(f"  kelajakda markaz - otgan (90 kun, juftlashgan): {dd.mean():+.2f} "
          f"(SE {se:.2f}) - {baho}")
    print("  shovqin std = 15 -> eng yaxshi mumkin MAE taxminan 12")


if __name__ == "__main__":
    main()

Natijaning muhim qismi:

text
=== 1. Kesh: kirish xeshi + parametr + kod versiyasi ===
  birinchi yurish   hisoblandi ['tozalash', 'belgilar', 'orgat'], keshdan []
  aynan takror      hisoblandi [], keshdan ['tozalash', 'belgilar', 'orgat']
  max_iter=200      hisoblandi ['orgat'], keshdan ['tozalash', 'belgilar']
  oyna=14           hisoblandi ['belgilar', 'orgat'], keshdan ['tozalash']
  1 qiymat ozgardi  hisoblandi ['tozalash', 'belgilar', 'orgat'], keshdan []
  kesh fayllari: 9 ta

=== 2. Quvurdagi sizish: bo'lish usuli va kelajakni ko'radigan belgi ===
  belgilar           bolish      val MAE  kelajak MAE
  asos               tasodifiy     13.60        26.78
  asos               vaqt          45.13        39.33
  asos + otgan       tasodifiy     14.90        13.89
  asos + otgan       vaqt          15.56        16.58
  asos + markaz      tasodifiy     14.21        19.91
  asos + markaz      vaqt          10.24        19.88
  kelajakda markaz - otgan (90 kun, juftlashgan): +3.30 (SE 1.61) - SEZILARLI yomon
  shovqin std = 15 -> eng yaxshi mumkin MAE taxminan 12

Nima ko'rsatdi: 2.7, 2.10-bo'limlar. Kesh aynan kerakli joyda ishladi: takroriy yurishda uchala bosqich keshdan olindi; max_iter faqat orgat ga tegdi; oyna belgilar va o'qitishni qayta hisoblatdi, tozalashni emas; bitta qiymatning o'zgarishi esa kirish xeshini o'zgartirdi va hammasini qayta hisoblatdi — eskirgan natija qaytmadi. Kalitda kod versiyasi borligi uchun bosqich kodini tuzatganda versiya ni oshirish kifoya.

Sizish qismida ikki xil xato ko'rindi. Birinchisi — bo'lish. Faqat kalendar belgilari bilan (asos) tasodifiy bo'lish validatsiyada juda yaxshi MAE ni va'da qildi: model har validatsiya kunining ikki tomonidagi qo'shnilarini ko'rgan, ya'ni interpolyatsiya qilgan. Haqiqiy kelajakda xato taxminan ikki baravar katta. Vaqt bo'yicha validatsiya bu modelni "yomon" deb to'g'ri ko'rsatdi. Ikkinchisi — belgi. markaz (markazlashgan oyna) vaqt bo'yicha validatsiyada eng yaxshi bo'lib ko'rindi — chunki u uch kunlik kelajakni biladi. Xizmatda esa kelajak yo'q, belgini faqat o'tmishdan hisoblash mumkin — va kelajakdagi xato halol otgan belgisiga nisbatan sezilarli yomon (+3.30, 2*SE = 3.22 — chegaraga yaqin, lekin undan katta). Validatsiyada eng yaxshi ko'ringan belgi ishlab chiqarishda halol belgidan yomonroq bo'lib chiqdi. Batafsil — 18.2 va 18.11-darslarda.


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

Noto'g'ri fikr To'g'risi
"Qayta ishga tushirish zararsiz" 2-misolda sodda INSERT uch yurishda jami summani uch baravar qildi
"Har qatordan keyin commit — ishonchliroq" Uzilishda jadval yarim holatda qoldi; tranzaksiya — hammasi yoki hech narsa
"Watermark kechikkan ma'lumotni ham oladi" Oddiy watermark 5 ta kechikkan yozuvni jim yo'qotdi; orqaga oyna + UPSERT kerak
"Tekshiruv chegarasini taxminan qo'yish mumkin" Chegaralar tarixdan olinadi va normal partiyalarda yolg'on signal o'lchanadi
"Buzilgan partiyani o'chirib tashlash kerak" Karantinda xom nusxa va sabab saqlanadi; tuzatib qayta tekshiriladi
"Barcha tekshiruvlar bir xil darajada" Sxema/diapazon — xato, yangi kategoriya/o'rtacha — ogohlantirish
"Kesh kaliti = fayl nomi" Kalit = kod versiyasi + kirish xeshi + parametr
"Orkestrator bo'lsa, quvur ishonchli" Orkestrator faqat ishga tushiradi va qayta uradi; bosqich idempotent bo'lishi kerak
"Validatsiyada eng yaxshi belgi — eng yaxshisi" markaz belgisi xizmatda mavjud emas; kelajakda sezilarli yomon

6. Keng tarqalgan xatolar va yechimlari

1. Asosiy kalitsiz INSERT

python
db.executemany("INSERT INTO t VALUES (?, ?, ?)", qatorlar)                   # ⚠️ dublikat
db.executemany("INSERT INTO t VALUES (?, ?, ?) ON CONFLICT(id) DO UPDATE SET "
               "summa=excluded.summa, holat=excluded.holat", qatorlar)        # ✅ UPSERT

2. Tranzaksiyasiz yuklash

python
for q in qatorlar: db.execute(SQL, q); db.commit()           # ⚠️ yarim holat
with db: db.executemany(SQL, qatorlar)                       # ✅ hammasi yoki hech narsa

3. Watermark ni alohida yangilash

python
yukla(yangi); db.execute("UPDATE holat SET wm=?", (m,))      # ⚠️ ikki tranzaksiya
with db: db.executemany(UPSERT, yangi); db.execute(...)      # ✅ bitta tranzaksiya

4. Oynasiz inkremental yuklash

python
yangi = [q for q in manba if q.vaqt > wm]                    # ⚠️ kechikkanlar yo'qoladi
yangi = [q for q in manba if q.vaqt > wm - oyna]             # ✅ + UPSERT

5. Sxemasiz diapazon tekshiruvi

python
(df["summa"] < 0).sum()                                      # ⚠️ summa matn bo'lsa - TypeError
if sxema(df).otdi: diapazon(df, "summa", 0, maks)            # ✅ sxema birinchi

6. Ulanishni yopmaslik

python
with sqlite3.connect(yol) as db: ...                         # ⚠️ faqat commit, ulanish ochiq
with closing(sqlite3.connect(yol)) as db, db: ...            # ✅ yopiladi va tranzaksiya

7. Kod versiyasisiz kesh kaliti

python
kalit = xesh(kirish)                                         # ⚠️ tuzatilgan kod - eski natija
kalit = xesh({"fn": nom, "v": versiya, "kirish": ..., "param": ...})   # ✅

7. Integratsiya — bu bilim qayerda kerak bo'ladi

  • 6-qism (o'tilgan): ma'lumotni tozalash — bu yerda u quvur bosqichi va tekshiruv bo'ldi
  • 18.2, 18.11-darslar (o'tilgan): vaqt bo'yicha bo'lish va validatsiyaga moslashish — quvurdagi sizish
  • 25.6-dars (o'tilgan): retry/backoff va kesh — xuddi shu naqshlar
  • 27.1, 27.2-darslar (o'tilgan): kirish shartnomasi va ma'lumot xeshi — bu yerda ular sifat tekshiruvi va kesh kaliti
  • 27.4-dars: har eksperiment yurishi qaysi ma'lumot versiyasidan o'qiganini qayd qiladi
  • 27.11, 27.12-darslar: sifat tekshiruvlari ishlab chiqarish monitoringi va drift aniqlashga aylanadi

8. Eng yaxshi amaliyotlar

  1. Xom qatlamni o'zgartirmasdan saqlang; qolgan hamma narsa undan qayta hisoblansin.

  2. Bog'liqliklarni DAG sifatida ochiq yozing; tartibni deterministik qiling.

  3. Har yuklash — tabiiy kalit bo'yicha UPSERT yoki bo'limni almashtirish, har doim tranzaksiyada.

  4. Inkremental yuklashda watermark va ma'lumotni bitta tranzaksiyada yangilang, orqaga qarash oynasini qo'ying.

  5. Sifat tekshiruvlarini ketma-ket bajaring (sxema birinchi), chegaralarni tarixdan oling va yolg'on signalni o'lchang.

  6. Buzilgan partiyani karantinga — xom nusxa va sabab bilan; chiqish faqat qayta tekshiruv orqali.

  7. Kesh kalitiga kod versiyasi, kirish xeshi va parametrlarni kiriting.

  8. Faqat vaqtinchalik xatolarga, faqat idempotent bosqichlarga, cheklangan marta qayta urining.


9. Amaliy topshiriq

Vazifa 1: Bashorat qiling

python
1.  # ETL va ELT farqi?
2.  # xom qatlam nega o'zgartirilmaydi?
3.  # DAG da sikl bo'lsa nima bo'ladi?
4.  # Kahn algoritmida tartib qanday deterministik qilinadi?
5.  # idempotent yuklashning uch usuli?
6.  # tranzaksiyasiz yuklash uzilsa nima qoladi?
7.  # watermark va ma'lumot nega bitta tranzaksiyada yangilanadi?
8.  # kechikkan yozuvlarni qanday olish mumkin?
9.  # nega sxema tekshiruvi birinchi?
10. # qaysi tekshiruvlar odatda ogohlantirish darajasida?
11. # kesh kalitida nima bo'lishi kerak?
12. # qaysi xatolarga qayta urinish kerak emas?
Javoblar
  1. ETL — omborga yuklashdan oldin o'zgartirish; ELT — avval xom yuklash, keyin ombor ichida o'zgartirish
  2. Transformatsiyadagi xato topilsa, hammasini qayta hisoblash uchun
  3. Hech bir bosqich boshlanolmaydi — ishga tushishdan oldin xato
  4. "Tayyor" navbatni har doim sorted qilib
  5. UPSERT, bo'limni almashtirish, atomik fayl yozish (os.replace)
  6. Yarim holat — qaysi qatorlar yozilgani noma'lum
  7. Aks holda yiqilishda ular bir-biridan ajraladi (yo'qotish yoki dublikat)
  8. Orqaga qarash oynasi + UPSERT
  9. Turi noto'g'ri bo'lsa, boshqa tekshiruvlar o'zi xato beradi
  10. Yangi kategoriya, taqsimot (o'rtacha) siljishi
  11. Bosqich nomi, kod versiyasi, kirish xeshlari, parametrlar
  12. Doimiy xatolar: sxema xatosi, 401, 404

Vazifa 2: Xatolarni tuzating

python
1.  df.to_csv("ombor.csv", mode="a", header=False)          # har kecha

2.  for q in qatorlar:
        db.execute("INSERT INTO t VALUES (?, ?)", q)
        db.commit()

3.  yangi = [q for q in manba if q["vaqt"] > watermark]

4.  with sqlite3.connect("ombor.db") as db:
        db.execute("INSERT ...")
    # keyin faylni o'chirish: PermissionError

5.  @retry(attempts=10)
    def yukla(df): db.executemany("INSERT INTO t VALUES (?, ?)", df.values)
Javoblar
python
1.  # bo'limni almashtirish: kun bo'yicha alohida fayl, atomik yozish
    tmp = papka / f"kun={sana}.csv.tmp"; df.to_csv(tmp, index=False)
    os.replace(tmp, papka / f"kun={sana}.csv")

2.  with db:
        db.executemany("INSERT INTO t VALUES (?, ?) ON CONFLICT(id) DO UPDATE "
                       "SET qiymat=excluded.qiymat", qatorlar)

3.  yangi = [q for q in manba if q["vaqt"] > watermark - oyna]   # + UPSERT

4.  with closing(sqlite3.connect("ombor.db")) as db, db:
        db.execute("INSERT ...")

5.  # faqat vaqtinchalik xatolarga, kam marta, IDEMPOTENT yuklash bilan
    qayta_urin(lambda: upsert_yukla(db, df), kut, urinish=4)

Vazifa 3: Freymvorkni kengaytirish

1-misoldagi Quvur ga:

  1. Yiqilgan joydan davom ettirish: har bosqich natijasini diskka saqlang va davom=True bo'lsa, bajarilganlarini o'tkazib yuboring
  2. Bosqich darajasidagi urinish parametri (dekoratorda)
  3. Har bosqich uchun virtual davomiylik va "kritik yo'l" (eng uzun zanjir) hisobi
  4. DAG ni matn ko'rinishida chizish (extract -> validate -> ...)

Vazifa 4: Idempotentlik

2-misol asosida:

  1. "Bo'limni almashtirish" usulini yozing: BEGIN; DELETE WHERE kun=?; INSERT ...; COMMIT; va uch marta ishga tushirib xeshni tekshiring
  2. Atomik CSV yozish funksiyasini yozing va yozish o'rtasida istisno bo'lganda eski fayl buzilmaganini ko'rsating
  3. Manbadagi o'chirilgan qatorlarni "yumshoq o'chirish" (ochirilgan=1) ustuni bilan qo'llab-quvvatlang
  4. Kechikishlar taqsimotini simulyatsiya qilib (masalan, eksponensial), 99.9% yozuvlarni ushlaydigan oyna kattaligini tanlang

Vazifa 5: Sifat tekshiruvlari

3-misolga qo'shing:

  1. Hajm tekshiruvi: qatorlar soni tarixiy [q01, q99] oralig'ida
  2. Yangilik tekshiruvi: eng so'nggi kun kutilgan kundan eski emas
  3. Munosabat tekshiruvi: har hudud alohida hududlar jadvalida bor
  4. 100 ta normal partiyada har tekshiruvning yolg'on signal ulushini hisoblang va chegaralarni sozlang

Vazifa 6: Kesh va sizish

4-misolni o'zgartiring:

  1. belgilar bosqichining kodini o'zgartiring, lekin versiya ni oshirmang — nima bo'ladi? Keyin oshiring
  2. Kesh hajmini cheklang: eng eski fayllarni o'chiradigan (LRU) siyosat
  3. markaz o'rniga "o'quv to'plamining butun o'rtachasi bilan to'ldirish" sizishini qo'shing — ta'siri qanchalik katta?
  4. Vaqt bo'yicha bir necha "kesish nuqtasi" (rolling origin) bilan baholang va natijani o'rtacha ± std bilan yozing

Vazifa 7: O'ylash

Kompaniyaning tungi quvuri: to'lov tizimidan tranzaksiyalarni oladi, firibgarlik modeli uchun belgilar hisoblaydi va jadvalga yozadi. Oxirgi oyda uch hodisa bo'ldi: (1) bir kecha quvur ikki marta ishga tushdi va belgilar jadvalida dublikat paydo bo'ldi; (2) to'lov tizimi yangilanishdan keyin amount ni tiyinlarda yubora boshladi va bu uch kun sezilmadi; (3) ba'zi tranzaksiyalar 2-3 soat kechikib keladi va hech qachon belgilar jadvaliga tushmaydi. Har biri uchun nima qilasiz va qaysi tartibda?

Javob

Qisqa javob: uchala hodisa ham bu darsdagi bitta himoya yo'qligidan: (1) idempotentlik, (2) sifat tekshiruvi, (3) orqaga qarash oynasi. Tartib — zarar va tuzatish narxi bo'yicha: avval (2), chunki u jim va eng qimmat; keyin (1) va (3).

1. Birlik o'zgarishi (eng xavfli — jim)

python
# kirish shartnomasi 27.1-bob + taqsimot tekshiruvi (Misol 3)
ortacha_chegara(df, "amount", m - 4 * s, m + 4 * s)   # 100 baravar siljish -> signal
# amount uchun XATO darajasi: yuqori chegaradan tashqaridagi ulush > 5% -> karantin

Uch kun sezilmagani — tekshiruv yo'qligi. Median nisbati yoki o'rtacha chegarasi birinchi partiyadayoq ushlardi. Shuningdek, to'lov tizimi jamoasi bilan ma'lumot shartnomasi: birlik va tur o'zgarishi oldindan e'lon qilinadi.

2. Dublikat (idempotentlik)

python
# belgilar jadvali: PRIMARY KEY (tranzaksiya_id)
# yuklash: UPSERT yoki "kun bo'limini almashtirish" - tranzaksiyada
with db:
    db.execute("DELETE FROM belgilar WHERE kun = ?", (kun,))
    db.executemany("INSERT INTO belgilar VALUES (...)", qatorlar)

Quvur ikki marta ishga tushsa ham natija bir xil bo'ladi. Qo'shimcha ravishda orkestratorda bir vaqtda faqat bitta yurish (max_active_runs=1) — lekin bu asosiy himoya emas.

3. Kechikkan tranzaksiyalar

python
# watermark - "kelgan vaqt" bo'yicha emas, manba vaqti bo'yicha, 6 soatlik oyna bilan
yangi = manba.query("vaqt > @watermark - 6 * 3600")
# + UPSERT: qayta o'qilgan qatorlar dublikat bermaydi

Oyna kattaligi kechikishlar taqsimotidan tanlanadi (masalan, 99.9-kvantil 3 soat bo'lsa — 6 soat zaxira bilan). Juda kech kelganlar uchun haftalik "tuzatish" yurishi.

4. Model uchun oqibat

  • Uch kunlik noto'g'ri amount bilan hisoblangan belgilar tuzatilgan kod bilan qayta hisoblanadi — xom qatlam saqlangani uchun bu mumkin (2.1)
  • Shu kunlardagi model qarorlari tahlil qilinadi (qancha firibgarlik o'tkazib yuborilgan)
  • O'quv ma'lumotiga bu kunlar buzilgan holda kirmasligi uchun ma'lumot versiyasi (xesh) qayd qilinadi

Rahbarga javob: "Uchala hodisaning ham arzon texnik yechimi bor. Birinchi navbatda kiruvchi ma'lumotga sifat tekshiruvi qo'yamiz — u birlik o'zgarishini birinchi partiyadayoq ushlaydi. Keyin yuklashni idempotent qilamiz va kechikkan yozuvlar uchun oyna qo'shamiz. Buzilgan uch kunning belgilarini xom ma'lumotdan qayta hisoblaymiz."

Nimani mustahkamlaydi: 2.3, 2.4, 2.5-bo'limlar.


Xulosa

Bu darsda qayta ishga tushirilganda xavfsiz, buzilgan ma'lumotni o'tkazmaydigan ma'lumot quvurini noldan qurdik.

Eng muhim uch fikr:

  1. Quvur — DAG, bosqichlari idempotent va tranzaksiyali. 1-misolda kichik freymvork bog'liqliklardan deterministik tartib hisobladi, sikl va noma'lum bog'liqlikni ishga tushishdan oldin ushladi va beqaror manbaga eksponensial kutish bilan qayta urindi — cheklangan marta, xatoni yashirmasdan. 2-misolda sodda INSERT har qayta yurishda jadvalni o'stirdi, UPSERT esa uch yurishda aynan bir xil jadval xeshini berdi; tranzaksiyasiz yuklash uzilishda yarim holat qoldirdi, tranzaksiya bilan — hammasi yoki hech narsa.

  2. Inkremental yuklash va sifat tekshiruvi jim yo'qotish va jim buzilishga qarshi. Watermark faqat yangi va yangilangan qatorlarni oldi, lekin kechikkan yozuvlarni yo'qotdi — orqaga qarash oynasi va UPSERT ularni dublikatsiz qaytardi. 3-misolda tarixdan olingan chegaralar 30 normal partiyada yolg'on signal bermadi, besh xil buzilishni karantinga yubordi, ikki shubhali o'zgarishni ogohlantirish bilan o'tkazdi; karantindagi partiya faqat tuzatilib, qayta tekshirilgandan keyin omborga tushdi.

  3. Kesh kaliti — kod versiyasi, kirish va parametr; sizish — quvurning o'zida. 4-misolda kesh faqat o'zgarish tekkan bosqichlarni qayta hisobladi. Tasodifiy bo'lish vaqt qatorida xatoni taxminan ikki baravar kam ko'rsatdi, kelajakni ko'radigan belgi esa validatsiyada eng yaxshi, ishlab chiqarishda esa sezilarli yomon bo'ldi.

Keyingi darsda Eksperiment kuzatuvi: nega model_final_v3_haqiqiy.pkl ishlamaydi, har yurishda nimani qayd qilish kerak va sqlite ustida MLflow ga o'xshash kichik kuzatuvchini noldan quramiz — giperparametr qidiruvi, eng yaxshi yurishni so'rov bilan topish va "g'olib la'nati".

Ulashish:Telegram'da

Izohlar (0)

Izoh yozish uchun kiring.

  • Hozircha izoh yo'q. Birinchi bo'ling!
27.3-dars: Ma'lumot quvuri — IlmHamroh