Перевірка CSV на 50 000 рядків за допомогою Python, asyncio і httpx
- Для кого
- Розробники та команди даних із великими CSV-списками
- Проблема
- Перевіряти по одній адресі за запит повільно, і так ви впираєтеся в ліміт запитів. Наївний масовий цикл може відкрити більше завдань, ніж дозволяє акаунт, заплатити двічі після збою або непомітно загубити рядки, результати яких так і не повернулися.
- Рішення
- Нормалізуйте колонку й приберіть дублікати, надсилайте масові завдання до 5 000 адрес, тримайте відкритими не більше п’яти, опитуйте кожне завдання до кінця, перечікуйте ліміти запитів, коректно зупиняйтеся, коли закінчуються кредити, і зберігайте id завдань, щоб повторний запуск продовжував роботу, а не платив знову.
- Що ви отримаєте
- Ваш CSV із трьома новими колонками (verify_status, verify_score, verify_reason) і невеликий файл .jobs.json, завдяки якому запуск можна продовжити.
Адреси, ідентифікатори та результати в прикладах ілюстративні. Домен example.com зарезервовано для документації, тож реальна перевірка цих адрес поверне invalid.
Приклади коду однакові для всіх мов, коментарі в них — англійською.
Як працюють масові ендпоінти
Скрипт спирається на такі правила Cold Leads API v1. Усім викликам потрібен секретний ключ (sk_…) акаунта Business у заголовку x-api-key (Authorization: Bearer теж працює).
- POST /api/v1/verify/bulk приймає JSON-тіло з масивом emails (і необов’язковою назвою завдання) та відповідає 202 з id завдання. Завдання вміщує до 5 000 адрес. API переводить список у нижній регістр, прибирає дублікати й ігнорує значення, які не є e-mail-адресами.
- Одночасно в черзі або в роботі може бути не більше п’яти завдань на акаунт, включно із запущеними у веб-застосунку. На спробу створити ще одне API відповідає 429 з помилкою too_many_jobs.
- GET /api/v1/verify/jobs/{id} повертає статус (queued, running, done або failed), done, total і лічильники для valid, risky та invalid. Кожен запит статусу перед відповіддю ще й перевіряє наступну порцію до 25 адрес, а обробник за розкладом просуває відкриті завдання у фоні, тож регулярне опитування завершує завдання швидше.
- Той самий запит із results=1 додає по одному запису на адресу: email, status, score і reason (останній код причини перевірки). Результати повертаються лише як JSON; CSV записує скрипт.
- Кредити перевіряються під час створення завдання (402 no_credits, якщо баланс менший за завдання) і списуються в міру обробки адрес. Завдання, якому забракло кредитів, зупиняється зі статусом failed і зберігає вже отримані результати.
Скрипт
Збережіть його як verify_csv.py. Потрібні Python 3.9 або новіший і пакет httpx; усе інше є в стандартній бібліотеці.
"""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()))Запуск
- Створіть секретний ключ у Cold Leads в Налаштування → API-ключі. Він починається з sk_, створити його може лише адміністратор акаунта Business, і показується він один раз.
- Встановіть httpx і експортуйте ключ як COLDLEADS_API_KEY.
- Запустіть скрипт, указавши вхідний файл, вихідний файл і назву колонки з адресами.
- Прочитайте перший рядок, який виводить скрипт: кількість рядків, унікальних адрес, завдань, потрібних кредитів і ваш баланс. Якщо балансу не вистачає, скрипт зупиняється, не створивши жодного завдання.
- Якщо запуск перервався, виконайте ту саму команду ще раз. id завдань зберігаються поруч із вихідним файлом (у прикладі — leads_verified.csv.jobs.json), тож наявні завдання опитуються, а не створюються й оплачуються знову.
pip install httpx
export COLDLEADS_API_KEY=sk_your_secret_key
python verify_csv.py leads.csv leads_verified.csv --column emailЩо робить скрипт, коли щось іде не так
| Ситуація | Відповідь API | Що робить скрипт |
|---|---|---|
| Забагато запитів для ключа | 429 rate_limited з Retry-After | Призупиняє всі запити на вказану кількість секунд (вікно триває 60 секунд від першого запиту в ньому), потім продовжує. Скрипт сам тримає темп не вище 100 запитів на хвилину, тож це трапляється лише тоді, коли той самий ключ використовують інші інструменти. |
| Уже відкрито п’ять завдань | 429 too_many_jobs | Чекає 30 секунд і знову пробує створити завдання. |
| Балансу не вистачає на наступне завдання | 402 no_credits | Не створює нових завдань, дає відкритим завершитися й позначає решту рядків як not_checked. |
| Кредити закінчилися під час завдання | Статус завдання failed | Завантажує результати, які вже має завдання, і позначає решту як not_checked. |
| Помилка сервера або мережі під час запиту статусу | 5xx або немає відповіді | Повторює через 10 секунд; запити статусу можна безпечно повторювати. |
| Помилка мережі під час створення завдання | Немає відповіді | Зупиняється замість повтору, бо завдання вже могло бути створене. Збережені id завдань дозволяють наступному запуску продовжити. |
| Неправильний ключ, API немає в тарифі, акаунт призупинено | 401 або 403 | Зупиняється й виводить код помилки. |
Ліміти та вартість
- 1 кредит за адресу в завданні, незалежно від результату, і навіть тоді, коли результат береться з 30-денного кешу результатів. Скрипт прибирає дублікати й клітинки, що не є адресами, ще до надсилання, тож вони нічого не коштують.
- 50 000 унікальних адрес коштують 50 000 кредитів. Тариф Business ($99 на місяць) містить API та 10 000 кредитів на місяць; решта 40 000 — це 40 додаткових пакетів по 1 000 кредитів по $5 кожен, тобто $200. Якщо з місячних 10 000 кредитів ще не використано жодного, список обійдеться в $299 за цей місяць; може додаватися податок.
- Місячні кредити рахуються за календарний місяць (UTC), а куплені пакети використовуються після них. Адміністратор купує пакети в Налаштування → Оплата — від 1 до 100 пакетів за одну покупку.
- До 5 000 адрес у завданні, п’ять завдань у черзі або в роботі на акаунт, 120 запитів на хвилину на ключ.
Як читати нові колонки
| verify_status | verify_reason | Значення |
|---|---|---|
| valid | ok | Поштовий сервер прийняв адресу через SMTP і не прийняв вигадану адресу на тому самому домені. |
| valid | smtp_unreachable | Синтаксис і MX-записи в порядку, але жоден поштовий сервер не був досяжний через SMTP, тож саму скриньку не перевірено. |
| risky | catch_all | Сервер приймає будь-яку адресу на домені, тож цю скриньку неможливо підтвердити. |
| risky | smtp_unknown | Поштовий сервер не дав однозначної відповіді, наприклад повернув тимчасову відповідь 4xx. |
| risky | ok, smtp_unreachable або smtp_unknown | Домен є у вбудованому списку одноразових поштових доменів; оцінка 40 або нижче. |
| invalid | mailbox_missing | Сервер відхилив одержувача відповіддю 5xx. |
| invalid | no_mx | Домен не має ні MX-, ні A-запису, тож не може приймати пошту. |
| invalid | syntax | Адреса порушує правила синтаксису, які перевіряє API (довжина, символи, крапки). |
| skipped | not_an_email | Додає скрипт: клітинка не є e-mail-адресою і не надсилалася, тож кредит не витрачено. |
| (порожньо) | not_checked | Додає скрипт: завдання зупинилося до цієї адреси, зазвичай тому, що закінчилися кредити. |
Рольові адреси на кшталт info@ чи sales@ лишаються у своєму кошику з нижчою оцінкою: 80 замість 97 для ok, 65 замість 75 для smtp_unreachable. Масові результати містять лише останній код причини; POST /api/v1/verify повертає повний список, зокрема disposable і role.
Питання
Чому б не викликати POST /api/v1/verify для кожного рядка?
Так теж можна, але тоді кожна адреса — окремий запит, а ключ дозволяє 120 запитів на хвилину: максимум 7 200 адрес на годину, якщо більше нічого не працює. Масове завдання приймає 5 000 адрес одним запитом, а подальші запити статусу виконують роботу порціями.
Чи списуються кредити двічі, якщо запустити скрипт двічі?
Ні, поки поруч із вихідним файлом є файл .jobs.json: другий запуск опитує ті самі завдання й завантажує їхні результати. Якщо його видалити, створюються нові завдання, і кожна адреса оплачується знову, навіть коли результати беруться з 30-денного кешу.
Чи може API повернути CSV-файл?
Ні. Результати завдання повертаються як JSON із GET /api/v1/verify/jobs/{id}?results=1; скрипт перетворює їх на колонки поруч із вашими вихідними рядками.
А як щодо списків понад 50 000 рядків?
Скрипт працює однаково для будь-якого розміру: ділить список на завдання по 5 000 і тримає п’ять відкритими. Вартість лишається 1 кредит за унікальну адресу, тож перевірте рядок про кредити, який скрипт виводить до старту першого завдання.