Cold Leads

Перевірка 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_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()))

Запуск

  1. Створіть секретний ключ у Cold Leads в Налаштування → API-ключі. Він починається з sk_, створити його може лише адміністратор акаунта Business, і показується він один раз.
  2. Встановіть httpx і експортуйте ключ як COLDLEADS_API_KEY.
  3. Запустіть скрипт, указавши вхідний файл, вихідний файл і назву колонки з адресами.
  4. Прочитайте перший рядок, який виводить скрипт: кількість рядків, унікальних адрес, завдань, потрібних кредитів і ваш баланс. Якщо балансу не вистачає, скрипт зупиняється, не створивши жодного завдання.
  5. Якщо запуск перервався, виконайте ту саму команду ще раз. id завдань зберігаються поруч із вихідним файлом (у прикладі — leads_verified.csv.jobs.json), тож наявні завдання опитуються, а не створюються й оплачуються знову.
bash
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_statusverify_reasonЗначення
validokПоштовий сервер прийняв адресу через SMTP і не прийняв вигадану адресу на тому самому домені.
validsmtp_unreachableСинтаксис і MX-записи в порядку, але жоден поштовий сервер не був досяжний через SMTP, тож саму скриньку не перевірено.
riskycatch_allСервер приймає будь-яку адресу на домені, тож цю скриньку неможливо підтвердити.
riskysmtp_unknownПоштовий сервер не дав однозначної відповіді, наприклад повернув тимчасову відповідь 4xx.
riskyok, smtp_unreachable або smtp_unknownДомен є у вбудованому списку одноразових поштових доменів; оцінка 40 або нижче.
invalidmailbox_missingСервер відхилив одержувача відповіддю 5xx.
invalidno_mxДомен не має ні MX-, ні A-запису, тож не може приймати пошту.
invalidsyntaxАдреса порушує правила синтаксису, які перевіряє API (довжина, символи, крапки).
skippednot_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 кредит за унікальну адресу, тож перевірте рядок про кредити, який скрипт виводить до старту першого завдання.