Cold Leads

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_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()))

Ausführen

  1. 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.
  2. Installieren Sie httpx und exportieren Sie den Schlüssel als COLDLEADS_API_KEY.
  3. Starten Sie das Skript mit der Eingabedatei, der Ausgabedatei und dem Namen der Spalte, die die Adressen enthält.
  4. 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.
  5. 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.
bash
pip install httpx
export COLDLEADS_API_KEY=sk_your_secret_key
python verify_csv.py leads.csv leads_verified.csv --column email

Was das Skript tut, wenn etwas schiefgeht

SituationAntwort der APIWas das Skript tut
Zu viele Anfragen für den Schlüssel429 rate_limited mit Retry-AfterPausiert 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 offen429 too_many_jobsWartet 30 Sekunden und versucht erneut, den Auftrag anzulegen.
Guthaben zu klein für den nächsten Auftrag402 no_creditsLegt 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 aufgebrauchtAuftragsstatus failedLädt die vorhandenen Ergebnisse des Auftrags herunter und markiert den Rest mit not_checked.
Server- oder Netzwerkfehler bei einer Statusanfrage5xx oder keine AntwortVersucht es nach 10 Sekunden erneut; Statusanfragen lassen sich gefahrlos wiederholen.
Netzwerkfehler beim Anlegen eines AuftragsKeine AntwortStoppt, 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 Konto401 oder 403Stoppt 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_statusverify_reasonBedeutung
validokEin Mailserver hat die Adresse per SMTP akzeptiert und eine erfundene Adresse derselben Domain nicht akzeptiert.
validsmtp_unreachableSyntax und MX-Einträge sind in Ordnung, aber kein Mailserver war per SMTP erreichbar, das Postfach selbst wurde also nicht geprüft.
riskycatch_allDer Server akzeptiert jede Adresse der Domain, dieses Postfach lässt sich also nicht bestätigen.
riskysmtp_unknownKeine eindeutige Antwort vom Mailserver, zum Beispiel eine vorübergehende 4xx-Antwort.
riskyok, smtp_unreachable oder smtp_unknownDie Domain steht auf der eingebauten Liste von Wegwerf-Mail-Domains; der Score ist 40 oder niedriger.
invalidmailbox_missingDer Server hat den Empfänger mit einer 5xx-Antwort abgelehnt.
invalidno_mxDie Domain hat weder einen MX- noch einen A-Eintrag und kann daher keine E-Mails empfangen.
invalidsyntaxDie Adresse verletzt die Syntaxregeln, die die API prüft (Länge, Zeichen, Punkte).
skippednot_an_emailVom Skript ergänzt: Die Zelle ist keine E-Mail-Adresse und wurde nicht gesendet, es wurde also kein Credit verbraucht.
(leer)not_checkedVom 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.