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 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í
- 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.
- Nainstalujte httpx a exportujte klíč jako COLDLEADS_API_KEY.
- Spusťte skript se vstupním souborem, výstupním souborem a názvem sloupce, ve kterém jsou adresy.
- 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.
- 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.
pip install httpx
export COLDLEADS_API_KEY=sk_your_secret_key
python verify_csv.py leads.csv leads_verified.csv --column emailCo skript dělá, když se něco pokazí
| Situace | Odpověď API | Co skript udělá |
|---|---|---|
| Příliš mnoho požadavků na klíč | 429 rate_limited s Retry-After | Pozastaví 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 úloh | 429 too_many_jobs | Počká 30 sekund a zkusí úlohu vytvořit znovu. |
| Zůstatek nestačí na další úlohu | 402 no_credits | Nevytváří další úlohy, nechá otevřené doběhnout a zbývající řádky označí not_checked. |
| Během úlohy došly kredity | Stav úlohy failed | Stáhne výsledky, které úloha má, a zbytek označí not_checked. |
| Chyba serveru nebo sítě při dotazu na stav | 5xx 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ý účet | 401 nebo 403 | Skončí 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_status | verify_reason | Význam |
|---|---|---|
| valid | ok | Poštovní server adresu přes SMTP přijal a smyšlenou adresu na stejné doméně nepřijal. |
| valid | smtp_unreachable | Syntaxe 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. |
| risky | catch_all | Server přijímá jakoukoli adresu na doméně, takže tuto schránku nelze potvrdit. |
| risky | smtp_unknown | Žádná jednoznačná odpověď poštovního serveru, například dočasná odpověď 4xx. |
| risky | ok, smtp_unreachable nebo smtp_unknown | Doména je na vestavěném seznamu jednorázových e-mailových domén; skóre je 40 nebo nižší. |
| invalid | mailbox_missing | Server příjemce odmítl odpovědí 5xx. |
| invalid | no_mx | Doména nemá MX ani A záznam, takže nemůže přijímat poštu. |
| invalid | syntax | Adresa porušuje syntaktická pravidla, která API kontroluje (délka, znaky, tečky). |
| skipped | not_an_email | Doplňuje skript: buňka není e-mailová adresa a nebyla odeslána, takže se nespotřeboval žádný kredit. |
| (prázdné) | not_checked | Doplň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.