Автоматизация 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_idsNode.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-ключ за полминуты, первый батч можно отправить из терминала за пять минут.
