AI Engineering Academy · Oppitunti

Eräkäsittely asynkronisesti ja jonojen avulla

Rakenna asyncioa ja tehtäväjonoa käyttävä asynkroninen poimintaputki, joka käsittelee tuhansia dokumentteja rinnakkain nopeusrajoituksia noudattaen ja etenemistä seuraten.

Oppitunti 3/413 vaihetta

Eräkäsittely asynkronisesti ja jonojen avulla on ilmainen AI Engineering Academy-oppitunti CoddyKitissä. Tämä on oppitunti 3/4. Voit lukea koko oppitunnin alta ilmaiseksi ja harjoitella sen jälkeen käytännössä selaimessa sisäänrakennetulla koodieditorilla ja ympäri vuorokauden käytettävissä olevan tekoälytuutorin avulla. Oppitunti kuuluu AI Engineering Academy-oppimispolkuun, ja edistymisesi synkronoituu verkon ja CoddyKit-sovelluksen välillä. AI Engineering Academy-kurssilla on yhteensä 4 oppituntia.

Miksi eräkäsittely on tärkeää

Tuhansien asiakirjojen käsittely yksi kerrallaan on tuotannossa liian hidasta. OpenAI API:a peräkkäin kutsuva synkroninen silmukka saattaa käsitellä yhden asiakirjan sekunnissa, jolloin 10 000 asiakirjan käsittely kestää lähes kolme tuntia. Asynkroninen eräkäsittely voi suorittaa samanaikaisesti satoja pyyntöjä, mikä lyhentää kokonaiskestoa suuruusluokan verran.

asyncion perusteet

Pythonin asyncio-tapahtumasilmukan avulla voitte suorittaa useita I/O-sidonnaisia tehtäviä samanaikaisesti ilman säikeitä. Kun yksi API-kutsu odottaa verkkovastausta, tapahtumasilmukka siirtyy käsittelemään toista kutsua. Kirjoitatte koodin async def- ja await-avainsanoilla, ja suoritusympäristö huolehtii ajoituksesta. Tämä sopii erinomaisesti LLM-kutsuille, jotka viettävät suurimman osan ajastaan palvelimen vastausta odottaen.

import asyncio
import instructor
from openai import AsyncOpenAI

async_client = instructor.from_openai(AsyncOpenAI())

async def extract_one(text: str) -> PersonExtract:
    return await async_client.chat.completions.create(
        model='gpt-4o-mini',
        response_model=PersonExtract,
        messages=[{'role': 'user', 'content': text}]
    )

Useiden poimintojen suorittaminen gatherilla

asyncio.gather suorittaa korutiinien luettelon samanaikaisesti ja palauttaa kaikki tulokset, kun viimeinen niistä valmistuu. Pienelle dokumenttierälle tämä riittää. Kokoa poimintakorutiinit list comprehension -rakenteella ja välitä ne gather-funktiolle. Kokonaisaika vastaa suunnilleen hitaimman yksittäisen kutsun aikaa, ei kaikkien kutsujen yhteenlaskettua aikaa.

async def batch_extract(texts: list) -> list:
    tasks = [extract_one(text) for text in texts]
    results = await asyncio.gather(*tasks, return_exceptions=True)
    # Filter out exceptions
    successes = [r for r in results if not isinstance(r, Exception)]
    failures = [r for r in results if isinstance(r, Exception)]
    print(f'Success: {len(successes)}, Failures: {len(failures)}')
    return successes

results = asyncio.run(batch_extract(documents))

Samanaikaisuuden hallinta semaforeilla

Tuhansien pyyntöjen lähettäminen samanaikaisesti ylittää nopeusrajoitukset ja aiheuttaa 429-virheitä. Käytä asyncio.Semaphore-semaforia samanaikaisten API-kutsujen määrän rajoittamiseen. Semaforin arvo 50 tarkoittaa, että käynnissä voi olla enintään 50 kutsua kerrallaan. Säädä tätä lukua OpenAI-tasosi ja kohdemallin nopeusrajoituksen perusteella.

import asyncio

sem = asyncio.Semaphore(50)  # max 50 concurrent calls

async def extract_with_limit(text: str, semaphore: asyncio.Semaphore):
    async with semaphore:
        return await extract_one(text)

async def batch_extract_limited(texts: list):
    tasks = [extract_with_limit(t, sem) for t in texts]
    return await asyncio.gather(*tasks, return_exceptions=True)

Eksponentiaalinen odotus nopeusrajoitusvirheissä

Semaforista huolimatta nopeusrajoitukset voivat tulla vastaan liikennepiikkien aikana. Toteuta eksponentiaalinen odotus: odota ennen uutta yritystä ensin 1 sekunti, sitten 2, 4 ja 8 sekuntia. Lisää viiveeseen satunnaisvaihtelua (pieni satunnainen poikkeama), jotta kaikki samanaikaiset kutsujat eivät yritä uudelleen täsmälleen samaan aikaan ja aiheuta uutta piikkiä. tenacity-kirjasto helpottaa tätä.

from tenacity import retry, wait_exponential, stop_after_attempt, retry_if_exception_type
from openai import RateLimitError

@retry(
    wait=wait_exponential(multiplier=1, min=1, max=60),
    stop=stop_after_attempt(5),
    retry=retry_if_exception_type(RateLimitError)
)
async def extract_with_retry(text: str):
    return await extract_one(text)

Työjonon käyttäminen suurissa erissä

Muutamaa tuhatta kohdetta suuremmille erille pysyvä työjono on parempi ratkaisu kuin asyncio.gather. Redis Queue (RQ)-, Celery- ja Dramatiq-tyyppiset jonot säilyttävät työt uudelleenkäynnistysten yli, mahdollistavat horisontaalisen skaalauksen useilla työntekijöillä ja tarjoavat näkyvyyden töiden tilaan ja virheisiin. Työntekijät noutavat työt jonosta ja kutsuvat API:a itsenäisesti.

# With Redis Queue (RQ)
from rq import Queue
from redis import Redis

redis_conn = Redis()
q = Queue('extractions', connection=redis_conn)

def enqueue_documents(doc_ids: list):
    for doc_id in doc_ids:
        q.enqueue(
            'workers.extract_document',
            doc_id,
            job_timeout=120,
            result_ttl=3600
        )

enqueue_documents(all_doc_ids)

Edistymisen seuraaminen tietokannassa

Pitkäkestoiset eräajot tarvitsevat edistymisen seurannan, jotta voit valvoa tilaa, tunnistaa jumiutuneet työt ja jatkaa käsittelyä virheiden jälkeen. Käytä tietokannassasi tilataulua, jossa on kentät dokumentin tunnisteelle, tilalle (pending, processing, completed, failed), aikaleimoille ja virheilmoituksille. Päivitä tila atomisesti jokaisen poimintakutsun ympärillä.

import asyncpg

async def process_document(pool, doc_id: str, text: str):
    async with pool.acquire() as conn:
        await conn.execute(
            'UPDATE extractions SET status=$1, started_at=NOW() WHERE doc_id=$2',
            'processing', doc_id
        )
        try:
            result = await extract_one(text)
            await conn.execute(
                'UPDATE extractions SET status=$1, result=$2, completed_at=NOW() WHERE doc_id=$3',
                'completed', result.model_dump_json(), doc_id
            )
        except Exception as e:
            await conn.execute(
                'UPDATE extractions SET status=$1, error=$2 WHERE doc_id=$3',
                'failed', str(e), doc_id
            )

Epäonnistuneiden töiden jatkaminen

Eräajon on oltava turvallisesti uudelleenkäynnistettävissä. Hae käynnistyksen yhteydessä tietokannasta dokumentit, joiden tila on pending tai failed, ja yritä niitä uudelleen. Käytä dokumenttikohtaista idempotency key -avainta, jotta saman dokumentin vahingossa tapahtuva kaksoisjonoitus havaitaan: toinen yritys tunnistaa valmiin tuloksen ja ohittaa uudelleenkäsittelyn. Tämä estää päällekkäiset kirjoitukset jatkokäsittelyjärjestelmiin.

async def get_pending_docs(pool) -> list:
    async with pool.acquire() as conn:
        rows = await conn.fetch(
            'SELECT doc_id, raw_text FROM extractions WHERE status IN ($1, $2)',
            'pending', 'failed'
        )
    return [dict(row) for row in rows]

async def resume_batch(pool):
    docs = await get_pending_docs(pool)
    print(f'Resuming {len(docs)} unprocessed documents')
    tasks = [process_document(pool, d['doc_id'], d['raw_text']) for d in docs]
    await asyncio.gather(*tasks, return_exceptions=True)

Erien käsittely OpenAI Batch API:lla

OpenAI:n Batch API mahdollistaa enintään 50 000 pyynnön lähettämisen yhdessä tiedostossa, ja tulokset saa 24 tunnin kuluessa 50 %:n alennuksella. Tämä sopii erinomaisesti ei-kiireellisiin poimintaputkiin, joissa kustannukset ovat tärkeämpiä kuin viive. Lähetät JSONL-tiedoston, joka sisältää pyynnöt, tarkistat valmistumista ja lataat tulostiedoston.

from openai import OpenAI
import json

client = OpenAI()

# Build JSONL batch file
with open('/tmp/batch_requests.jsonl', 'w') as f:
    for i, text in enumerate(documents):
        request = {
            'custom_id': f'doc_{i}',
            'method': 'POST',
            'url': '/v1/chat/completions',
            'body': {
                'model': 'gpt-4o-mini',
                'messages': [{'role': 'user', 'content': text}]
            }
        }
        f.write(json.dumps(request) + '\n')

# Upload and submit
batch_file = client.files.create(file=open('/tmp/batch_requests.jsonl', 'rb'), purpose='batch')
batch = client.batches.create(input_file_id=batch_file.id, endpoint='/v1/chat/completions', completion_window='24h')
print(batch.id)

Suorituskyvyn ja kustannusten seuranta

Seuraa eräajojen aikana poiminnan suorituskykyä (dokumenttia minuutissa) ja dokumenttikohtaista kustannusta. Jaa API:n kokonaiskulut käsiteltyjen dokumenttien määrällä saadaksesi kustannusvertailuarvon. Skaalatessasi etsi lineaarista kustannusten kasvua — superlineaarinen kasvu viittaa siihen, että tuhlaat tokeneita tarpeettoman pitkiin kehotteisiin. Yksinkertainen mittaristo auttaa havaitsemaan tehottomuudet ennen kuin ne kumuloituvat.

import time

class BatchMetrics:
    def __init__(self):
        self.start_time = time.time()
        self.processed = 0
        self.total_tokens = 0
        self.cost = 0.0

    def record(self, usage):
        self.processed += 1
        self.total_tokens += usage.total_tokens
        self.cost += usage.prompt_tokens * 0.00000015 + usage.completion_tokens * 0.0000006

    def report(self):
        elapsed = time.time() - self.start_time
        print(f'{self.processed} docs in {elapsed:.1f}s = {self.processed/elapsed:.1f} docs/sec')
        print(f'Cost: ${self.cost:.4f} = ${self.cost/self.processed:.6f} per doc')

Tokenien minuuttikohtaisten nopeusrajoitusten noudattaminen

OpenAI:n nopeusrajoitukset koskevat sekä pyyntöjä minuutissa (RPM) että tokeneita minuutissa (TPM). 50 samanaikaisen kutsun lähettäminen on RPM:n kannalta sopivaa, mutta jos jokainen kutsu käyttää 2 000 tokenia, 50 kutsua vastaa 100 000 tokenia minuutissa — tämä ylittää helposti tason 1 rajoitukset. Laske odotettujen tokenien määrä ennen lähettämistä tiktokenilla ja toteuta tokenibudjetti samanaikaisuutta rajoittavan semaforin rinnalle.

import tiktoken

enc = tiktoken.encoding_for_model('gpt-4o-mini')

def estimate_tokens(text: str) -> int:
    return len(enc.encode(text)) + 300  # +300 for schema + response

# TPM_LIMIT = 200_000  # Tier 2 limit
# Only submit a batch if estimated total tokens fits within budget
def fits_in_budget(texts: list, tpm_limit: int = 200_000) -> bool:
    total = sum(estimate_tokens(t) for t in texts)
    return total <= tpm_limit

Pikatarkistus

Testaa ymmärryksesi dokumenttien poiminnan asynkronisesta eräkäsittelystä.

Oppitunnin kertaus

Tässä oppitunnissa opit, että asyncio ja Semaphore mahdollistavat samanaikaiset API-kutsut nopeusrajoituksia noudattaen, työjonot ja tilataulut tekevät suurista eräajoista jatkettavia ja seurattavia ja OpenAI Batch API tarjoaa 50 %:n kustannussäästöt ei-kiireellisissä kuormissa 24 tunnin viiveen kustannuksella. Seuraavaksi käsittelemme skeeman muuttumista pitkäkestoisissa poimintaputkissa.

Aloita maksutta

Opi Python tekoälytuutorin avulla — ilmaiseksi

Kirjoita ja suorita oikeaa koodia selaimessa, saa välitöntä apua tekoälytuutorilta ympäri vuorokauden ja jatka siitä, mihin jäit, verkossa tai sovelluksessa.

Kurssit
30
Oppitunnit
120

Usein kysytyt kysymykset

Onko oppitunti ”Eräkäsittely asynkronisesti ja jonojen avulla” ilmainen?

Kyllä – oppitunnin ”Eräkäsittely asynkronisesti ja jonojen avulla” koko tekstin voi lukea täällä verkossa ilmaiseksi. Jos haluat harjoitella interaktiivisesti sisäänrakennetulla koodieditorilla ja ympäri vuorokauden käytettävissä olevan tekoälytuutorin avulla sekä avata koko AI Engineering Academy-kurssin, päivitä CoddyKit PROhon. AI Engineering Academy-kurssilla on yhteensä 4 oppituntia.

Mitä opin oppitunnilla ”Eräkäsittely asynkronisesti ja jonojen avulla”?

Rakenna asyncioa ja tehtäväjonoa käyttävä asynkroninen poimintaputki, joka käsittelee tuhansia dokumentteja rinnakkain nopeusrajoituksia noudattaen ja etenemistä seuraten. Harjoittelet AI Engineering Academy-aihetta koodilla, jonka suoritat suoraan selaimessa. Ympäri vuorokauden käytettävissä oleva tekoälytuutori vastaa kysymyksiisi oppitunnin aikana.

Tarvitsenko kokemusta aloittaakseni AI Engineering Academy-opiskelun?

Aiempi kokemus ei ole tarpeen. CoddyKitin AI Engineering Academy-oppimispolku sopii vasta-alkajista edistyneisiin, joten voit aloittaa tästä tai alusta ja edetä omaan tahtiisi. Tämä on oppitunti 3/4.

Kuinka kauan ”Eräkäsittely asynkronisesti ja jonojen avulla”-oppitunnin suorittaminen kestää?

Useimmat CoddyKitin oppitunnit kestävät noin 5–10 minuuttia. Jokainen oppitunti on lyhyt ja interaktiivinen, joten edistyt tasaisesti ja voit jatkaa siitä, mihin jäit – sekä verkossa että sovelluksessa.

Voinko kirjoittaa ja suorittaa koodia tällä AI Engineering Academy-oppitunnilla?

Kyllä. Jokainen AI Engineering Academy-oppitunti sisältää sisäänrakennetun koodieditorin, joten voit kirjoittaa ja suorittaa oikeaa koodia suoraan selaimessa ja saada välitöntä palautetta tekoälyltä – paikallista asennusta ei tarvita.

Kaikki tämän kurssin oppitunnit

  1. Instructor: tyypitetty poiminta Pydanticilla
  2. Osittaisen ja puuttuvan datan käsittely
  3. Eräkäsittely asynkronisesti ja jonojen avulla
  4. Skeeman kehitys ja taaksepäin yhteensopivuus
← Takaisin: AI Engineering Academy