Verify a 50,000-row CSV with Python, asyncio and httpx
- Who it is for
- Developers and data teams with large CSV lists
- The problem
- Checking one address per request is slow and runs into the rate limit. A naive bulk loop can open more jobs than the account allows, pay twice after a crash, or quietly drop rows whose results never came back.
- The solution
- Normalise and de-duplicate the column, send bulk jobs of up to 5,000 addresses, keep at most five open, poll every job to the end, wait out rate limits, stop cleanly when credits run out, and save the job ids so that a second run resumes instead of paying again.
- What you get
- Your CSV with three new columns (verify_status, verify_score, verify_reason) and a small .jobs.json file that makes the run resumable.
Addresses, IDs and results in the examples are illustrative. example.com is reserved for documentation, so a real check of these addresses returns invalid.
How the bulk endpoints behave
The script is built on these rules of the Cold Leads API v1. All calls need a secret key (sk_…) of a Business account, sent in the x-api-key header (Authorization: Bearer works too).
- POST /api/v1/verify/bulk takes a JSON body with an emails array (and an optional job name) and answers 202 with a job id. A job holds up to 5,000 addresses. The API lower-cases and de-duplicates the list and ignores strings that are not e-mail addresses.
- At most five jobs per account can be queued or running at the same time, including jobs started in the web app. One more answers 429 with the error too_many_jobs.
- GET /api/v1/verify/jobs/{id} returns the status (queued, running, done or failed), done, total and counts for valid, risky and invalid. Every status request also checks the next batch of up to 25 addresses before it answers, and a scheduled worker advances open jobs in the background, so steady polling finishes a job sooner.
- The same request with results=1 adds one entry per address: email, status, score and reason (the last reason code of the check). Results come back as JSON only; the script writes the CSV.
- Credits are checked when a job is created (402 no_credits if the balance is smaller than the job) and charged as addresses are processed. A job that runs out of credits stops with the status failed and keeps the results it already has.
The script
Save it as verify_csv.py. It needs Python 3.9 or newer and the httpx package; everything else is in the standard library.
"""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()))Run it
- Create a secret key in Cold Leads under Settings → API keys. It starts with sk_, only an admin of a Business account can create it, and it is shown once.
- Install httpx and export the key as COLDLEADS_API_KEY.
- Run the script with the input file, the output file and the name of the column that holds the addresses.
- Read the first line it prints: rows, unique addresses, jobs, the credits needed and your balance. If the balance is short, it stops before creating any job.
- If the run is interrupted, start the same command again. The job ids are stored next to the output (leads_verified.csv.jobs.json in the example), so existing jobs are polled instead of being created and charged again.
pip install httpx
export COLDLEADS_API_KEY=sk_your_secret_key
python verify_csv.py leads.csv leads_verified.csv --column emailWhat the script does when something goes wrong
| Situation | API answer | What the script does |
|---|---|---|
| Too many requests for the key | 429 rate_limited with Retry-After | Pauses every request for the given seconds (the window lasts 60 seconds from its first request), then continues. It paces itself to at most 100 requests a minute, so this only happens when other tools use the same key. |
| Five jobs already open | 429 too_many_jobs | Waits 30 seconds and tries to create the job again. |
| Balance too small for the next job | 402 no_credits | Creates no further jobs, lets the open ones finish and marks the remaining rows not_checked. |
| Credits ran out during a job | Job status failed | Downloads the results the job has and marks the rest not_checked. |
| Server or network error on a status request | 5xx or no answer | Tries again after 10 seconds; status requests are safe to repeat. |
| Network error while creating a job | No answer | Stops instead of retrying, because the job may already exist. The saved job ids let the next run continue. |
| Wrong key, no API in the plan, suspended account | 401 or 403 | Stops and prints the error code. |
Limits and costs
- 1 credit per address in a job, whatever the result, and also when the result comes from the 30-day result cache. The script removes duplicates and cells that are not addresses before sending, so those cost nothing.
- 50,000 unique addresses cost 50,000 credits. The Business plan ($99 a month) includes the API and 10,000 credits a month; the other 40,000 are 40 extra packs of 1,000 credits at $5 each, which is $200. If none of the month's 10,000 credits were used yet, the list costs $299 that month; tax may apply.
- Monthly credits are counted per calendar month (UTC), and purchased packs are used after them. An admin buys packs in Settings → Billing, from 1 to 100 packs per purchase.
- Up to 5,000 addresses per job, five queued or running jobs per account, 120 requests per minute per key.
Reading the new columns
| verify_status | verify_reason | Meaning |
|---|---|---|
| valid | ok | A mail server accepted the address over SMTP and did not accept a made-up address on the same domain. |
| valid | smtp_unreachable | Syntax and MX records are fine, but no mail server could be reached over SMTP, so the mailbox itself was not checked. |
| risky | catch_all | The server accepts any address on the domain, so this mailbox cannot be confirmed. |
| risky | smtp_unknown | No definite answer from the mail server, for example a temporary 4xx reply. |
| risky | ok, smtp_unreachable or smtp_unknown | The domain is on the built-in list of disposable-mail domains; the score is 40 or lower. |
| invalid | mailbox_missing | The server rejected the recipient with a 5xx reply. |
| invalid | no_mx | The domain has neither an MX nor an A record, so it cannot receive mail. |
| invalid | syntax | The address breaks the syntax rules the API checks (length, characters, dots). |
| skipped | not_an_email | Added by the script: the cell is not an e-mail address and was not sent, so no credit was used. |
| (empty) | not_checked | Added by the script: the job stopped before this address, usually because credits ran out. |
Role addresses such as info@ or sales@ stay in their bucket with a lower score: 80 instead of 97 for ok, 65 instead of 75 for smtp_unreachable. Bulk results carry only the last reason code; POST /api/v1/verify returns the full list, including disposable and role.
FAQ
Why not call POST /api/v1/verify once per row?
It works, but every address is then its own request, and a key allows 120 requests per minute: at most 7,200 addresses an hour with nothing else running. A bulk job takes 5,000 addresses in one request, and the status requests that follow do the work in batches.
Does running the script twice cost credits twice?
Not while the .jobs.json file next to the output exists: the second run polls the same jobs and downloads their results. If you delete it, new jobs are created and every address is charged again, even when the results come from the 30-day cache.
Can the API return a CSV file?
No. Job results come back as JSON from GET /api/v1/verify/jobs/{id}?results=1; the script turns them into columns next to your original rows.
What about lists larger than 50,000 rows?
The script works the same for any size: it splits the list into jobs of 5,000 and keeps five open. The cost stays at 1 credit per unique address, so check the credit line it prints before the first job starts.