Data Engineering и ETL Коннекторы и приём данных: SaaS API, файлы, вебхуки и дрейф схемы
0%

Коннекторы и приём данных: SaaS API, файлы, вебхуки и дрейф схемы

Коннекторы и приём данных: SaaS API, файлы, вебхуки и дрейф схемы

В первой статье трека мы разобрали приём данных из базы: инкремент по водяному знаку, нахлёст, CDC через журнал транзакций, идемпотентный MERGE. Это честно описывает источник, который принадлежит вам: у него есть схема, транзакции, монотонное время и админ, которому можно написать в личку.

Проблема в том, что в реальной платформе таких источников меньшинство. Типичный список подключений зрелой компании — это 5–10 своих баз и 40–80 всего остального: биллинг, CRM, рекламные кабинеты, почтовая рассылка, поддержка, HR-система, антифрод-провайдер, партнёрские выгрузки по SFTP, вебхуки платёжного шлюза, выгрузка из 1С в CSV раз в сутки. У этой части нет транзакций, нет журнала изменений, нет гарантии, что вчерашние числа сегодня те же, а половина документации написана в 2019 году и с тех пор не обновлялась.

Слой приёма из внешних систем — отдельная инженерная дисциплина со своим набором отказов. Именно здесь возникает большинство инцидентов «данные приехали не полностью», и именно здесь их труднее всего заметить: внешний источник почти никогда не сообщает об ошибке — он молча отдаёт меньше строк.


1. Пять классов источников и их физика

Прежде чем писать код, полезно понять, к какому классу относится источник: от этого зависят и гарантии, и способ инкремента, и то, что вообще придётся мониторить.

Класс Инкремент Гарантия полноты Главный отказ Задержка
Реляционная БД водяной знак / CDC высокая (транзакции) лаг слота, DDL в источнике секунды–минуты
SaaS API updated_since или курсор средняя, зависит от вендора лимиты, пагинация, ретроспективные правки минуты–часы
Файловый дроп новый файл = новая порция низкая (файл может не приехать) частичная выкладка, дубли, кодировки часы–сутки
Вебхуки поток событий низкая (доставка не гарантирована навсегда) потеря при недоступности приёмника секунды
Ручной ввод версия документа никакой человек дни

Из таблицы следует главное правило проектирования приёма: чем ниже гарантия полноты, тем обязательнее независимая сверка. Для CDC достаточно мониторить лаг. Для вебхуков сверка с периодическим полным опросом API — не «хорошая практика», а часть механизма: без неё вы гарантированно потеряете события в первый же инцидент сети.


2. Анатомия коннектора

Любой коннектор, независимо от источника, состоит из одних и тех же шагов. Если в вашем скрипте нет какого-то из них — это не упрощение, это отложенный инцидент.

Ключевых инвариантов три, и все три нарушают в первой же самописной интеграции:

  1. Курсор двигается только после подтверждённой записи данных. Не «прочитали — записали курсор — пишем данные»: падение между шагами создаёт дыру, которую никто не заметит месяцами. Порядок строго обратный, и в идеале курсор лежит в той же транзакции (или в том же коммите табличного формата), что и данные.
  2. Состояние коннектора — это данные, а не переменная процесса. Курсор хранится в таблице или в бэкенде оркестратора, а не в файле на диске воркера, который завтра пересоздадут.
  3. Запуск идемпотентен. Повтор того же интервала не должен ни удваивать строки, ни терять их; способ — тот же MERGE по ключу или перезапись партиции из главы 01.

3. Пагинация: самая тихая потеря данных

Внешний API почти никогда не отдаёт всё сразу. Способов нарезки три, и они принципиально разного качества.

Offset-пагинация (?limit=100&offset=300) — самая распространённая и самая опасная. Она предполагает, что между запросом страницы N и страницы N+1 набор данных не меняется. На живом источнике это неправда.

Как OFFSET-пагинация теряет строки на изменяющемся источнике

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. Жизненный цикл подключения

У коннектора больше состояний, чем «работает / не работает», и хороший приёмный слой различает их явно.

Различие между «сломан» и «деградация» стоит дороже, чем кажется. Если оба состояния поднимают один и тот же алерт, дежурный за месяц перестаёт их читать; а именно «деградация, которую никто не разобрал» превращается в «инкремент отстал на четверо суток, и вендор уже удалил историю».


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, но не навсегда; порядок не гарантирован; повторы возможны; если ваш эндпоинт лежал полчаса, вендор, скорее всего, отступится после нескольких попыток.

Четыре правила, которые стоит записать в шаблон приёмника:

  1. Отвечать 200 как можно раньше, до бизнес-обработки. Вендор считает медленный ответ ошибкой и начинает повторять — вы получаете лавину дублей ровно в момент пиковой нагрузки.
  2. Хранить сырое тело. Разбор может оказаться неверным; повторить разбор можно только по оригиналу. Это тот же принцип неприкосновенного RAW, что и в обзорной статье.
  3. Дедуплицировать по идентификатору события, а не по содержимому: одинаковые по полям события бывают разными фактами.
  4. Не считать вебхуки источником правды. Периодический опрос API — обязательная страховка; расхождение между «пришло вебхуком» и «есть в API» — метрика, за которой следят.

Про сами гарантии доставки и способы жить с ними подробно написано в «Идемпотентность и семантика доставки» и в статье про обмен сообщениями.


9. Дрейф схемы: что делать, когда источник изменился

Схема внешнего источника меняется без вашего согласия, и это нормальная жизнь. Ненормально — узнавать об этом от аналитика.

Изменение Опасность Разумная реакция по умолчанию
Добавлена колонка низкая принять автоматически, добавить в RAW, уведомить владельца
Удалена колонка высокая не удалять в приёмнике, заполнять NULL, алерт: витрины могут молча обнулиться
Переименована колонка высокая равно «удалена + добавлена»; нужен человек
Расширен тип (intbigint) низкая принять
Сужен тип (bigintint, строка короче) высокая карантин: часть значений не влезет
Сменился смысл значения критическая тестами не ловится; ловится контрактом и общением
Новое значение перечисления средняя принять, но 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. Счётчики. «Сколько объектов создано за 1 июля» из API против COUNT(*) в RAW. Дёшево, ловит грубые пропуски.
  2. Суммы и контрольные агрегаты. Сумма платежей за день, число уникальных клиентов. Ловит дубликаты и частичную загрузку.
  3. Выборочная сверка ключей. Случайные 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. Типичные ошибки

  1. Курсор сдвигается до записи данных. Классическая дыра: падение между шагами теряет батч навсегда, и никакой ретрай его не вернёт.
  2. Offset-пагинация на изменяющемся источнике. Тихая потеря строк, которую находят через полгода при сверке с бухгалтерией.
  3. Отсутствие окна пересинхронизации. Числа в источнике уточняются, у вас остаются первые версии — расхождение растёт монотонно.
  4. Бесконечный ретрай на 401. Сожжённый лимит, письмо от вендора, отключённый ключ.
  5. Парсинг CSV через split(','). Работает ровно до первой запятой внутри кавычек — обычно в названии компании.
  6. Вебхук как единственный канал. Полчаса недоступности приёмника = навсегда потерянные события.
  7. Схема-жёсткость в RAW. Источник добавил поле — приём упал; источник удалил поле — витрина молча обнулилась.
  8. Секреты в коде DAG’а. Ротация ключа превращается в поиск по репозиторию.
  9. Нет журнала приёмов. Невозможно ответить на вопрос «этот файл мы уже грузили?» иначе как сравнением данных.
  10. Приём с трансформацией. «Заодно приведём типы и отфильтруем мусор» — и через месяц никто не может доказать, что именно прислал источник. Приём отделён от трансформации не из эстетики, а ради разрешимости споров.

14. Как это выглядит в проде

Зрелый приёмный слой выглядит скучно, и это его достоинство. Каждое подключение описано конфигом в репозитории (объекты, курсор, окно пересинхронизации, лимиты, владелец, допуски сверки), а не живёт в UI. Секреты приезжают из хранилища секретов. Все запуски пишут в общий журнал: строк, вызовов, длительность, ошибки. Сверка идёт по расписанию отдельным DAG’ом и заводит инцидент сама. Новый источник подключается по чек-листу и попадает в общий дашборд лагов автоматически — никто не «не забыл добавить мониторинг», потому что мониторинг следует из конфига.

Признак незрелого слоя тоже узнаваем: набор скриптов с разными подходами к курсорам, ключи в переменных окружения воркера, «мониторинг» в виде почтовых уведомлений о падениях и ежеквартальное обнаружение того, что один из источников не обновлялся с мая.


15. Мини-итог

  • Внешние источники не дают транзакций и монотонного времени — гарантии приходится строить самому.
  • Коннектор — это цикл «состояние → план → выборка → приземление → сдвиг курсора → сверка», и порядок шагов важнее, чем язык, на котором он написан.
  • Offset-пагинация ломается на живых данных; правильный курсор — позиция по уникальному неизменяемому ключу.
  • Ретраи разделяют повторяемые и неповторяемые ошибки, уважают Retry-After и используют джиттер.
  • Данные в источнике уточняются задним числом: окно пересинхронизации подбирается измерением, а не интуицией.
  • Файлы требуют журнала приёмов, вебхуки — дедупликации и страховочного опроса API.
  • RAW принимает любую схему; строгость начинается на staging.
  • Единственный способ узнать о потере — независимая сверка с источником по счётчикам и суммам.
  • Готовый коннектор берут по умолчанию; свой пишут, когда цена простоя источника это оправдывает.

Источники

Что дальше

Данные из внешнего мира приезжают в платформу, и с этого момента всё, что с ними происходит, — ваш код: SQL-модели, DAG’и, конфиги коннекторов. Код надо менять, а менять его в системе, где каждая ошибка обнаруживается на проде через сутки и стоит пересчёта витрин, страшно. Следующая статья — про то, как сделать изменения дешёвыми: тесты трансформаций, среды разработки поверх хранилища, сравнение результата до и после и безопасный релиз моделей.

Дальше: Тестирование и CI/CD дата-пайплайнов.

Нашли неточность? Выделите фрагмент текста — рядом появится жучок.

Нужен разбор именно вашей ситуации?

Статья описывает общий случай. Если у вас частный — можно разобрать его отдельно, платно. А если не хватает целого материала, предложите тему: её оплачивают вскладчину, и она выходит открытой для всех.

Доска запросов