uCheckeruChecker
13 мин чтения

Автоматизация email-гигиены через API: архитектура и код

Ручная очистка email-базы работает ровно до того момента, когда база перерастает тысячу адресов. Дальше начинается рутина: экспорт CSV, загрузка в сервис, ожидание, скачивание, импорт обратно. Каждый раз одно и то же. API позволяет убрать человека из этой цепочки и собрать пайплайн, который чистит базу по расписанию, без ручных действий.


Почему ручная гигиена не масштабируется

Email-база деградирует постоянно. Люди меняют работу, провайдеры отключают неактивные ящики, домены истекают. По разным оценкам, 2-3% адресов становятся невалидными каждый месяц. Для базы в 50 000 контактов это тысяча-полторы новых мёртвых адресов ежемесячно.

Маркетолог, который раз в квартал вспоминает про очистку, обнаруживает 5-8% невалидных адресов. К этому моменту bounce rate уже навредил репутации домена, ESP начал throttling, а часть писем улетела в спам даже у тех подписчиков, чьи адреса в порядке.

Автоматизация решает эту проблему: скрипт запускается по cron, отправляет адреса на проверку через API, получает результаты и обновляет базу. Без CSV, без веб-интерфейса, без участия человека.

Архитектура пайплайна

Типовой пайплайн автоматической гигиены состоит из четырёх этапов. Извлечение адресов из базы данных или CRM. Отправка на валидацию через API. Получение результатов. Обновление статусов в базе.

Между этапами есть нюансы. Не все адреса нужно проверять каждый раз - только те, которые давно не валидировались или недавно bounced. Результаты нужно не просто получить, а правильно интерпретировать: «risky» - не то же самое, что «invalid». И обновление базы должно быть идемпотентным - повторный запуск скрипта не должен ломать данные.

Рассмотрим каждый этап с кодом на Python и Node.js.

Этап 1. Выборка адресов для проверки

Проверять всю базу каждый раз - расточительно. Если адрес валидировался три дня назад и оказался good, его статус почти наверняка не изменился. Разумная стратегия - проверять адреса, которые не валидировались дольше N дней, плюс адреса, по которым был bounce с последней рассылки.

# Python: выборка адресов, не проверявшихся 30+ дней
import psycopg2
from datetime import datetime, timedelta

def get_stale_emails(conn, days: int = 30, limit: int = 50_000) -> list[str]:
    """Fetch emails that haven't been validated recently."""
    cutoff = datetime.utcnow() - timedelta(days=days)
    with conn.cursor() as cur:
        cur.execute("""
            SELECT email FROM subscribers
            WHERE status = 'active'
              AND (last_validated_at IS NULL OR last_validated_at < %s)
            ORDER BY last_validated_at ASC NULLS FIRST
            LIMIT %s
        """, (cutoff, limit))
        return [row[0] for row in cur.fetchall()]

То же на Node.js с PostgreSQL:

// Node.js: выборка адресов для проверки
const { Pool } = require("pg");
const pool = new Pool();

async function getStaleEmails(days = 30, limit = 50000) {
  const cutoff = new Date(Date.now() - days * 86400000);
  const { rows } = await pool.query(
    `SELECT email FROM subscribers
     WHERE status = 'active'
       AND (last_validated_at IS NULL OR last_validated_at < $1)
     ORDER BY last_validated_at ASC NULLS FIRST
     LIMIT $2`,
    [cutoff, limit]
  );
  return rows.map((r) => r.email);
}

Лимит в 50 000 адресов за запуск - разумный потолок. Если база больше, скрипт обработает следующую порцию при следующем запуске. За неделю ежедневных запусков пройдёт вся база в 350К адресов.

Сортировка NULLS FIRST гарантирует, что адреса, которые никогда не проверялись, уйдут на валидацию раньше тех, которые хотя бы раз прошли проверку.

Этап 2. Отправка на валидацию

У API валидации обычно два эндпоинта: single (один адрес за запрос) и bulk (массив адресов). Для автоматической гигиены bulk - единственный разумный вариант. Один HTTP-запрос вместо пятидесяти тысяч.

Bulk-эндпоинт принимает массив адресов и возвращает task_id. Валидация происходит асинхронно: сервер проверяет адреса в фоне, результаты забираются отдельным запросом. Размер батча ограничен - обычно 50 000 адресов. Если список больше, его нужно разбить.

# Python: отправка батчей на валидацию
import os
import time
import requests
from requests.adapters import HTTPAdapter
from urllib3.util.retry import Retry

API_KEY = os.environ["UCHECKER_API_KEY"]
BASE = "https://api.uchecker.net"
BATCH_SIZE = 10_000

session = requests.Session()
retry = Retry(total=3, backoff_factor=2, status_forcelist=[429, 500, 502, 503])
session.mount("https://", HTTPAdapter(max_retries=retry))

def chunk(lst: list, size: int):
    for i in range(0, len(lst), size):
        yield lst[i : i + size]

def submit_for_validation(emails: list[str]) -> list[int]:
    """Submit emails in batches, return list of task IDs."""
    task_ids = []
    batches = list(chunk(emails, BATCH_SIZE))

    for i, batch in enumerate(batches):
        resp = session.post(
            f"{BASE}/api/v1/validate/bulk",
            headers={"x-api-key": API_KEY},
            json={"emails": batch},
        )
        resp.raise_for_status()
        data = resp.json()
        task_ids.append(data["task_id"])
        print(f"Batch {i+1}/{len(batches)}: task_id={data['task_id']}, "
              f"queued={data['valid_emails']}")
        time.sleep(1)  # avoid rate limits between batches

    return task_ids

Node.js-версия:

// Node.js: отправка батчей на валидацию
const API_KEY = process.env.UCHECKER_API_KEY;
const BASE = "https://api.uchecker.net";
const BATCH_SIZE = 10_000;

function chunk(arr, size) {
  const result = [];
  for (let i = 0; i < arr.length; i += size) {
    result.push(arr.slice(i, i + size));
  }
  return result;
}

async function submitForValidation(emails) {
  const batches = chunk(emails, BATCH_SIZE);
  const taskIds = [];

  for (let i = 0; i < batches.length; i++) {
    const resp = await fetch(`${BASE}/api/v1/validate/bulk`, {
      method: "POST",
      headers: {
        "x-api-key": API_KEY,
        "Content-Type": "application/json",
      },
      body: JSON.stringify({ emails: batches[i] }),
    });

    if (!resp.ok) throw new Error(`HTTP ${resp.status}: ${await resp.text()}`);
    const data = await resp.json();
    taskIds.push(data.task_id);
    console.log(`Batch ${i + 1}/${batches.length}: task_id=${data.task_id}`);

    if (i < batches.length - 1) {
      await new Promise((r) => setTimeout(r, 1000));
    }
  }
  return taskIds;
}

Retry-логика на уровне HTTP обязательна. Сеть ненадёжна, серверы иногда отвечают 500 под нагрузкой. Без retry скрипт упадёт на третьем батче из двадцати, и вы узнаете об этом утром из логов cron. В Python это решается через urllib3 Retry. В Node.js - через обёртку с повтором или библиотеку вроде p-retry.

Этап 3. Получение результатов

После отправки батча сервер ставит задачу в очередь. Проверка 10 000 адресов занимает 3-5 минут. Результаты забираются через polling или webhook.

Polling проще в реализации: периодически опрашиваете эндпоинт задачи, пока статус не станет completed. С экспоненциальным backoff polling не нагружает API: начинаете с 5 секунд, увеличиваете интервал до минуты.

# Python: polling задачи с экспоненциальным backoff
def poll_task(task_id: int, timeout: int = 900) -> dict:
    """Wait for task completion. Returns results on success."""
    interval = 5
    max_interval = 60
    elapsed = 0

    while elapsed < timeout:
        resp = session.get(
            f"{BASE}/api/v1/tasks/{task_id}",
            headers={"x-api-key": API_KEY},
        )
        data = resp.json()

        if data["status"] == "completed":
            return data
        if data["status"] == "failed":
            raise RuntimeError(f"Task {task_id} failed: {data.get('error')}")

        time.sleep(interval)
        elapsed += interval
        interval = min(interval * 1.5, max_interval)

    raise TimeoutError(f"Task {task_id} timed out after {timeout}s")
// Node.js: polling задачи
async function pollTask(taskId, timeoutMs = 900000) {
  let interval = 5000;
  const maxInterval = 60000;
  const start = Date.now();

  while (Date.now() - start < timeoutMs) {
    const resp = await fetch(`${BASE}/api/v1/tasks/${taskId}`, {
      headers: { "x-api-key": API_KEY },
    });
    const data = await resp.json();

    if (data.status === "completed") return data;
    if (data.status === "failed") {
      throw new Error(`Task ${taskId} failed: ${data.error}`);
    }

    await new Promise((r) => setTimeout(r, interval));
    interval = Math.min(interval * 1.5, maxInterval);
  }
  throw new Error(`Task ${taskId} timed out`);
}

Когда задача завершена, забираем детальные результаты:

# Python: получение результатов задачи
def fetch_results(task_id: int) -> list[dict]:
    """Get detailed validation results for a completed task."""
    resp = session.get(
        f"{BASE}/api/v1/tasks/{task_id}/results",
        headers={"x-api-key": API_KEY},
        params={"format": "json"},
    )
    resp.raise_for_status()
    return resp.json()["data"]
    # Each item: {"email": "...", "validation_result": "good"|"bad"|"risky",
    #             "reason": "mailbox_not_found"|"disposable"|...}

Если задач несколько, опрашивайте их параллельно. В Python - через concurrent.futures.ThreadPoolExecutor. В Node.js - через Promise.allSettled. Последовательный опрос пяти задач по 5 минут каждая - 25 минут. Параллельный - 5 минут.

// Node.js: параллельный polling нескольких задач
async function pollAllTasks(taskIds) {
  const results = await Promise.allSettled(
    taskIds.map((id) => pollTask(id))
  );
  const completed = [];
  for (const [i, result] of results.entries()) {
    if (result.status === "fulfilled") {
      completed.push(taskIds[i]);
    } else {
      console.error(`Task ${taskIds[i]}: ${result.reason.message}`);
    }
  }
  return completed;
}

Этап 4. Обновление базы данных

Результаты валидации нужно записать обратно в базу. Здесь важно не просто пометить адрес как «валидный» или «невалидный», а сохранить причину и дату проверки. Это даёт возможность фильтровать по типу проблемы и отслеживать динамику.

# Python: запись результатов в базу
from datetime import datetime

def update_subscriber_statuses(conn, results: list[dict]):
    """Update subscribers based on validation results."""
    now = datetime.utcnow()
    good, bad, risky = 0, 0, 0

    with conn.cursor() as cur:
        for r in results:
            email = r["email"]
            verdict = r["validation_result"]  # "good", "bad", "risky"
            reason = r.get("reason", "")

            if verdict == "bad":
                cur.execute("""
                    UPDATE subscribers
                    SET status = 'invalid', validation_reason = %s,
                        last_validated_at = %s
                    WHERE email = %s
                """, (reason, now, email))
                bad += 1
            elif verdict == "risky":
                cur.execute("""
                    UPDATE subscribers
                    SET status = 'risky', validation_reason = %s,
                        last_validated_at = %s
                    WHERE email = %s
                """, (reason, now, email))
                risky += 1
            else:
                cur.execute("""
                    UPDATE subscribers
                    SET status = 'active', validation_reason = NULL,
                        last_validated_at = %s
                    WHERE email = %s
                """, (now, email))
                good += 1

    conn.commit()
    print(f"Updated: {good} good, {bad} bad, {risky} risky")
    return {"good": good, "bad": bad, "risky": risky}

Для больших объёмов поштучный UPDATE медленный. Быстрее собрать результаты в три списка и обновить пакетно:

# Python: пакетное обновление через executemany
def update_batch(conn, results: list[dict]):
    """Batch update - faster for large result sets."""
    now = datetime.utcnow()
    bad_emails = [(r["reason"], now, r["email"])
                  for r in results if r["validation_result"] == "bad"]
    risky_emails = [(r["reason"], now, r["email"])
                    for r in results if r["validation_result"] == "risky"]
    good_emails = [(now, r["email"])
                   for r in results if r["validation_result"] == "good"]

    with conn.cursor() as cur:
        if bad_emails:
            cur.executemany("""
                UPDATE subscribers
                SET status='invalid', validation_reason=%s, last_validated_at=%s
                WHERE email=%s
            """, bad_emails)
        if risky_emails:
            cur.executemany("""
                UPDATE subscribers
                SET status='risky', validation_reason=%s, last_validated_at=%s
                WHERE email=%s
            """, risky_emails)
        if good_emails:
            cur.executemany("""
                UPDATE subscribers
                SET status='active', validation_reason=NULL, last_validated_at=%s
                WHERE email=%s
            """, good_emails)

    conn.commit()
    print(f"Batch update: {len(good_emails)} good, "
          f"{len(bad_emails)} bad, {len(risky_emails)} risky")

Отдельный вопрос - что делать с «risky». Это адреса, которые формально существуют, но несут повышенный риск: catch-all домены, одноразовые адреса, ящики с подозрительными паттернами. Удалять их сразу - потеряете часть живых подписчиков. Оставлять без внимания - получите bounce. Рабочий компромисс: выделить risky в отдельный сегмент и отправлять им рассылку с пониженной частотой. Если адрес bounce - перевести в invalid.

Собираем пайплайн

Полный скрипт, который можно запустить вручную или по cron:

#!/usr/bin/env python3
"""Automated email hygiene via uChecker API.

Run daily via cron:
  0 3 * * * /usr/bin/python3 /opt/scripts/email_hygiene.py >> /var/log/email_hygiene.log 2>&1
"""
import os
import sys
import psycopg2
from datetime import datetime

# ... functions from above: get_stale_emails, submit_for_validation,
#     poll_task, fetch_results, update_batch

def run():
    conn = psycopg2.connect(os.environ["DATABASE_URL"])

    # 1. Get emails that need checking
    emails = get_stale_emails(conn, days=30, limit=50_000)
    if not emails:
        print(f"{datetime.utcnow()}: No stale emails. Exiting.")
        return

    print(f"{datetime.utcnow()}: Found {len(emails)} emails to validate")

    # 2. Check credit balance
    balance = session.get(
        f"{BASE}/api/v1/account/balance",
        headers={"x-api-key": API_KEY},
    ).json()["credits_remaining"]

    if balance < len(emails):
        print(f"Insufficient credits: {balance} < {len(emails)}")
        sys.exit(1)

    # 3. Submit for validation
    task_ids = submit_for_validation(emails)

    # 4. Wait for results
    all_results = []
    for tid in task_ids:
        poll_task(tid)
        results = fetch_results(tid)
        all_results.extend(results)

    # 5. Update database
    update_batch(conn, all_results)

    conn.close()
    print(f"{datetime.utcnow()}: Done. Processed {len(all_results)} emails.")

if __name__ == "__main__":
    run()

Node.js-эквивалент для тех, у кого бэкенд на JavaScript:

// email-hygiene.js - Node.js version
const { Pool } = require("pg");
const pool = new Pool();

async function run() {
  const emails = await getStaleEmails(30, 50000);
  if (!emails.length) {
    console.log("No stale emails. Exiting.");
    return;
  }
  console.log(`Found ${emails.length} emails to validate`);

  const taskIds = await submitForValidation(emails);
  const completedIds = await pollAllTasks(taskIds);

  for (const tid of completedIds) {
    const results = await fetchResults(tid);
    await updateSubscribers(pool, results);
  }

  console.log("Done.");
  await pool.end();
}

run().catch((err) => {
  console.error(err);
  process.exit(1);
});

Настройка cron и мониторинг

Запуск в 3 часа ночи ежедневно - стандартный выбор. Нагрузка на базу данных минимальна, API не перегружен дневным трафиком, результаты готовы к утру.

# crontab -e
0 3 * * * /usr/bin/python3 /opt/scripts/email_hygiene.py >> /var/log/email_hygiene.log 2>&1

# Или раз в неделю, если база маленькая
0 3 * * 1 /usr/bin/python3 /opt/scripts/email_hygiene.py >> /var/log/email_hygiene.log 2>&1

Скрипт, который молча падает, хуже скрипта, которого нет. Добавьте уведомления о сбоях. Простейший вариант - отправка алерта в Slack при ненулевом exit code:

# Обёртка для cron с алертом в Slack
#!/bin/bash
SLACK_WEBHOOK="https://hooks.slack.com/services/YOUR/WEBHOOK/URL"

/usr/bin/python3 /opt/scripts/email_hygiene.py >> /var/log/email_hygiene.log 2>&1
EXIT_CODE=$?

if [ $EXIT_CODE -ne 0 ]; then
  curl -s -X POST "$SLACK_WEBHOOK" \
    -H "Content-Type: application/json" \
    -d "{\"text\":\"Email hygiene script failed with exit code $EXIT_CODE\"}"
fi

Для полноценного мониторинга имеет смысл отправлять метрики после каждого запуска: сколько адресов проверено, сколько оказалось невалидными, сколько кредитов потрачено, сколько времени заняла обработка. Если процент невалидных адресов внезапно вырос с 3% до 15% - это сигнал. Либо что-то сломалось в сборе адресов, либо произошла массовая деактивация ящиков на конкретном провайдере.

Webhook вместо polling

Polling удобен для скриптов, которые запускаются по cron и ждут результата. Но если у вас есть веб-сервер, webhook эффективнее: API сам отправит POST-запрос на ваш эндпоинт, когда задача завершится. Не нужно держать процесс и опрашивать статус.

# Отправка батча с webhook
resp = session.post(
    f"{BASE}/api/v1/validate/bulk",
    headers={"x-api-key": API_KEY},
    json={
        "emails": batch,
        "webhook_url": "https://your-app.com/webhooks/email-validation",
    },
)
// Express: обработка webhook
const express = require("express");
const app = express();

app.post("/webhooks/email-validation", express.json(), async (req, res) => {
  const { task_id, status } = req.body;

  if (status === "completed") {
    const results = await fetchResults(task_id);
    await updateSubscribers(pool, results);
    console.log(`Webhook: task ${task_id} processed, ${results.length} emails`);
  }

  res.sendStatus(200);
});

В продакшене webhook и polling стоит комбинировать. Webhook - основной канал. Polling с большим интервалом (раз в минуту) - fallback на случай, если webhook не дошёл из-за сетевого сбоя или перезапуска вашего сервера.

Валидация на входе: не только по расписанию

Пайплайн по cron чистит существующую базу. Но новые невалидные адреса попадают туда прямо сейчас - через формы подписки, импорт из CRM, интеграции. Если проверять их только через 30 дней, вы месяц будете слать письма на мёртвые ящики.

Single-эндпоинт API решает эту задачу. Одинвызов - один адрес - ответ за секунду. Встраивается в форму подписки или обработчик импорта.

# Python: проверка адреса в момент подписки
def validate_on_signup(email: str) -> dict:
    """Validate a single email at signup. Returns verdict."""
    resp = session.post(
        f"{BASE}/api/v1/validate/single",
        headers={"x-api-key": API_KEY},
        json={"email": email},
        timeout=5,
    )
    resp.raise_for_status()
    data = resp.json()
    return {
        "valid": data["result"] != "invalid",
        "reason": data.get("reason"),
        "risk_score": data.get("risk_score", 0),
    }
// Node.js: middleware для проверки email при регистрации
async function validateEmail(email) {
  const resp = await fetch(`${BASE}/api/v1/validate/single`, {
    method: "POST",
    headers: {
      "x-api-key": API_KEY,
      "Content-Type": "application/json",
    },
    body: JSON.stringify({ email }),
  });
  const data = await resp.json();
  return {
    valid: data.result !== "invalid",
    reason: data.reason ?? null,
    riskScore: data.risk_score ?? 0,
  };
}

// Usage in Express route:
app.post("/subscribe", async (req, res) => {
  const { email } = req.body;
  const check = await validateEmail(email);

  if (!check.valid) {
    return res.status(400).json({ error: "Invalid email address" });
  }
  // ... save subscriber
});

Два уровня вместе - проверка на входе и периодическая перевалидация - закрывают задачу гигиены полностью. На входе не пропускаете мусор. По cron отлавливаете адреса, которые стали невалидными после подписки.

Типичные ошибки

Перевалидация всей базы каждый раз. Если адрес проверялся неделю назад и оказался good, тратить на него кредит бессмысленно. Поле last_validated_at решает проблему - проверяйте только stale-адреса.

Удаление risky-адресов без разбора. Catch-all домены попадают в risky, но на них могут сидеть реальные подписчики. Переводите risky в отдельный сегмент, не удаляйте сразу.

Отсутствие чекпоинтов. Скрипт отправил 5 батчей, упал на шестом. При перезапуске отправляет всё заново - двойной расход кредитов. Сохраняйте task_id в файл или базу после каждой отправки.

Игнорирование rate limits. Десять POST-запросов в секунду - и API начнёт отвечать 429. Пауза между батчами и retry с backoff решают проблему до её появления.

Скрипт без уведомлений. Cron-задача, которая падает молча, создаёт ложное ощущение работающей системы. Добавьте алерт при любом ненулевом exit code.

Автоматическая гигиена - это не одноразовая настройка. Это конвейер, который должен работать каждый день, на любом объёме, при любых сбоях. Разница между скриптом и конвейером - обработка ошибок, чекпоинты и мониторинг.

Попробуйте автоматизировать гигиену своей базы с uChecker API - бесплатные кредиты при регистрации, API-ключ за полминуты, первый батч можно отправить из терминала за пять минут.

автоматизация email гигиеныemail validation APIочистка базы emailPython email validationNode.js email verificationcron валидацияuchecker api