Коннекторы и приём данных: SaaS API, файлы, вебхуки и дрейф схемы
В первой статье трека мы разобрали приём данных из базы: инкремент
по водяному знаку, нахлёст, CDC через журнал транзакций, идемпотентный MERGE. Это честно описывает
источник, который принадлежит вам: у него есть схема, транзакции, монотонное время и админ, которому
можно написать в личку.
Проблема в том, что в реальной платформе таких источников меньшинство. Типичный список подключений зрелой компании — это 5–10 своих баз и 40–80 всего остального: биллинг, CRM, рекламные кабинеты, почтовая рассылка, поддержка, HR-система, антифрод-провайдер, партнёрские выгрузки по SFTP, вебхуки платёжного шлюза, выгрузка из 1С в CSV раз в сутки. У этой части нет транзакций, нет журнала изменений, нет гарантии, что вчерашние числа сегодня те же, а половина документации написана в 2019 году и с тех пор не обновлялась.
Слой приёма из внешних систем — отдельная инженерная дисциплина со своим набором отказов. Именно здесь возникает большинство инцидентов «данные приехали не полностью», и именно здесь их труднее всего заметить: внешний источник почти никогда не сообщает об ошибке — он молча отдаёт меньше строк.
1. Пять классов источников и их физика
Прежде чем писать код, полезно понять, к какому классу относится источник: от этого зависят и гарантии, и способ инкремента, и то, что вообще придётся мониторить.
| Класс | Инкремент | Гарантия полноты | Главный отказ | Задержка |
|---|---|---|---|---|
| Реляционная БД | водяной знак / CDC | высокая (транзакции) | лаг слота, DDL в источнике | секунды–минуты |
| SaaS API | updated_since или курсор |
средняя, зависит от вендора | лимиты, пагинация, ретроспективные правки | минуты–часы |
| Файловый дроп | новый файл = новая порция | низкая (файл может не приехать) | частичная выкладка, дубли, кодировки | часы–сутки |
| Вебхуки | поток событий | низкая (доставка не гарантирована навсегда) | потеря при недоступности приёмника | секунды |
| Ручной ввод | версия документа | никакой | человек | дни |
Из таблицы следует главное правило проектирования приёма: чем ниже гарантия полноты, тем обязательнее независимая сверка. Для CDC достаточно мониторить лаг. Для вебхуков сверка с периодическим полным опросом API — не «хорошая практика», а часть механизма: без неё вы гарантированно потеряете события в первый же инцидент сети.
2. Анатомия коннектора
Любой коннектор, независимо от источника, состоит из одних и тех же шагов. Если в вашем скрипте нет какого-то из них — это не упрощение, это отложенный инцидент.
курсор предыдущего запуска] B --> C{Состояние есть?} C -->|нет| D[Полный снимок:
исторический период] C -->|да| E[План инкремента:
окно от курсора − нахлёст] D --> F[Выборка страницами
с ретраями и лимитами] E --> F F --> G[Приземление в RAW
как есть, без трансформаций] G --> H{Батч дописан
полностью?} H -->|нет| I[Откат батча,
курсор НЕ двигаем] H -->|да| J[Сдвиг курсора
атомарно с фиксацией батча] J --> K[Метрики: строк, вызовов,
ошибок, длительность] K --> L{Расхождение
со сверкой?} L -->|да| M[Алерт и пересинхронизация] L -->|нет| N[Успех] I --> M
Ключевых инвариантов три, и все три нарушают в первой же самописной интеграции:
- Курсор двигается только после подтверждённой записи данных. Не «прочитали — записали курсор — пишем данные»: падение между шагами создаёт дыру, которую никто не заметит месяцами. Порядок строго обратный, и в идеале курсор лежит в той же транзакции (или в том же коммите табличного формата), что и данные.
- Состояние коннектора — это данные, а не переменная процесса. Курсор хранится в таблице или в бэкенде оркестратора, а не в файле на диске воркера, который завтра пересоздадут.
- Запуск идемпотентен. Повтор того же интервала не должен ни удваивать строки, ни терять их;
способ — тот же
MERGEпо ключу или перезапись партиции из главы 01.
3. Пагинация: самая тихая потеря данных
Внешний API почти никогда не отдаёт всё сразу. Способов нарезки три, и они принципиально разного качества.
Offset-пагинация (?limit=100&offset=300) — самая распространённая и самая опасная. Она
предполагает, что между запросом страницы N и страницы N+1 набор данных не меняется. На живом
источнике это неправда.
Keyset-пагинация (курсор по ключу): WHERE id < :last_seen_id ORDER BY id DESC LIMIT 100.
Курсор — это позиция в данных, а не номер строки, поэтому вставки и удаления его не сдвигают.
Требование одно: сортировка по уникальному неизменяемому ключу. Если сортируете по времени, ключ
должен быть парой (created_at, id) — иначе строки с одинаковой меткой времени на границе страницы
теряются или дублируются.
Непрозрачный курсор вендора (next_page_token, has_more + starting_after) — лучший вариант,
когда он есть: вендор сам отвечает за консистентность. Два правила: токен нельзя парсить (он
меняется без предупреждения) и у него есть срок жизни — обход в 40 минут может упасть на середине с
400 invalid cursor, а значит, обход надо уметь начинать заново, а не «дочитывать».
Отдельный класс — экспортные задания: вы просите вендора подготовить выгрузку, получаете
job_id, опрашиваете статус, скачиваете файл. Так работают крупные рекламные кабинеты и часть
биллингов. Здесь коннектор — уже не «цикл по страницам», а маленький асинхронный протокол с
собственными таймаутами и состоянием.
4. Лимиты, ретраи и бюджет вызовов
Внешний API — разделяемый ресурс, и он защищается. Практический минимум, который обязан уметь коннектор:
- Различать классы ошибок.
429и5xx— повторяемые.4xx(кроме 429) — не повторяемые: бесконечный ретрай на401только сожжёт лимит и разбудит дежурного вендора. - Уважать
Retry-After. Если сервер сказал, через сколько вернуться, экспоненциальная формула не нужна — нужен указанный интервал. - Экспоненциальная задержка с джиттером. Без джиттера сто параллельных задач синхронно повторяют запрос и создают ту же волну, из-за которой их и отсекли. Классический разбор — Exponential Backoff and Jitter в блоге AWS: «full jitter» (случайная пауза от нуля до текущей границы) выигрывает у чистой экспоненты и по времени завершения, и по нагрузке на сервер.
- Ограничивать собственную скорость. Токен-бакет на стороне коннектора дешевле, чем ловить
429: вы платите не только временем, но и репутацией у вендора, а иногда и деньгами за вызов. - Считать бюджет. Если лимит 10 000 вызовов в сутки, страница — 100 записей, а объект меняется у 200 000 записей в день, полный обход невозможен в принципе. Это надо обнаружить на этапе проектирования, а не на третьем месяце эксплуатации.
"""Ядро коннектора к внешнему REST API.
Свойства: keyset-курсор, честная обработка 429/5xx, ограничение своей скорости,
чекпоинт после подтверждённой записи, метрики на каждый запуск.
"""
import random
import time
from dataclasses import dataclass, field
from typing import Any, Iterator
import requests
RETRYABLE = {429, 500, 502, 503, 504}
@dataclass
class Budget:
"""Простейший токен-бакет: не более rate вызовов в секунду."""
rate: float
_next_slot: float = field(default_factory=time.monotonic)
def acquire(self) -> None:
now = time.monotonic()
if self._next_slot > now:
time.sleep(self._next_slot - now)
self._next_slot = max(now, self._next_slot) + 1.0 / self.rate
def fetch_page(session: requests.Session, url: str, params: dict, budget: Budget,
max_attempts: int = 6) -> dict[str, Any]:
"""Один запрос страницы с ретраями. Возвращает разобранный JSON."""
for attempt in range(max_attempts):
budget.acquire()
resp = session.get(url, params=params, timeout=30)
if resp.status_code == 200:
return resp.json()
if resp.status_code not in RETRYABLE:
# 401/403/404 повторять бессмысленно: это ошибка конфигурации, а не сети.
resp.raise_for_status()
# Сервер знает лучше нас, когда возвращаться.
retry_after = resp.headers.get("Retry-After")
if retry_after is not None:
delay = float(retry_after)
else:
# full jitter: пауза равномерна на [0, 2^attempt), верхняя граница — минута.
delay = random.uniform(0, min(60.0, 2 ** attempt))
time.sleep(delay)
raise RuntimeError(f"не удалось получить страницу за {max_attempts} попыток: {url}")
def iter_records(session: requests.Session, url: str, cursor: str | None,
budget: Budget) -> Iterator[tuple[list[dict], str]]:
"""Генератор батчей: отдаёт (записи страницы, курсор ПОСЛЕ этой страницы).
Курсор отдаётся вместе с данными, чтобы вызывающий код мог сохранить его
строго после успешной записи батча, а не до неё.
"""
params = {"limit": 500, "order": "updated_at.asc"}
while True:
if cursor:
params["after"] = cursor
page = fetch_page(session, url, params, budget)
records = page.get("data", [])
if not records:
return
cursor = page.get("next_cursor") or records[-1]["id"]
yield records, cursor
if not page.get("has_more", True):
return
Асимптотика тут скучная и потому важная: время работы — O(N / page_size) сетевых запросов, а
не O(N), и именно поэтому размер страницы — первая настройка, которую стоит поднять до
максимально разрешённой вендором. Память — O(page_size): полный список записей в оперативке
не собирают никогда, батч пишется в хранилище и забывается.
5. Инкремент, которого нет: ретроспективные правки
С базой данных всё честно: строка изменилась — updated_at вырос. Внешние системы регулярно нарушают
это допущение, и нарушают тихо:
- Правки задним числом. Возврат по заказу двухмесячной давности меняет запись за март, но
updated_atу неё может остаться мартовским (особенно если правка сделана импортом). - Пересчёт агрегатов. Рекламные кабинеты и системы аналитики уточняют вчерашние цифры ещё 3–14 дней: дедупликация показов, фильтрация фрода, конверсии с отложенной атрибуцией. Значение «показы за 1 июля», забранное 2 июля и 20 июля, будет разным — и второе правильнее.
- Отсутствие поля времени изменения. Есть только
created_at, и апдейты невидимы. - Мягкие удаления без времени. Запись «исчезла» из выдачи, и обнаружить это можно только сверкой множеств ключей.
Практическая стратегия — скользящее окно пересинхронизации: каждый запуск забирает не только
новое, но и последние N дней целиком, а приземление идемпотентно (MERGE по ключу или перезапись
партиции). Размер окна берётся не из головы, а из измерения: раз в неделю сравнивают значения,
забранные с задержкой в 1, 3, 7, 14 и 30 дней, и смотрят, когда цифры перестают меняться.
-- Проверка «когда данные стабилизируются»: сравниваем один и тот же день,
-- забранный в разные даты. Пока разница не нулевая, окно пересинхронизации мало.
WITH snapshots AS (
SELECT report_date,
_extracted_date,
SUM(impressions) AS impressions
FROM raw.ads_daily_report
GROUP BY report_date, _extracted_date
)
SELECT report_date,
DATEDIFF('day', report_date, _extracted_date) AS age_days,
impressions,
impressions - FIRST_VALUE(impressions) OVER (
PARTITION BY report_date ORDER BY _extracted_date
) AS drift_from_first
FROM snapshots
WHERE report_date >= CURRENT_DATE - 45
ORDER BY report_date, age_days;
Правило: окно пересинхронизации должно быть больше периода стабилизации данных, а витрины, построенные поверх, обязаны уметь пересчитываться за этот период. Дашборд, который показывает «выручку за вчера» как окончательную, когда источник уточняет её ещё неделю, врёт по построению — см. разговор про доверие к числам в главе о качестве данных.
6. Жизненный цикл подключения
У коннектора больше состояний, чем «работает / не работает», и хороший приёмный слой различает их явно.
проверен доступ Настроен --> Дискаверинг: читаем список объектов
и их схемы Дискаверинг --> ПолныйСнимок: первая загрузка ПолныйСнимок --> Инкремент: снимок завершён,
курсор зафиксирован ПолныйСнимок --> Деградация: лимит исчерпан
на середине снимка Инкремент --> Инкремент: обычный запуск Инкремент --> Деградация: 429 / 5xx
дольше порога Деградация --> Инкремент: источник ожил,
догоняем отставание Инкремент --> Пересинхронизация: сверка нашла
расхождение Пересинхронизация --> Инкремент: расхождение закрыто Инкремент --> Сломан: схема изменилась
несовместимо Сломан --> Дискаверинг: схема принята,
модель обновлена Инкремент --> Отключён: источник выведен
из эксплуатации Отключён --> [*] note right of Деградация Деградация — не отказ. Данные не теряются, растёт лаг. Алерт нужен другой: не «упало», а «отстаём N часов». end note
Различие между «сломан» и «деградация» стоит дороже, чем кажется. Если оба состояния поднимают один и тот же алерт, дежурный за месяц перестаёт их читать; а именно «деградация, которую никто не разобрал» превращается в «инкремент отстал на четверо суток, и вендор уже удалил историю».
7. Файловые источники: скучно и опасно
Партнёр выкладывает orders_20260716.csv на SFTP. Что может пойти не так — почти всё.
| Дефект | Как проявляется | Защита |
|---|---|---|
| Частичная выкладка | файл виден, пока пишется; прочитали половину | ждать маркер .done или загрузку во временное имя + переименование |
| Дубль файла | тот же контент под новым именем | хэш содержимого в журнале приёмов, отказ от повторного приёма |
| Пропуск дня | файла просто нет | проверка ожидаемого расписания, алерт на отсутствие |
| Сдвиг колонок | добавили колонку в середину | приём по имени колонки из заголовка, а не по позиции |
| Кодировка и разделители | windows-1251, ;, десятичная запятая |
явно фиксировать диалект в конфиге источника, а не угадывать |
| Экранирование | перенос строки внутри поля рвёт CSV | предпочитать Parquet/JSONL; для CSV — строгий парсер, не split(',') |
| Тихая перезаписка | партнёр выложил тот же день заново с другими числами | версионировать приёмы, хранить все версии в RAW |
Дисциплина, которая закрывает большинство пунктов: журнал приёмов как таблица. Каждый принятый файл — строка с именем, размером, хэшем, временем приёма, числом строк и статусом. Дальше приём становится идемпотентным по построению: файл с известным хэшем не обрабатывается повторно, а недостающий день виден запросом, а не глазами.
-- Пропуски в ежедневной выкладке видны одним запросом,
-- если журнал приёмов существует. Без журнала — не видны никак.
SELECT d.expected_date
FROM UNNEST(SEQUENCE(CURRENT_DATE - 30, CURRENT_DATE - 1, INTERVAL '1' DAY)) AS d(expected_date)
LEFT JOIN ingest.file_log f
ON f.source = 'partner_sftp_orders'
AND f.business_date = d.expected_date
AND f.status = 'loaded'
WHERE f.business_date IS NULL
ORDER BY d.expected_date;
8. Вебхуки: событие приехало один раз и больше не приедет
Вебхук — это HTTP-запрос от источника к вам. Отсюда все свойства: доставка at-least-once, но не навсегда; порядок не гарантирован; повторы возможны; если ваш эндпоинт лежал полчаса, вендор, скорее всего, отступится после нескольких попыток.
не создаёт вторую строку loop раз в сутки W->>A: полный опрос за последние N дней A-->>W: список объектов W->>R: сверка: чего нет в RAW — дозабрать end
Четыре правила, которые стоит записать в шаблон приёмника:
- Отвечать
200как можно раньше, до бизнес-обработки. Вендор считает медленный ответ ошибкой и начинает повторять — вы получаете лавину дублей ровно в момент пиковой нагрузки. - Хранить сырое тело. Разбор может оказаться неверным; повторить разбор можно только по оригиналу. Это тот же принцип неприкосновенного RAW, что и в обзорной статье.
- Дедуплицировать по идентификатору события, а не по содержимому: одинаковые по полям события бывают разными фактами.
- Не считать вебхуки источником правды. Периодический опрос API — обязательная страховка; расхождение между «пришло вебхуком» и «есть в API» — метрика, за которой следят.
Про сами гарантии доставки и способы жить с ними подробно написано в «Идемпотентность и семантика доставки» и в статье про обмен сообщениями.
9. Дрейф схемы: что делать, когда источник изменился
Схема внешнего источника меняется без вашего согласия, и это нормальная жизнь. Ненормально — узнавать об этом от аналитика.
| Изменение | Опасность | Разумная реакция по умолчанию |
|---|---|---|
| Добавлена колонка | низкая | принять автоматически, добавить в RAW, уведомить владельца |
| Удалена колонка | высокая | не удалять в приёмнике, заполнять NULL, алерт: витрины могут молча обнулиться |
| Переименована колонка | высокая | равно «удалена + добавлена»; нужен человек |
Расширен тип (int → bigint) |
низкая | принять |
Сужен тип (bigint → int, строка короче) |
высокая | карантин: часть значений не влезет |
| Сменился смысл значения | критическая | тестами не ловится; ловится контрактом и общением |
| Новое значение перечисления | средняя | принять, но CASE в витринах должен иметь ветку «прочее» |
Три опоры, на которых держится устойчивость к дрейфу:
Реестр схем. Для событийных источников — Schema Registry с режимом совместимости. BACKWARD
(режим по умолчанию в Confluent) означает: новая версия схемы умеет читать данные, записанные старой.
Это позволяет обновлять потребителей раньше производителей. FORWARD — наоборот. FULL — оба
одновременно, дороже всего в согласовании. Формулировки и таблица допустимых изменений — в
документации Confluent.
RAW принимает всё. В сыром слое схема максимально терпима: полезная нагрузка хранится как JSON
плюс несколько служебных колонок (_source, _extracted_at, _batch_id, _payload_hash).
Тогда добавление поля в источнике физически не может уронить приём. Строгая типизация начинается на
слое staging, где падение — уже осознанный сигнал, а не потеря данных.
Проверка схемы отдельно от загрузки. Коннектор сравнивает обнаруженную схему с зафиксированной и принимает решение: совместимо — принять и записать в журнал изменений; несовместимо — остановиться с внятной ошибкой. Это ровно тот механизм, который в главе про качество назван контрактом данных, только со стороны приёмника.
# contracts/sources/billing_invoices.v3.yaml — контракт на ВНЕШНИЙ источник.
# Отличие от контракта на свой сервис: обещания даём мы сами, а не вендор.
source: billing_saas
object: invoices
version: 3
extraction:
mode: incremental
cursor_field: updated_at
cursor_type: timestamp
resync_window_days: 14 # вендор уточняет суммы до 10 дней, берём с запасом
page_size: 500
rate_limit_rps: 4
schema:
required: [id, customer_id, amount_cents, currency, status, updated_at]
types:
id: string
amount_cents: integer # деньги только в минорных единицах, не float
currency: string
status: string
enum:
status: [draft, open, paid, void, uncollectible]
on_drift:
new_column: accept # принять и уведомить
removed_column: alert # не ломать приём, но разбудить владельца
type_narrowing: quarantine # в карантин, разбирает человек
new_enum_value: accept # витрины обязаны иметь ветку «прочее»
reconciliation:
method: count_and_sum
keys: [report_date]
tolerance_pct: 0.1
owner: team-finance-data
10. Сверка с источником: единственный способ узнать правду
Все предыдущие механизмы уменьшают вероятность потери. Обнаружить потерю может только независимая проверка — сравнение вашего приёмника с источником по агрегату, который источник умеет посчитать сам.
Три уровня, по возрастанию стоимости:
- Счётчики. «Сколько объектов создано за 1 июля» из API против
COUNT(*)в RAW. Дёшево, ловит грубые пропуски. - Суммы и контрольные агрегаты. Сумма платежей за день, число уникальных клиентов. Ловит дубликаты и частичную загрузку.
- Выборочная сверка ключей. Случайные 1000 идентификаторов из источника — есть ли они у нас, и совпадают ли поля. Ловит расхождение по содержимому, которое агрегаты маскируют.
-- Ежедневная сверка: расхождение больше допуска — инцидент, а не строчка в логе.
WITH ours AS (
SELECT DATE(created_at) AS d, COUNT(*) AS cnt, SUM(amount_cents) AS amt
FROM raw.billing_invoices
WHERE created_at >= CURRENT_DATE - 7
GROUP BY 1
),
theirs AS ( -- заполняется отдельной задачей из отчётного эндпоинта вендора
SELECT report_date AS d, invoice_count AS cnt, invoice_amount_cents AS amt
FROM raw.billing_vendor_totals
WHERE report_date >= CURRENT_DATE - 7
)
SELECT t.d,
t.cnt AS src_cnt, o.cnt AS our_cnt,
t.amt AS src_amt, o.amt AS our_amt,
ROUND(100.0 * (o.cnt - t.cnt) / NULLIF(t.cnt, 0), 3) AS cnt_diff_pct,
ROUND(100.0 * (o.amt - t.amt) / NULLIF(t.amt, 0), 3) AS amt_diff_pct
FROM theirs t
LEFT JOIN ours o USING (d)
WHERE o.cnt IS NULL
OR ABS(o.cnt - t.cnt) > 0
OR ABS(o.amt - t.amt) > t.amt * 0.001 -- допуск 0,1 %
ORDER BY t.d DESC;
Допуск — важная часть конструкции. Нулевой допуск для денег обязателен, а для событийной телеметрии он породит ежедневный ложный алерт: часть событий всегда приходит с задержкой. Допуск должен быть явным числом в контракте, а не негласной привычкой дежурного «ну, полпроцента — это нормально».
11. Своё против готового
Соблазн написать «простой скрипт на 200 строк» велик, потому что первая версия действительно занимает день. Стоимость коннектора — не в первой версии, а в поддержке: смена API вендора, новые объекты, лимиты, ретроспективные правки, дежурство.
| Вариант | Что даёт | Чем платите | Когда разумно |
|---|---|---|---|
| Управляемый сервис (Fivetran, Stitch) | сотни готовых источников, чужое дежурство | цена растёт с объёмом строк, чужие решения о схеме | много типовых SaaS-источников, мало людей |
| Открытая платформа (Airbyte, Meltano) | те же коннекторы, свой контур | вы эксплуатируете платформу и чините коннекторы | есть требования к контуру данных, есть команда |
| Библиотека (dlt, Singer-таргеты) | код в вашем репозитории, обычный CI | пишете инкремент и схему сами | 5–20 нетиповых источников |
| Полностью своё | точный контроль, минимум зависимостей | поддержка навсегда | источник уникален или критичен настолько, что чужому не доверяют |
| CDC-платформа (Debezium) | журнал БД в реальном времени | эксплуатация Kafka Connect, риск для источника | свои базы, нужна низкая задержка |
Практический критерий, который экономит месяцы: сколько стоит день простоя этого источника. Если ответ «никто не заметит» — берите готовое и не думайте. Если «встанет расчёт выплат» — коннектор должен быть вашим, покрытым тестами и сверкой, независимо от того, есть ли готовый.
Отдельно про доступы: ключи внешних API — такие же секреты, как пароли к продакшн-базе, и живут они в хранилище секретов, а не в переменных DAG’а. Разбор практик — в статье «Управление секретами». Если вы принимаете персональные данные — до первого запуска стоит прочитать «Приватность и соответствие требованиям»: приём — это тот момент, когда лишние поля попадают в платформу и остаются в ней навсегда.
12. Наблюдаемость приёма
Минимальный набор метрик на каждое подключение, без которого эксплуатация превращается в гадание:
| Метрика | Зачем | Типичный алерт |
|---|---|---|
Лаг источника (now − max(cursor)) |
главный индикатор здоровья | больше двух интервалов расписания |
| Строк за запуск | ловит тихое обнуление | падение более чем на 50 % от медианы за 14 дней |
Вызовов и доля 429 |
приближение к лимиту | доля 429 выше 5 % |
| Длительность запуска | предсказывает отказ по таймауту | рост в 2 раза к скользящей медиане |
| Расхождение со сверкой | единственная метрика полноты | вне допуска контракта |
| Стоимость (вызовы, трафик, строки) | приём тоже стоит денег | рост без роста объёма данных |
Обратите внимание на асимметрию: «пайплайн упал» видно всем и сразу, «пайплайн отработал успешно и привёз 0 строк» не видно никому — формально всё зелёное. Поэтому метрика объёма и сверка важнее метрики успешности запуска. Как из этих метрик строить SLO на свежесть данных, разобрано в главе про качество; общая теория SLI и SLO — в треке SRE.
13. Типичные ошибки
- Курсор сдвигается до записи данных. Классическая дыра: падение между шагами теряет батч навсегда, и никакой ретрай его не вернёт.
- Offset-пагинация на изменяющемся источнике. Тихая потеря строк, которую находят через полгода при сверке с бухгалтерией.
- Отсутствие окна пересинхронизации. Числа в источнике уточняются, у вас остаются первые версии — расхождение растёт монотонно.
- Бесконечный ретрай на
401. Сожжённый лимит, письмо от вендора, отключённый ключ. - Парсинг CSV через
split(','). Работает ровно до первой запятой внутри кавычек — обычно в названии компании. - Вебхук как единственный канал. Полчаса недоступности приёмника = навсегда потерянные события.
- Схема-жёсткость в RAW. Источник добавил поле — приём упал; источник удалил поле — витрина молча обнулилась.
- Секреты в коде DAG’а. Ротация ключа превращается в поиск по репозиторию.
- Нет журнала приёмов. Невозможно ответить на вопрос «этот файл мы уже грузили?» иначе как сравнением данных.
- Приём с трансформацией. «Заодно приведём типы и отфильтруем мусор» — и через месяц никто не может доказать, что именно прислал источник. Приём отделён от трансформации не из эстетики, а ради разрешимости споров.
14. Как это выглядит в проде
Зрелый приёмный слой выглядит скучно, и это его достоинство. Каждое подключение описано конфигом в репозитории (объекты, курсор, окно пересинхронизации, лимиты, владелец, допуски сверки), а не живёт в UI. Секреты приезжают из хранилища секретов. Все запуски пишут в общий журнал: строк, вызовов, длительность, ошибки. Сверка идёт по расписанию отдельным DAG’ом и заводит инцидент сама. Новый источник подключается по чек-листу и попадает в общий дашборд лагов автоматически — никто не «не забыл добавить мониторинг», потому что мониторинг следует из конфига.
Признак незрелого слоя тоже узнаваем: набор скриптов с разными подходами к курсорам, ключи в переменных окружения воркера, «мониторинг» в виде почтовых уведомлений о падениях и ежеквартальное обнаружение того, что один из источников не обновлялся с мая.
15. Мини-итог
- Внешние источники не дают транзакций и монотонного времени — гарантии приходится строить самому.
- Коннектор — это цикл «состояние → план → выборка → приземление → сдвиг курсора → сверка», и порядок шагов важнее, чем язык, на котором он написан.
- Offset-пагинация ломается на живых данных; правильный курсор — позиция по уникальному неизменяемому ключу.
- Ретраи разделяют повторяемые и неповторяемые ошибки, уважают
Retry-Afterи используют джиттер. - Данные в источнике уточняются задним числом: окно пересинхронизации подбирается измерением, а не интуицией.
- Файлы требуют журнала приёмов, вебхуки — дедупликации и страховочного опроса API.
- RAW принимает любую схему; строгость начинается на
staging. - Единственный способ узнать о потере — независимая сверка с источником по счётчикам и суммам.
- Готовый коннектор берут по умолчанию; свой пишут, когда цена простоя источника это оправдывает.
Источники
- Debezium Documentation — CDC-коннекторы, снапшоты, обработка DDL: https://debezium.io/documentation/reference/stable/
- Confluent Schema Registry — режимы совместимости схем и правила эволюции: https://docs.confluent.io/platform/current/schema-registry/fundamentals/schema-evolution.html
- Marc Brooker. Exponential Backoff and Jitter, AWS Architecture Blog — почему джиттер обязателен: https://aws.amazon.com/blogs/architecture/exponential-backoff-and-jitter/
- Stripe API Reference — эталон курсорной пагинации, идемпотентных ключей и подписи вебхуков: https://docs.stripe.com/api/pagination
- Google Cloud. API Design Guide: List Pagination — почему токен непрозрачен: https://cloud.google.com/apis/design/design_patterns#list_pagination
- dlt (data load tool) — библиотека приёма с состоянием, инкрементом и эволюцией схемы: https://dlthub.com/docs/intro
- Airbyte Protocol — модель дискаверинга, синка и состояния коннектора: https://docs.airbyte.com/understanding-airbyte/airbyte-protocol
- RFC 9110, §10.2.3
Retry-After— что именно обещает заголовок: https://www.rfc-editor.org/rfc/rfc9110#field.retry-after - Joe Reis, Matt Housley. Fundamentals of Data Engineering, O’Reilly, 2022 — глава Ingestion как систематический разбор источников.
Что дальше
Данные из внешнего мира приезжают в платформу, и с этого момента всё, что с ними происходит, — ваш код: SQL-модели, DAG’и, конфиги коннекторов. Код надо менять, а менять его в системе, где каждая ошибка обнаруживается на проде через сутки и стоит пересчёта витрин, страшно. Следующая статья — про то, как сделать изменения дешёвыми: тесты трансформаций, среды разработки поверх хранилища, сравнение результата до и после и безопасный релиз моделей.
Дальше: Тестирование и CI/CD дата-пайплайнов.