Cold Leads

Ověřte CSV s 50 000 řádky v Pythonu s asyncio a httpx

Pro koho
Vývojáři a datové týmy s velkými seznamy v CSV
Problém
Kontrola jedné adresy na požadavek je pomalá a naráží na limit požadavků. Naivní hromadná smyčka může otevřít víc úloh, než účet povoluje, po pádu zaplatit dvakrát nebo potichu vynechat řádky, jejichž výsledky se nikdy nevrátily.
Řešení
Sloupec normalizujte a odstraňte z něj duplicity, posílejte hromadné úlohy až po 5 000 adresách, držte otevřených nejvýš pět, dotazujte se na každou úlohu až do jejího konce, přečkejte limity požadavků, při vyčerpání kreditů čistě skončete a uložte id úloh, aby druhé spuštění navázalo, místo aby platilo znovu.
Co získáte
Vaše CSV se třemi novými sloupci (verify_status, verify_score, verify_reason) a malý soubor .jobs.json, díky kterému lze běh po přerušení obnovit.

Adresy, identifikátory a výsledky v příkladech jsou ilustrativní. Doména example.com je vyhrazená pro dokumentaci, takže skutečná kontrola těchto adres vrátí invalid.

Ukázky kódu jsou ve všech jazycích stejné, komentáře v nich jsou anglicky.

Jak se chovají hromadné endpointy

Skript stojí na těchto pravidlech API Cold Leads v1. Všechna volání potřebují tajný klíč (sk_…) účtu s tarifem Business, posílaný v hlavičce x-api-key (funguje i Authorization: Bearer).

  • POST /api/v1/verify/bulk přijímá JSON tělo s polem emails (a volitelným názvem úlohy) a odpovídá 202 s id úlohy. Úloha pojme až 5 000 adres. API seznam převede na malá písmena, odstraní duplicity a ignoruje řetězce, které nejsou e-mailovými adresami.
  • Ve frontě nebo v běhu může být současně nejvýš pět úloh na účet, včetně úloh spuštěných ve webové aplikaci. Na další API odpoví 429 s chybou too_many_jobs.
  • GET /api/v1/verify/jobs/{id} vrací status (queued, running, done nebo failed), done, total a počty pro valid, risky a invalid. Každý dotaz na stav navíc před odpovědí zkontroluje další dávku až 25 adres a plánovaný worker posouvá otevřené úlohy na pozadí, takže pravidelné dotazování úlohu dokončí dřív.
  • Stejný požadavek s results=1 přidá jednu položku na adresu: email, status, score a reason (poslední kód důvodu kontroly). Výsledky přicházejí jen jako JSON; CSV zapisuje skript.
  • Kredity se kontrolují při vytvoření úlohy (402 no_credits, pokud je zůstatek menší než úloha) a strhávají se průběžně, jak se adresy zpracovávají. Úloha, které dojdou kredity, skončí se stavem failed a výsledky, které už má, si ponechá.

Skript

Uložte ho jako verify_csv.py. Potřebuje Python 3.9 nebo novější a balíček httpx; vše ostatní je ve standardní knihovně.

verify_csv.py
"""Verify a large CSV with the Cold Leads bulk API (Python 3.9+, httpx).

  pip install httpx
  export COLDLEADS_API_KEY=sk_...
  python verify_csv.py leads.csv leads_verified.csv --column email

Addresses are normalised and de-duplicated, sent in jobs of up to 5,000 with at most
5 jobs open at a time, and every job is polled until it is done. The output is the
input CSV plus verify_status, verify_score and verify_reason. Job ids are saved next
to the output file, so a second run resumes the same jobs instead of paying twice.
"""
from __future__ import annotations

import argparse
import asyncio
import csv
import hashlib
import json
import os
import re
import sys
import time
from pathlib import Path

import httpx

API_BASE = os.environ.get("COLDLEADS_API_BASE", "https://coldleads.app")
JOB_SIZE = 5000    # addresses per job (API maximum)
MAX_JOBS = 5       # queued or running jobs per account (API maximum)
POLL_PAUSE = 3.0   # seconds between two polls of the same job
MIN_GAP = 0.6      # seconds between any two requests: at most 100 a minute (limit: 120 per key)
# The pattern the API uses to pick addresses out of a list. Cells that do not match are
# skipped here and cost nothing.
EMAIL_RE = re.compile(r"[a-z0-9!#$%&'*+/=?^_`{|}~.-]+@[a-z0-9.-]+\.[a-z]{2,}")


class ApiError(Exception):
    pass


class Api:
    """HTTP calls paced across all tasks, with the documented 429 and 5xx handling."""

    def __init__(self, client: httpx.AsyncClient):
        self.client = client
        self.lock = asyncio.Lock()
        self.next_at = 0.0

    async def call(self, method: str, path: str, **kwargs) -> tuple[int, dict]:
        while True:
            async with self.lock:
                delay = self.next_at - time.monotonic()
                if delay > 0:
                    await asyncio.sleep(delay)
                self.next_at = time.monotonic() + MIN_GAP
            try:
                r = await self.client.request(method, path, **kwargs)
            except httpx.TransportError as e:
                if method != "GET":
                    # a POST may have reached the server: retrying could open a second job
                    raise ApiError(f"network error on {method} {path}: {e!r}") from e
                print(f"network error on {path}, retrying in 10 s", file=sys.stderr)
                await asyncio.sleep(10)
                continue
            try:
                data = r.json()
            except ValueError:
                data = {}
            if r.status_code == 429 and data.get("error") == "rate_limited":
                wait = int(r.headers.get("retry-after") or data.get("retry_after_seconds") or 60)
                print(f"rate limited, pausing {wait} s", file=sys.stderr)
                self.next_at = time.monotonic() + wait  # every task waits out the window
                continue
            if r.status_code >= 500 and method == "GET":
                await asyncio.sleep(10)
                continue
            return r.status_code, data


async def run_job(api, slots, stop, idx, emails, state, save_state, results):
    async with slots:
        job_id = state["jobs"].get(str(idx))
        while not job_id:
            if stop.is_set():
                return
            status, data = await api.call("POST", "/api/v1/verify/bulk", json={"emails": emails, "name": f"{state['name']} #{idx + 1}"})
            if status == 202:
                job_id = data["id"]
                state["jobs"][str(idx)] = job_id
                save_state()
            elif status == 429 and data.get("error") == "too_many_jobs":
                await asyncio.sleep(30)  # other jobs of this account are still open
            elif status == 402:
                print(f"job {idx + 1}: not enough credits ({data.get('credits')}), stopping", file=sys.stderr)
                stop.set()
            else:
                raise ApiError(f"job {idx + 1}: HTTP {status} {data}")
        shown = -1
        while True:
            status, job = await api.call("GET", f"/api/v1/verify/jobs/{job_id}")
            if status != 200:
                raise ApiError(f"job {job_id}: HTTP {status} {job}")
            if job["status"] in ("done", "failed"):
                break
            if job["done"] // 500 != shown // 500:
                shown = job["done"]
                print(f"job {idx + 1}: {job['done']}/{job['total']} {job['counts']}")
            await asyncio.sleep(POLL_PAUSE)  # each poll also processes the next addresses
        status, job = await api.call("GET", f"/api/v1/verify/jobs/{job_id}", params={"results": "1"})
        if status != 200:
            raise ApiError(f"job {job_id}: HTTP {status} {job}")
        for r in job.get("results", []):
            results[r["email"]] = r
        if job["status"] == "failed":
            print(f"job {idx + 1} stopped early (usually: credits ran out); kept {len(job.get('results', []))} results", file=sys.stderr)


async def main(argv=None) -> int:
    p = argparse.ArgumentParser(description="Verify a CSV column of e-mail addresses with Cold Leads")
    p.add_argument("input")
    p.add_argument("output")
    p.add_argument("--column", default="email")
    args = p.parse_args(argv)
    key = os.environ.get("COLDLEADS_API_KEY", "").strip()
    if not key.startswith("sk_"):
        print("Set COLDLEADS_API_KEY to a secret key (sk_...) from Settings -> API keys", file=sys.stderr)
        return 2

    with open(args.input, newline="", encoding="utf-8-sig") as f:
        reader = csv.DictReader(f)
        rows = list(reader)
        fields = list(reader.fieldnames or [])
    col = next((c for c in fields if c.strip().lower() == args.column.lower()), None)
    if col is None:
        print(f"column {args.column!r} not found in {fields}", file=sys.stderr)
        return 2
    norm = [(row.get(col) or "").strip().lower().strip("<>") for row in rows]
    unique = list(dict.fromkeys(e for e in norm if EMAIL_RE.fullmatch(e)))
    chunks = [unique[i:i + JOB_SIZE] for i in range(0, len(unique), JOB_SIZE)]

    state_path = Path(args.output + ".jobs.json")
    fingerprint = hashlib.sha256("\n".join(unique).encode()).hexdigest()
    state = json.loads(state_path.read_text()) if state_path.exists() else {"fingerprint": fingerprint, "name": Path(args.input).stem[:60], "jobs": {}}
    if state["fingerprint"] != fingerprint:
        print(f"{state_path} belongs to a different list; delete it to start over", file=sys.stderr)
        return 2

    def save_state():
        state_path.write_text(json.dumps(state, indent=2))

    headers = {"x-api-key": key, "User-Agent": "coldleads-csv-example/1.0"}
    timeout = httpx.Timeout(120.0, connect=10.0)  # a poll verifies a batch before it answers
    results: dict[str, dict] = {}
    async with httpx.AsyncClient(base_url=API_BASE, headers=headers, timeout=timeout) as client:
        api = Api(client)
        status, credits = await api.call("GET", "/api/v1/credits")
        if status != 200:
            print(f"API error {status}: {credits}", file=sys.stderr)
            return 1
        needed = sum(len(c) for i, c in enumerate(chunks) if str(i) not in state["jobs"])
        print(f"{len(rows)} rows, {len(unique)} unique addresses in {len(chunks)} jobs; {needed} credits needed, {credits['left']} left")
        if credits["left"] < needed:
            print(f"Not enough credits: short by {needed - credits['left']}. Add packs in Settings -> Billing.", file=sys.stderr)
            return 2
        slots, stop = asyncio.Semaphore(MAX_JOBS), asyncio.Event()
        try:
            await asyncio.gather(*(run_job(api, slots, stop, i, c, state, save_state, results) for i, c in enumerate(chunks)))
        except ApiError as e:
            print(f"{e}\nJob ids are saved in {state_path}: run the same command again to resume.", file=sys.stderr)
            return 1

    out_fields = fields + [c for c in ("verify_status", "verify_score", "verify_reason") if c not in fields]
    counts: dict[str, int] = {}
    with open(args.output, "w", newline="", encoding="utf-8") as f:
        w = csv.DictWriter(f, fieldnames=out_fields, extrasaction="ignore")
        w.writeheader()
        for row, email in zip(rows, norm):
            r = results.get(email)
            if r:
                row.update(verify_status=r["status"], verify_score=r["score"], verify_reason=r["reason"])
            elif not EMAIL_RE.fullmatch(email):
                row.update(verify_status="skipped", verify_score="", verify_reason="not_an_email")
            else:
                row.update(verify_status="", verify_score="", verify_reason="not_checked")
            counts[row["verify_status"] or "not_checked"] = counts.get(row["verify_status"] or "not_checked", 0) + 1
            w.writerow(row)
    print(f"wrote {args.output}: {counts}")
    return 0


if __name__ == "__main__":
    sys.exit(asyncio.run(main()))

Spuštění

  1. V Cold Leads vytvořte tajný klíč v Nastavení → API klíče. Začíná na sk_, vytvořit ho může jen admin účtu s tarifem Business a zobrazí se jen jednou.
  2. Nainstalujte httpx a exportujte klíč jako COLDLEADS_API_KEY.
  3. Spusťte skript se vstupním souborem, výstupním souborem a názvem sloupce, ve kterém jsou adresy.
  4. Přečtěte si první řádek, který vypíše: řádky, unikátní adresy, úlohy, potřebné kredity a váš zůstatek. Pokud zůstatek nestačí, skončí dřív, než vytvoří jakoukoli úlohu.
  5. Pokud se běh přeruší, spusťte stejný příkaz znovu. Id úloh se ukládají vedle výstupu (v příkladu leads_verified.csv.jobs.json), takže se skript dotazuje na existující úlohy, místo aby je znovu vytvořil a znovu zaplatil.
bash
pip install httpx
export COLDLEADS_API_KEY=sk_your_secret_key
python verify_csv.py leads.csv leads_verified.csv --column email

Co skript dělá, když se něco pokazí

SituaceOdpověď APICo skript udělá
Příliš mnoho požadavků na klíč429 rate_limited s Retry-AfterPozastaví všechny požadavky na uvedený počet sekund (okno trvá 60 sekund od svého prvního požadavku) a pak pokračuje. Sám se drží nejvýš na 100 požadavcích za minutu, takže k tomu dojde, jen když stejný klíč používají i jiné nástroje.
Už je otevřeno pět úloh429 too_many_jobsPočká 30 sekund a zkusí úlohu vytvořit znovu.
Zůstatek nestačí na další úlohu402 no_creditsNevytváří další úlohy, nechá otevřené doběhnout a zbývající řádky označí not_checked.
Během úlohy došly kredityStav úlohy failedStáhne výsledky, které úloha má, a zbytek označí not_checked.
Chyba serveru nebo sítě při dotazu na stav5xx nebo žádná odpověďZkusí to znovu po 10 sekundách; dotazy na stav lze bezpečně opakovat.
Chyba sítě při vytváření úlohyŽádná odpověďMísto opakování skončí, protože úloha už může existovat. Uložená id úloh umožní dalšímu spuštění pokračovat.
Špatný klíč, tarif bez API, pozastavený účet401 nebo 403Skončí a vypíše kód chyby.

Limity a náklady

  • 1 kredit za adresu v úloze bez ohledu na výsledek, i když výsledek pochází z 30denní cache výsledků. Skript před odesláním odstraní duplicity a buňky, které nejsou adresami, takže ty nic nestojí.
  • 50 000 unikátních adres stojí 50 000 kreditů. Tarif Business ($99 měsíčně) zahrnuje API a 10 000 kreditů měsíčně; zbylých 40 000 pokryje 40 dalších balíčků po 1 000 kreditech, každý za $5 – celkem $200. Pokud z 10 000 měsíčních kreditů ještě žádný nebyl využit, stojí seznam ten měsíc $299; může se připočíst daň.
  • Měsíční kredity se počítají za kalendářní měsíc (UTC) a zakoupené balíčky se čerpají až po nich. Balíčky kupuje admin v Nastavení → Platby, 1 až 100 balíčků na jeden nákup.
  • Až 5 000 adres na úlohu, pět úloh ve frontě nebo v běhu na účet, 120 požadavků za minutu na klíč.

Jak číst nové sloupce

verify_statusverify_reasonVýznam
validokPoštovní server adresu přes SMTP přijal a smyšlenou adresu na stejné doméně nepřijal.
validsmtp_unreachableSyntaxe a MX záznamy jsou v pořádku, ale přes SMTP se nepodařilo spojit s žádným poštovním serverem, takže samotná schránka se nekontrolovala.
riskycatch_allServer přijímá jakoukoli adresu na doméně, takže tuto schránku nelze potvrdit.
riskysmtp_unknownŽádná jednoznačná odpověď poštovního serveru, například dočasná odpověď 4xx.
riskyok, smtp_unreachable nebo smtp_unknownDoména je na vestavěném seznamu jednorázových e-mailových domén; skóre je 40 nebo nižší.
invalidmailbox_missingServer příjemce odmítl odpovědí 5xx.
invalidno_mxDoména nemá MX ani A záznam, takže nemůže přijímat poštu.
invalidsyntaxAdresa porušuje syntaktická pravidla, která API kontroluje (délka, znaky, tečky).
skippednot_an_emailDoplňuje skript: buňka není e-mailová adresa a nebyla odeslána, takže se nespotřeboval žádný kredit.
(prázdné)not_checkedDoplňuje skript: úloha se zastavila před touto adresou, obvykle proto, že došly kredity.

Rolové adresy jako info@ nebo sales@ zůstávají ve své skupině s nižším skóre: 80 místo 97 u ok, 65 místo 75 u smtp_unreachable. Hromadné výsledky nesou jen poslední kód důvodu; POST /api/v1/verify vrací celý seznam včetně disposable a role.

Otázky

Proč nevolat POST /api/v1/verify pro každý řádek zvlášť?

Funguje to, ale každá adresa je pak samostatný požadavek a klíč povoluje 120 požadavků za minutu: nejvýš 7 200 adres za hodinu, pokud nic jiného neběží. Hromadná úloha vezme 5 000 adres jedním požadavkem a následné dotazy na stav odvedou práci po dávkách.

Stojí dvojí spuštění skriptu dvakrát tolik kreditů?

Ne, dokud vedle výstupu existuje soubor .jobs.json: druhé spuštění se dotazuje na stejné úlohy a stáhne jejich výsledky. Pokud ho smažete, vytvoří se nové úlohy a každá adresa se účtuje znovu, i když výsledky pocházejí z 30denní cache.

Umí API vrátit soubor CSV?

Ne. Výsledky úloh přicházejí jako JSON z GET /api/v1/verify/jobs/{id}?results=1; skript z nich udělá sloupce vedle vašich původních řádků.

A co seznamy větší než 50 000 řádků?

Skript funguje stejně pro jakoukoli velikost: rozdělí seznam na úlohy po 5 000 a drží pět otevřených. Cena zůstává 1 kredit za unikátní adresu, takže před startem první úlohy zkontrolujte řádek s kredity, který vypíše.