Cold Leads

Verifique un CSV de 50 000 filas con Python, asyncio y httpx

Para quién
Desarrolladores y equipos de datos con listas CSV grandes
El problema
Comprobar una dirección por petición es lento y choca con el límite de peticiones. Un bucle masivo ingenuo puede abrir más tareas de las que permite la cuenta, pagar dos veces tras un fallo o perder en silencio filas cuyos resultados nunca llegaron.
La solución
Normalice la columna y elimine duplicados, envíe tareas masivas de hasta 5 000 direcciones, mantenga como máximo cinco abiertas, consulte cada tarea hasta el final, espere cuando salte el límite de peticiones, deténgase limpiamente cuando se acaben los créditos y guarde los id de las tareas para que una segunda ejecución continúe en lugar de volver a pagar.
Qué obtiene
Su CSV con tres columnas nuevas (verify_status, verify_score, verify_reason) y un pequeño archivo .jobs.json que permite reanudar la ejecución.

Las direcciones, los identificadores y los resultados de los ejemplos son ilustrativos. example.com está reservado para documentación, así que una comprobación real de estas direcciones devuelve invalid.

Los ejemplos de código son iguales en todos los idiomas; sus comentarios están en inglés.

Cómo se comportan los endpoints masivos

El script se basa en estas reglas de la API v1 de Cold Leads. Todas las llamadas necesitan una clave secreta (sk_…) de una cuenta Business, enviada en la cabecera x-api-key (Authorization: Bearer también funciona).

  • POST /api/v1/verify/bulk recibe un cuerpo JSON con un array emails (y un nombre de tarea opcional) y responde 202 con un id de tarea. Una tarea admite hasta 5 000 direcciones. La API pasa la lista a minúsculas, elimina duplicados e ignora las cadenas que no son direcciones de correo.
  • Como máximo cinco tareas por cuenta pueden estar en cola o en curso a la vez, incluidas las iniciadas en la aplicación web. Una más responde 429 con el error too_many_jobs.
  • GET /api/v1/verify/jobs/{id} devuelve el estado (queued, running, done o failed), done, total y los recuentos de valid, risky e invalid. Cada petición de estado también comprueba el siguiente lote de hasta 25 direcciones antes de responder, y un proceso programado hace avanzar las tareas abiertas en segundo plano, así que un sondeo constante termina antes una tarea.
  • La misma petición con results=1 añade una entrada por dirección: email, status, score y reason (el último código de motivo de la comprobación). Los resultados solo llegan en JSON; el CSV lo escribe el script.
  • Los créditos se comprueban al crear una tarea (402 no_credits si el saldo es menor que la tarea) y se cobran a medida que se procesan las direcciones. Una tarea que se queda sin créditos se detiene con el estado failed y conserva los resultados que ya tiene.

El script

Guárdelo como verify_csv.py. Necesita Python 3.9 o posterior y el paquete httpx; todo lo demás está en la biblioteca estándar.

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

Ejecútelo

  1. Cree una clave secreta en Cold Leads, en Ajustes → Claves de API. Empieza por sk_, solo un admin de una cuenta Business puede crearla y se muestra una sola vez.
  2. Instale httpx y exporte la clave como COLDLEADS_API_KEY.
  3. Ejecute el script con el archivo de entrada, el archivo de salida y el nombre de la columna que contiene las direcciones.
  4. Lea la primera línea que imprime: filas, direcciones únicas, tareas, los créditos necesarios y su saldo. Si el saldo no alcanza, se detiene antes de crear ninguna tarea.
  5. Si la ejecución se interrumpe, lance de nuevo el mismo comando. Los id de las tareas se guardan junto a la salida (leads_verified.csv.jobs.json en el ejemplo), así que las tareas existentes se consultan en lugar de crearse y cobrarse otra vez.
bash
pip install httpx
export COLDLEADS_API_KEY=sk_your_secret_key
python verify_csv.py leads.csv leads_verified.csv --column email

Qué hace el script cuando algo falla

SituaciónRespuesta de la APIQué hace el script
Demasiadas peticiones para la clave429 rate_limited con Retry-AfterPausa todas las peticiones durante los segundos indicados (la ventana dura 60 segundos desde su primera petición) y luego continúa. Se limita a sí mismo a 100 peticiones por minuto como máximo, así que esto solo ocurre cuando otras herramientas usan la misma clave.
Ya hay cinco tareas abiertas429 too_many_jobsEspera 30 segundos y vuelve a intentar crear la tarea.
Saldo insuficiente para la siguiente tarea402 no_creditsNo crea más tareas, deja terminar las abiertas y marca las filas restantes como not_checked.
Los créditos se acabaron durante una tareaEstado de la tarea failedDescarga los resultados que tiene la tarea y marca el resto como not_checked.
Error del servidor o de red en una petición de estado5xx o sin respuestaVuelve a intentarlo tras 10 segundos; las peticiones de estado se pueden repetir sin riesgo.
Error de red al crear una tareaSin respuestaSe detiene en lugar de reintentar, porque puede que la tarea ya exista. Los id de tarea guardados permiten continuar en la siguiente ejecución.
Clave incorrecta, plan sin API, cuenta suspendida401 o 403Se detiene e imprime el código de error.

Límites y costes

  • 1 crédito por dirección de una tarea, sea cual sea el resultado, también cuando el resultado sale de la caché de resultados de 30 días. El script elimina los duplicados y las celdas que no son direcciones antes de enviar, así que esos no cuestan nada.
  • 50 000 direcciones únicas cuestan 50 000 créditos. El plan Business ($99 al mes) incluye la API y 10 000 créditos al mes; los otros 40 000 son 40 paquetes extra de 1 000 créditos a $5 cada uno, es decir, $200. Si aún no se ha usado ninguno de los 10 000 créditos del mes, la lista cuesta $299 ese mes; pueden aplicarse impuestos.
  • Los créditos mensuales se cuentan por mes natural (UTC), y los paquetes comprados se usan después de ellos. Un admin compra paquetes en Ajustes → Facturación, de 1 a 100 paquetes por compra.
  • Hasta 5 000 direcciones por tarea, cinco tareas en cola o en curso por cuenta, 120 peticiones por minuto por clave.

Cómo leer las columnas nuevas

verify_statusverify_reasonSignificado
validokUn servidor de correo aceptó la dirección por SMTP y no aceptó una dirección inventada del mismo dominio.
validsmtp_unreachableLa sintaxis y los registros MX son correctos, pero no se pudo conectar con ningún servidor de correo por SMTP, así que el propio buzón no se comprobó.
riskycatch_allEl servidor acepta cualquier dirección del dominio, así que este buzón no se puede confirmar.
riskysmtp_unknownSin respuesta concluyente del servidor de correo, por ejemplo una respuesta temporal 4xx.
riskyok, smtp_unreachable o smtp_unknownEl dominio está en la lista integrada de dominios de correo desechable; la puntuación es de 40 o menos.
invalidmailbox_missingEl servidor rechazó al destinatario con una respuesta 5xx.
invalidno_mxEl dominio no tiene registro MX ni registro A, así que no puede recibir correo.
invalidsyntaxLa dirección incumple las reglas de sintaxis que comprueba la API (longitud, caracteres, puntos).
skippednot_an_emailLo añade el script: la celda no es una dirección de correo y no se envió, así que no se usó ningún crédito.
(vacío)not_checkedLo añade el script: la tarea se detuvo antes de esta dirección, normalmente porque se acabaron los créditos.

Las direcciones de rol como info@ o sales@ se quedan en su grupo con una puntuación menor: 80 en lugar de 97 para ok, 65 en lugar de 75 para smtp_unreachable. Los resultados masivos solo llevan el último código de motivo; POST /api/v1/verify devuelve la lista completa, incluidos disposable y role.

Preguntas frecuentes

¿Por qué no llamar a POST /api/v1/verify una vez por fila?

Funciona, pero entonces cada dirección es una petición propia, y una clave admite 120 peticiones por minuto: como máximo 7 200 direcciones por hora sin nada más en marcha. Una tarea masiva recibe 5 000 direcciones en una sola petición, y las peticiones de estado que siguen hacen el trabajo por lotes.

¿Ejecutar el script dos veces cuesta el doble de créditos?

No mientras exista el archivo .jobs.json junto a la salida: la segunda ejecución consulta las mismas tareas y descarga sus resultados. Si lo borra, se crean tareas nuevas y cada dirección se cobra de nuevo, aunque los resultados salgan de la caché de 30 días.

¿Puede la API devolver un archivo CSV?

No. Los resultados de las tareas llegan en JSON desde GET /api/v1/verify/jobs/{id}?results=1; el script los convierte en columnas junto a sus filas originales.

¿Y las listas de más de 50 000 filas?

El script funciona igual con cualquier tamaño: divide la lista en tareas de 5 000 y mantiene cinco abiertas. El coste sigue siendo 1 crédito por dirección única, así que revise la línea de créditos que imprime antes de que empiece la primera tarea.