Eine CSV mit 50.000 Zeilen per Python, asyncio und httpx prüfen
- Für wen
- Entwickler und Datenteams mit großen CSV-Listen
- Das Problem
- Eine Adresse pro Anfrage zu prüfen ist langsam und stößt an das Ratenlimit. Eine naive Bulk-Schleife kann mehr Aufträge öffnen, als das Konto erlaubt, nach einem Absturz doppelt bezahlen oder stillschweigend Zeilen verlieren, deren Ergebnisse nie zurückkamen.
- Die Lösung
- Normalisieren und deduplizieren Sie die Spalte, senden Sie Bulk-Aufträge mit bis zu 5.000 Adressen, halten Sie höchstens fünf offen, fragen Sie jeden Auftrag bis zum Ende ab, warten Sie Ratenlimits ab, stoppen Sie sauber, wenn die Credits ausgehen, und speichern Sie die Auftrags-IDs, damit ein zweiter Lauf fortsetzt, statt erneut zu bezahlen.
- Was Sie bekommen
- Ihre CSV mit drei neuen Spalten (verify_status, verify_score, verify_reason) und eine kleine .jobs.json-Datei, die den Lauf fortsetzbar macht.
Adressen, IDs und Ergebnisse in den Beispielen sind illustrativ. example.com ist für Dokumentation reserviert, eine echte Prüfung dieser Adressen liefert daher invalid.
Die Codebeispiele sind in allen Sprachen gleich, ihre Kommentare sind auf Englisch.
Wie sich die Bulk-Endpunkte verhalten
Das Skript beruht auf diesen Regeln der Cold Leads API v1. Alle Aufrufe brauchen einen geheimen Schlüssel (sk_…) eines Business-Kontos, gesendet im Header x-api-key (Authorization: Bearer funktioniert ebenfalls).
- POST /api/v1/verify/bulk nimmt einen JSON-Body mit einem Array emails (und einem optionalen Auftragsnamen) entgegen und antwortet mit 202 und einer Auftrags-ID. Ein Auftrag fasst bis zu 5.000 Adressen. Die API schreibt die Liste klein, entfernt Duplikate und ignoriert Zeichenketten, die keine E-Mail-Adressen sind.
- Pro Konto können höchstens fünf Aufträge gleichzeitig warten oder laufen, einschließlich in der Web-App gestarteter Aufträge. Ein weiterer wird mit 429 und dem Fehler too_many_jobs beantwortet.
- GET /api/v1/verify/jobs/{id} liefert den Status (queued, running, done oder failed), done, total und counts für valid, risky und invalid. Jede Statusanfrage prüft vor ihrer Antwort auch den nächsten Stapel von bis zu 25 Adressen, und ein geplanter Worker treibt offene Aufträge im Hintergrund voran; stetiges Abfragen schließt einen Auftrag also schneller ab.
- Dieselbe Anfrage mit results=1 ergänzt einen Eintrag pro Adresse: email, status, score und reason (der letzte Grundcode der Prüfung). Ergebnisse kommen nur als JSON zurück; die CSV schreibt das Skript.
- Credits werden beim Anlegen eines Auftrags geprüft (402 no_credits, wenn das Guthaben kleiner als der Auftrag ist) und abgebucht, während die Adressen verarbeitet werden. Ein Auftrag, dem die Credits ausgehen, endet mit dem Status failed und behält die Ergebnisse, die er schon hat.
Das Skript
Speichern Sie es als verify_csv.py. Es braucht Python 3.9 oder neuer und das Paket httpx; alles andere steckt in der Standardbibliothek.
"""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()))Ausführen
- Erstellen Sie in Cold Leads unter Einstellungen → API-Schlüssel einen geheimen Schlüssel. Er beginnt mit sk_, nur ein Admin eines Business-Kontos kann ihn erstellen, und er wird einmal angezeigt.
- Installieren Sie httpx und exportieren Sie den Schlüssel als COLDLEADS_API_KEY.
- Starten Sie das Skript mit der Eingabedatei, der Ausgabedatei und dem Namen der Spalte, die die Adressen enthält.
- Lesen Sie die erste Zeile, die es ausgibt: Zeilen, eindeutige Adressen, Aufträge, benötigte Credits und Ihr Guthaben. Reicht das Guthaben nicht, stoppt es, bevor es einen Auftrag anlegt.
- Wird der Lauf unterbrochen, starten Sie denselben Befehl erneut. Die Auftrags-IDs liegen neben der Ausgabe (im Beispiel leads_verified.csv.jobs.json), sodass bestehende Aufträge abgefragt statt erneut angelegt und berechnet werden.
pip install httpx
export COLDLEADS_API_KEY=sk_your_secret_key
python verify_csv.py leads.csv leads_verified.csv --column emailWas das Skript tut, wenn etwas schiefgeht
| Situation | Antwort der API | Was das Skript tut |
|---|---|---|
| Zu viele Anfragen für den Schlüssel | 429 rate_limited mit Retry-After | Pausiert alle Anfragen für die angegebenen Sekunden (das Fenster dauert 60 Sekunden ab seiner ersten Anfrage) und macht dann weiter. Es drosselt sich selbst auf höchstens 100 Anfragen pro Minute, das passiert also nur, wenn andere Tools denselben Schlüssel nutzen. |
| Bereits fünf Aufträge offen | 429 too_many_jobs | Wartet 30 Sekunden und versucht erneut, den Auftrag anzulegen. |
| Guthaben zu klein für den nächsten Auftrag | 402 no_credits | Legt keine weiteren Aufträge an, lässt die offenen zu Ende laufen und markiert die restlichen Zeilen mit not_checked. |
| Credits während eines Auftrags aufgebraucht | Auftragsstatus failed | Lädt die vorhandenen Ergebnisse des Auftrags herunter und markiert den Rest mit not_checked. |
| Server- oder Netzwerkfehler bei einer Statusanfrage | 5xx oder keine Antwort | Versucht es nach 10 Sekunden erneut; Statusanfragen lassen sich gefahrlos wiederholen. |
| Netzwerkfehler beim Anlegen eines Auftrags | Keine Antwort | Stoppt, statt es erneut zu versuchen, weil der Auftrag bereits existieren könnte. Mit den gespeicherten Auftrags-IDs setzt der nächste Lauf fort. |
| Falscher Schlüssel, keine API im Tarif, gesperrtes Konto | 401 oder 403 | Stoppt und gibt den Fehlercode aus. |
Limits und Kosten
- 1 Credit pro Adresse in einem Auftrag, unabhängig vom Ergebnis, auch wenn das Ergebnis aus dem 30-Tage-Ergebnis-Cache kommt. Das Skript entfernt Duplikate und Zellen, die keine Adressen sind, vor dem Senden, diese kosten also nichts.
- 50.000 eindeutige Adressen kosten 50.000 Credits. Der Business-Tarif ($99 im Monat) enthält die API und 10.000 Credits im Monat; die übrigen 40.000 sind 40 zusätzliche Pakete zu 1.000 Credits, jedes für $5 und zusammen $200. Wenn von den 10.000 Credits des Monats noch keine verbraucht sind, kostet die Liste in diesem Monat $299; es können Steuern anfallen.
- Monatliche Credits zählen pro Kalendermonat (UTC), gekaufte Pakete werden danach verbraucht. Ein Admin kauft Pakete unter Einstellungen → Abrechnung, 1 bis 100 Pakete pro Kauf.
- Bis zu 5.000 Adressen pro Auftrag, fünf wartende oder laufende Aufträge pro Konto, 120 Anfragen pro Minute und Schlüssel.
Die neuen Spalten lesen
| verify_status | verify_reason | Bedeutung |
|---|---|---|
| valid | ok | Ein Mailserver hat die Adresse per SMTP akzeptiert und eine erfundene Adresse derselben Domain nicht akzeptiert. |
| valid | smtp_unreachable | Syntax und MX-Einträge sind in Ordnung, aber kein Mailserver war per SMTP erreichbar, das Postfach selbst wurde also nicht geprüft. |
| risky | catch_all | Der Server akzeptiert jede Adresse der Domain, dieses Postfach lässt sich also nicht bestätigen. |
| risky | smtp_unknown | Keine eindeutige Antwort vom Mailserver, zum Beispiel eine vorübergehende 4xx-Antwort. |
| risky | ok, smtp_unreachable oder smtp_unknown | Die Domain steht auf der eingebauten Liste von Wegwerf-Mail-Domains; der Score ist 40 oder niedriger. |
| invalid | mailbox_missing | Der Server hat den Empfänger mit einer 5xx-Antwort abgelehnt. |
| invalid | no_mx | Die Domain hat weder einen MX- noch einen A-Eintrag und kann daher keine E-Mails empfangen. |
| invalid | syntax | Die Adresse verletzt die Syntaxregeln, die die API prüft (Länge, Zeichen, Punkte). |
| skipped | not_an_email | Vom Skript ergänzt: Die Zelle ist keine E-Mail-Adresse und wurde nicht gesendet, es wurde also kein Credit verbraucht. |
| (leer) | not_checked | Vom Skript ergänzt: Der Auftrag endete vor dieser Adresse, meist weil die Credits aufgebraucht waren. |
Rollenadressen wie info@ oder sales@ bleiben mit niedrigerem Score in ihrem Bereich: 80 statt 97 bei ok, 65 statt 75 bei smtp_unreachable. Bulk-Ergebnisse enthalten nur den letzten Grundcode; POST /api/v1/verify liefert die vollständige Liste, einschließlich disposable und role.
FAQ
Warum nicht POST /api/v1/verify einmal pro Zeile aufrufen?
Das funktioniert, aber dann ist jede Adresse eine eigene Anfrage, und ein Schlüssel erlaubt 120 Anfragen pro Minute: höchstens 7.200 Adressen pro Stunde, wenn sonst nichts läuft. Ein Bulk-Auftrag nimmt 5.000 Adressen in einer Anfrage entgegen, und die folgenden Statusanfragen erledigen die Arbeit in Stapeln.
Kostet ein zweiter Lauf des Skripts doppelt Credits?
Nicht, solange die .jobs.json-Datei neben der Ausgabe existiert: Der zweite Lauf fragt dieselben Aufträge ab und lädt ihre Ergebnisse herunter. Löschen Sie sie, werden neue Aufträge angelegt und jede Adresse wird erneut berechnet, auch wenn die Ergebnisse aus dem 30-Tage-Cache kommen.
Kann die API eine CSV-Datei liefern?
Nein. Auftragsergebnisse kommen als JSON von GET /api/v1/verify/jobs/{id}?results=1; das Skript macht daraus Spalten neben Ihren ursprünglichen Zeilen.
Was ist mit Listen über 50.000 Zeilen?
Das Skript funktioniert bei jeder Größe gleich: Es teilt die Liste in Aufträge zu 5.000 auf und hält fünf offen. Die Kosten bleiben bei 1 Credit pro eindeutiger Adresse; prüfen Sie also die Credit-Zeile, die es ausgibt, bevor der erste Auftrag startet.