Зачем коннекторы и чем они отличаются от загрузчиков
Загрузчик из предыдущего урока работает с локальным файлом: открыл, прочитал, закрыл. Коннектор решает другую задачу — он работает с живым удалённым источником:
- Аутентификация — нужен токен, который может протухнуть.
- Пагинация — API отдаёт по 100 записей, а страниц тысячи.
- Rate limiting — нельзя слать 1000 запросов в секунду.
- Иерархия — страница вложена в раздел, раздел — в пространство.
- Инкрементальная синхронизация — при повторном запуске брать только изменения.
- Разные форматы ответа — Notion отдаёт блоки JSON, Confluence — HTML, Drive — файлы.
На выходе коннектор, как и загрузчик, должен отдавать
list[Document] — единый интерфейс для остального пайплайна.
Аутентификация: токены, OAuth и сервисные аккаунты
Первое, что нужно сделать при работе с любым API — разобраться с аутентификацией. Каждый сервис использует свой метод, и ошибки здесь приводят к самым непонятным падениям.
secret_. Страницы нужно явно расшарить интеграции.email:token в base64. Cloud и Server — разные URL.OAuth 2.0 Client Credentials: как это работает
SharePoint (и многие корпоративные API) используют OAuth 2.0 в режиме Client Credentials — это flow для серверных приложений без участия пользователя. Понимать его важно, чтобы правильно настроить Azure AD и не путаться с токенами.
client_id и tenant_id. Создаём
client_secret (действует 1–2 года, нужно обновлять).
https://login.microsoftonline.com/{tenant_id}/oauth2/v2.0/token
с параметрами grant_type=client_credentials,
scope=https://graph.microsoft.com/.default.
В ответ — JWT-токен со сроком жизни 3600 секунд.
Authorization: Bearer <access_token>.
Токен кэшируем — не запрашиваем новый на каждый вызов.
ConfidentialClientApplication.
Базовый класс коннектора
Прежде чем писать под каждый сервис, вынесем общую логику: rate limiting, retry с exponential backoff и пагинацию. Эти паттерны одинаковы для всех API.
# connectors/base.py
import asyncio
import time
import logging
from abc import ABC, abstractmethod
from dataclasses import dataclass, field
from typing import AsyncIterator
from pathlib import Path
import httpx
from core.document import Document
logger = logging.getLogger(__name__)
@dataclass
class SyncState:
"""
Состояние синхронизации — сохраняем между запусками,
чтобы следующий раз брать только изменения.
"""
last_sync_at: str = "" # ISO 8601 timestamp
processed_ids: set[str] = field(default_factory=set)
cursor: str = "" # cursor/page_token для resume
def save(self, path: Path) -> None:
import json
data = {
"last_sync_at": self.last_sync_at,
"processed_ids": list(self.processed_ids),
"cursor": self.cursor,
}
path.write_text(json.dumps(data, ensure_ascii=False, indent=2))
@classmethod
def load(cls, path: Path) -> "SyncState":
import json
if not path.exists():
return cls()
data = json.loads(path.read_text())
return cls(
last_sync_at=data.get("last_sync_at", ""),
processed_ids=set(data.get("processed_ids", [])),
cursor=data.get("cursor", ""),
)
class RateLimiter:
"""Токен-баккет: пропускает не более max_rps запросов в секунду."""
def __init__(self, max_rps: float = 3.0):
self._interval = 1.0 / max_rps
self._last_call = 0.0
async def acquire(self) -> None:
now = time.monotonic()
wait = self._interval - (now - self._last_call)
if wait > 0:
await asyncio.sleep(wait)
self._last_call = time.monotonic()
class BaseConnector(ABC):
"""
Базовый класс для всех коннекторов.
Предоставляет: rate limiting, retry, общий HTTP-клиент.
"""
MAX_RETRIES = 3
RETRY_STATUSES = {429, 500, 502, 503, 504}
def __init__(self, max_rps: float = 3.0):
self._limiter = RateLimiter(max_rps)
self._client: httpx.AsyncClient | None = None
async def __aenter__(self):
self._client = httpx.AsyncClient(timeout=30)
return self
async def __aexit__(self, *args):
if self._client:
await self._client.aclose()
async def _get(self, url: str, **kwargs) -> dict:
"""GET с rate limiting и retry."""
await self._limiter.acquire()
for attempt in range(self.MAX_RETRIES):
response = await self._client.get(url, **kwargs)
if response.status_code == 200:
return response.json()
if response.status_code in self.RETRY_STATUSES:
# 429 может содержать Retry-After заголовок
retry_after = float(
response.headers.get("Retry-After", 2 ** attempt)
)
logger.warning(
f"HTTP {response.status_code} from {url}, "
f"retry {attempt + 1}/{self.MAX_RETRIES} "
f"after {retry_after:.1f}s"
)
await asyncio.sleep(retry_after)
else:
response.raise_for_status()
raise RuntimeError(f"Max retries exceeded for {url}")
async def _post(self, url: str, **kwargs) -> dict:
"""POST с rate limiting и retry (нужен для Notion API)."""
await self._limiter.acquire()
for attempt in range(self.MAX_RETRIES):
response = await self._client.post(url, **kwargs)
if response.status_code in (200, 201):
return response.json()
if response.status_code in self.RETRY_STATUSES:
await asyncio.sleep(2 ** attempt)
else:
response.raise_for_status()
raise RuntimeError(f"Max retries exceeded for {url}")
@abstractmethod
async def load(self, state: SyncState | None = None) -> list[Document]:
"""Загрузить документы. state=None → полная загрузка."""
...
Notion: блочная модель контента
Notion хранит содержимое страниц не как текст, а как дерево блоков. Страница — корневой блок. Внутри — параграфы, заголовки, списки, таблицы, вложенные страницы — каждый элемент является отдельным блоком с типом. Это мощная модель для редактирования, но для RAG её нужно «развернуть» в плоский текст.
Page (block type: "page")
├── heading_1: "Введение"
├── paragraph: "Текст вводного раздела…"
├── bulleted_list_item: "Пункт 1"
├── bulleted_list_item: "Пункт 2"
├── child_page: "Вложенная страница" ← рекурсивно
└── table
├── table_row: ["Имя", "Значение"]
└── table_row: ["foo", "bar"]
# connectors/notion_connector.py
import os
from datetime import datetime, timezone
from core.document import Document
from connectors.base import BaseConnector, SyncState
NOTION_VERSION = "2022-06-28"
# Типы блоков, которые содержат текст
TEXT_BLOCK_TYPES = {
"paragraph", "heading_1", "heading_2", "heading_3",
"bulleted_list_item", "numbered_list_item", "toggle",
"quote", "callout", "code",
}
HEADING_PREFIX = {
"heading_1": "# ", "heading_2": "## ", "heading_3": "### ",
}
def _extract_rich_text(rich_text_list: list) -> str:
"""Извлекаем текст из Notion rich_text массива."""
return "".join(rt.get("plain_text", "") for rt in rich_text_list)
def _block_to_text(block: dict) -> str:
"""Преобразуем блок в строку Markdown."""
btype = block.get("type", "")
content = block.get(btype, {})
if btype in TEXT_BLOCK_TYPES:
text = _extract_rich_text(content.get("rich_text", []))
prefix = HEADING_PREFIX.get(btype, "")
if btype == "bulleted_list_item":
prefix = "- "
elif btype == "code":
lang = content.get("language", "")
return f"```{lang}\n{text}\n```"
return f"{prefix}{text}" if text else ""
if btype == "table_row":
cells = [
_extract_rich_text(cell)
for cell in content.get("cells", [])
]
return "| " + " | ".join(cells) + " |"
return "" # image, divider, embed — пропускаем
class NotionConnector(BaseConnector):
"""
Загружает страницы Notion.
Обходит все дочерние блоки рекурсивно, собирает текст,
поддерживает инкрементальную синхронизацию по last_edited_time.
"""
BASE_URL = "https://api.notion.com/v1"
def __init__(self, token: str | None = None, max_rps: float = 3.0):
super().__init__(max_rps=max_rps)
self._token = token or os.environ["NOTION_TOKEN"]
@property
def _headers(self) -> dict:
return {
"Authorization": f"Bearer {self._token}",
"Notion-Version": NOTION_VERSION,
"Content-Type": "application/json",
}
async def load(self, state: SyncState | None = None) -> list[Document]:
"""Загружает все расшаренные с интеграцией страницы."""
pages = await self._search_all_pages(
edited_after=state.last_sync_at if state else ""
)
documents = []
for page in pages:
page_id = page["id"]
if state and page_id in state.processed_ids:
continue
doc = await self._load_page(page)
if doc:
documents.append(doc)
return documents
async def load_page(self, page_id: str) -> Document | None:
"""Загружает одну страницу по ID."""
page = await self._get(
f"{self.BASE_URL}/pages/{page_id}",
headers=self._headers,
)
return await self._load_page(page)
async def _search_all_pages(self, edited_after: str = "") -> list[dict]:
"""Получаем все страницы через Search API с пагинацией."""
pages = []
cursor = None
payload: dict = {
"filter": {"property": "object", "value": "page"},
"page_size": 100,
"sort": {"direction": "descending", "timestamp": "last_edited_time"},
}
if edited_after:
# Notion Search API не поддерживает фильтр по дате напрямую,
# фильтруем на стороне клиента
pass
while True:
if cursor:
payload["start_cursor"] = cursor
data = await self._post(
f"{self.BASE_URL}/search",
headers=self._headers,
json=payload,
)
results = data.get("results", [])
# Фильтрация по дате (инкрементальная синхронизация)
for page in results:
edited = page.get("last_edited_time", "")
if edited_after and edited <= edited_after:
return pages # отсортировано по убыванию — дальше только старые
pages.append(page)
if not data.get("has_more"):
break
cursor = data.get("next_cursor")
return pages
async def _load_page(self, page: dict) -> Document | None:
"""Получает все блоки страницы и собирает текст."""
page_id = page["id"]
title = _extract_rich_text(
page.get("properties", {})
.get("title", page.get("properties", {}).get("Name", {}))
.get("title", [])
)
blocks = await self._fetch_blocks(page_id)
lines = [_block_to_text(b) for b in blocks]
content = "\n\n".join(l for l in lines if l.strip())
if not content.strip():
return None
# URL страницы: https://notion.so/{page_id без дефисов}
url = "https://notion.so/" + page_id.replace("-", "")
return Document(
page_content=content,
metadata={
"source": url,
"page_id": page_id,
"title": title,
"last_edited": page.get("last_edited_time", ""),
"created": page.get("created_time", ""),
"url": url,
"connector": "notion",
}
)
async def _fetch_blocks(self, block_id: str, depth: int = 0) -> list[dict]:
"""Рекурсивно получаем все блоки (макс. глубина 3 уровня)."""
if depth > 3:
return []
blocks = []
cursor = None
while True:
url = f"{self.BASE_URL}/blocks/{block_id}/children"
params = {"page_size": 100}
if cursor:
params["start_cursor"] = cursor
data = await self._get(url, headers=self._headers, params=params)
for block in data.get("results", []):
blocks.append(block)
# Рекурсивно загружаем дочерние блоки
if block.get("has_children") and block["type"] != "child_page":
children = await self._fetch_blocks(block["id"], depth + 1)
blocks.extend(children)
if not data.get("has_more"):
break
cursor = data.get("next_cursor")
return blocks
Confluence: REST API и HTML-контент
Confluence отдаёт тело страницы в двух форматах: storage (внутренний XML)
и view (HTML для браузера). Для RAG удобнее view — это
стандартный HTML, который мы уже умеем парсить. Но нужно убрать
навигационные элементы Confluence (панель редактирования, макросы).
# connectors/confluence_connector.py
import os
import base64
from bs4 import BeautifulSoup
from core.document import Document
from connectors.base import BaseConnector, SyncState
class ConfluenceConnector(BaseConnector):
"""
Загружает страницы из Confluence Cloud через REST API v2.
Аутентификация: Basic Auth (email + API token).
"""
def __init__(
self,
base_url: str | None = None, # https://your-domain.atlassian.net
email: str | None = None,
token: str | None = None,
space_key: str | None = None, # если None — загружаем всё
max_rps: float = 5.0,
):
super().__init__(max_rps=max_rps)
self._base_url = (base_url or os.environ["CONFLUENCE_URL"]).rstrip("/")
_email = email or os.environ["CONFLUENCE_EMAIL"]
_token = token or os.environ["CONFLUENCE_TOKEN"]
# Basic Auth = base64(email:token)
credentials = base64.b64encode(f"{_email}:{_token}".encode()).decode()
self._auth_header = f"Basic {credentials}"
self._space_key = space_key
@property
def _headers(self) -> dict:
return {
"Authorization": self._auth_header,
"Accept": "application/json",
}
async def load(self, state: SyncState | None = None) -> list[Document]:
pages = await self._fetch_all_pages(
modified_after=state.last_sync_at if state else ""
)
documents = []
for page in pages:
doc = await self._load_page(page)
if doc:
documents.append(doc)
return documents
async def _fetch_all_pages(self, modified_after: str = "") -> list[dict]:
"""Получаем список страниц с пагинацией (cursor-based в API v2)."""
pages = []
url = f"{self._base_url}/wiki/api/v2/pages"
params: dict = {
"limit": 250,
"status": "current",
"body-format": "view", # сразу получаем HTML-тело
}
if self._space_key:
params["space-key"] = self._space_key
if modified_after:
# Confluence v2 поддерживает фильтр по версии, но не по дате напрямую
# Используем metadata.modified для фильтрации после загрузки
pass
while True:
data = await self._get(url, headers=self._headers, params=params)
results = data.get("results", [])
for page in results:
version_at = page.get("version", {}).get("createdAt", "")
if modified_after and version_at <= modified_after:
continue
pages.append(page)
# Cursor-based пагинация через Link header
links = data.get("_links", {})
next_url = links.get("next")
if not next_url:
break
# next_url — относительный путь
url = f"{self._base_url}{next_url}"
params = {} # параметры уже в next_url
return pages
async def _load_page(self, page: dict) -> Document | None:
"""Парсим HTML-тело страницы и формируем Document."""
page_id = page["id"]
title = page.get("title", "")
html = page.get("body", {}).get("view", {}).get("value", "")
if not html:
return None
text = self._html_to_text(html, title)
if not text.strip():
return None
space = page.get("spaceId", "")
version = page.get("version", {})
web_url = (
self._base_url
+ page.get("_links", {}).get("webui", f"/wiki/spaces/{space}/pages/{page_id}")
)
return Document(
page_content=text,
metadata={
"source": web_url,
"page_id": page_id,
"title": title,
"space_id": space,
"version": version.get("number", 1),
"last_modified": version.get("createdAt", ""),
"author": version.get("authorId", ""),
"url": web_url,
"connector": "confluence",
}
)
@staticmethod
def _html_to_text(html: str, title: str) -> str:
"""
Конвертируем HTML Confluence в чистый текст.
Удаляем confluence-специфичные элементы (макросы, панели).
"""
soup = BeautifulSoup(html, "lxml")
# Удаляем confluence-специфичные макросы и boilerplate
for tag in soup.find_all(class_=lambda c: c and (
"conf-macro" in c or "panel" in c or
"aui-" in c or "confluence-information-macro" in c
)):
tag.decompose()
parts = [f"# {title}"]
for element in soup.find_all(
["h1", "h2", "h3", "h4", "p", "li", "td", "th", "pre", "code"]
):
text = element.get_text(separator=" ", strip=True)
if not text:
continue
if element.name == "h1":
parts.append(f"# {text}")
elif element.name == "h2":
parts.append(f"## {text}")
elif element.name in ("h3", "h4"):
parts.append(f"### {text}")
elif element.name == "li":
parts.append(f"- {text}")
elif element.name in ("th", "td"):
parts.append(text) # таблицы — упрощённо
elif element.name in ("pre", "code"):
parts.append(f"`{text}`")
else:
parts.append(text)
return "\n\n".join(p for p in parts if p.strip())
Google Drive: Drive API + Docs API
Google Drive — это файловое хранилище, но «файлы» там бывают двух видов: обычные (PDF, DOCX, изображения) и нативные Google-форматы (Docs, Sheets, Slides). Нативные форматы нельзя скачать напрямую — их нужно экспортировать через отдельный метод.
# connectors/gdrive_connector.py
import os
import io
from pathlib import Path
from google.oauth2 import service_account
from googleapiclient.discovery import build
from googleapiclient.http import MediaIoBaseDownload
from core.document import Document
from connectors.base import SyncState
from loaders.document_dispatcher import DocumentDispatcher
# MIME-типы Google-форматов и их экспортный формат
GOOGLE_MIME_EXPORT = {
"application/vnd.google-apps.document": "text/plain",
"application/vnd.google-apps.spreadsheet": "text/csv",
"application/vnd.google-apps.presentation": "text/plain",
}
# Обычные файлы, которые умеем читать
SUPPORTED_MIME = {
"application/pdf",
"application/vnd.openxmlformats-officedocument.wordprocessingml.document",
"text/html", "text/markdown", "text/plain",
}
class GoogleDriveConnector:
"""
Загружает файлы из Google Drive.
Поддерживает: Google Docs, Sheets, Slides (экспорт),
PDF, DOCX, HTML, Markdown (прямое скачивание).
Аутентификация: Service Account (JSON key file).
"""
SCOPES = [
"https://www.googleapis.com/auth/drive.readonly",
]
def __init__(
self,
service_account_file: str | None = None,
folder_id: str | None = None, # None → весь Drive
modified_after: str = "", # ISO 8601, для инкр. синхронизации
):
key_file = service_account_file or os.environ["GOOGLE_SERVICE_ACCOUNT_FILE"]
credentials = service_account.Credentials.from_service_account_file(
key_file, scopes=self.SCOPES
)
self._drive = build("drive", "v3", credentials=credentials)
self._folder_id = folder_id
self._modified_after = modified_after
self._dispatcher = DocumentDispatcher()
def load(self, state: SyncState | None = None) -> list[Document]:
files = self._list_files(
modified_after=state.last_sync_at if state else self._modified_after
)
documents = []
for file_meta in files:
try:
docs = self._load_file(file_meta)
documents.extend(docs)
except Exception as e:
print(f" ✗ {file_meta['name']}: {e}")
return documents
def _list_files(self, modified_after: str = "") -> list[dict]:
"""Список файлов с пагинацией."""
query_parts = ["trashed = false"]
if self._folder_id:
query_parts.append(f"'{self._folder_id}' in parents")
if modified_after:
query_parts.append(f"modifiedTime > '{modified_after}'")
# Фильтр по поддерживаемым MIME
supported = list(GOOGLE_MIME_EXPORT) + list(SUPPORTED_MIME)
mime_filter = " or ".join(f"mimeType = '{m}'" for m in supported)
query_parts.append(f"({mime_filter})")
query = " and ".join(query_parts)
files = []
page_token = None
while True:
response = self._drive.files().list(
q=query,
fields="nextPageToken, files(id, name, mimeType, modifiedTime, webViewLink, parents)",
pageSize=200,
pageToken=page_token,
).execute()
files.extend(response.get("files", []))
page_token = response.get("nextPageToken")
if not page_token:
break
return files
def _load_file(self, file_meta: dict) -> list[Document]:
"""Скачивает файл и загружает через DocumentDispatcher."""
file_id = file_meta["id"]
file_name = file_meta["name"]
mime_type = file_meta["mimeType"]
web_url = file_meta.get("webViewLink", "")
# Определяем: экспорт или прямое скачивание
if mime_type in GOOGLE_MIME_EXPORT:
export_mime = GOOGLE_MIME_EXPORT[mime_type]
content = self._export_file(file_id, export_mime)
ext = ".txt" if "text/plain" in export_mime else ".csv"
else:
content = self._download_file(file_id)
ext = Path(file_name).suffix or ".bin"
# Сохраняем во временный файл и загружаем через dispatcher
import tempfile
with tempfile.NamedTemporaryFile(suffix=ext, delete=False) as tmp:
tmp.write(content)
tmp_path = Path(tmp.name)
try:
docs = self._dispatcher.load(tmp_path)
finally:
tmp_path.unlink(missing_ok=True)
# Обогащаем metadata Drive-специфичными полями
for doc in docs:
doc.metadata.update({
"source": web_url or f"gdrive://{file_id}",
"file_id": file_id,
"filename": file_name,
"drive_mime": mime_type,
"modified_time": file_meta.get("modifiedTime", ""),
"connector": "gdrive",
})
return docs
def _export_file(self, file_id: str, mime_type: str) -> bytes:
request = self._drive.files().export_media(
fileId=file_id, mimeType=mime_type
)
buf = io.BytesIO()
downloader = MediaIoBaseDownload(buf, request)
done = False
while not done:
_, done = downloader.next_chunk()
return buf.getvalue()
def _download_file(self, file_id: str) -> bytes:
request = self._drive.files().get_media(fileId=file_id)
buf = io.BytesIO()
downloader = MediaIoBaseDownload(buf, request)
done = False
while not done:
_, done = downloader.next_chunk()
return buf.getvalue()
SharePoint: Microsoft Graph API
SharePoint доступен через Microsoft Graph API — единый REST-шлюз для всех Microsoft 365 сервисов. Для серверного приложения используем Client Credentials flow через MSAL (Microsoft Authentication Library).
# connectors/sharepoint_connector.py
import os
from pathlib import Path
from msal import ConfidentialClientApplication
from core.document import Document
from connectors.base import BaseConnector, SyncState
from loaders.document_dispatcher import DocumentDispatcher
class SharePointConnector(BaseConnector):
"""
Загружает файлы из SharePoint через Microsoft Graph API.
Аутентификация: OAuth 2.0 Client Credentials (MSAL).
"""
GRAPH_URL = "https://graph.microsoft.com/v1.0"
SCOPE = ["https://graph.microsoft.com/.default"]
def __init__(
self,
tenant_id: str | None = None,
client_id: str | None = None,
client_secret: str | None = None,
site_id: str | None = None, # ID сайта SharePoint
drive_id: str | None = None, # ID библиотеки документов
folder_path: str = "/", # путь внутри библиотеки
max_rps: float = 5.0,
):
super().__init__(max_rps=max_rps)
self._site_id = site_id or os.environ["SHAREPOINT_SITE_ID"]
self._drive_id = drive_id or os.environ.get("SHAREPOINT_DRIVE_ID", "")
self._folder = folder_path
# MSAL кэширует токены автоматически
self._msal_app = ConfidentialClientApplication(
client_id=client_id or os.environ["AZURE_CLIENT_ID"],
client_credential=client_secret or os.environ["AZURE_CLIENT_SECRET"],
authority=f"https://login.microsoftonline.com/"
f"{tenant_id or os.environ['AZURE_TENANT_ID']}",
)
self._dispatcher = DocumentDispatcher()
def _get_token(self) -> str:
"""Получаем токен из кэша или запрашиваем новый."""
result = self._msal_app.acquire_token_silent(self.SCOPE, account=None)
if not result:
result = self._msal_app.acquire_token_for_client(scopes=self.SCOPE)
if "access_token" not in result:
raise RuntimeError(f"MSAL error: {result.get('error_description')}")
return result["access_token"]
@property
def _auth_header(self) -> dict:
return {"Authorization": f"Bearer {self._get_token()}"}
async def load(self, state: SyncState | None = None) -> list[Document]:
items = await self._list_items(
modified_after=state.last_sync_at if state else ""
)
documents = []
for item in items:
try:
docs = await self._load_item(item)
documents.extend(docs)
except Exception as e:
print(f" ✗ {item.get('name')}: {e}")
return documents
async def _list_items(self, modified_after: str = "") -> list[dict]:
"""Рекурсивный обход папок через Graph API."""
items = []
# Если drive_id не задан — берём default drive сайта
if self._drive_id:
base = f"{self.GRAPH_URL}/drives/{self._drive_id}"
else:
base = f"{self.GRAPH_URL}/sites/{self._site_id}/drive"
path = self._folder.strip("/")
if path:
list_url = f"{base}/root:/{path}:/children"
else:
list_url = f"{base}/root/children"
await self._collect_items(list_url, items, modified_after)
return items
async def _collect_items(
self, url: str, items: list, modified_after: str
) -> None:
"""Рекурсивно собираем файлы, пропуская папки."""
params = {"$top": 200, "$select": "id,name,file,folder,lastModifiedDateTime,webUrl,size"}
if modified_after:
params["$filter"] = f"lastModifiedDateTime gt {modified_after}"
data = await self._get(url, headers=self._auth_header, params=params)
for item in data.get("value", []):
if "folder" in item:
# Рекурсивно входим в подпапку
folder_url = f"{url.rsplit('/children', 1)[0]}/{item['id']}/children"
await self._collect_items(folder_url, items, modified_after)
elif "file" in item:
items.append(item)
# Пагинация через @odata.nextLink
next_link = data.get("@odata.nextLink")
if next_link:
await self._collect_items(next_link, items, modified_after)
async def _load_item(self, item: dict) -> list[Document]:
"""Скачивает файл и обрабатывает через DocumentDispatcher."""
item_id = item["id"]
filename = item["name"]
web_url = item.get("webUrl", "")
# Ссылка на скачивание
if self._drive_id:
download_url_endpoint = (
f"{self.GRAPH_URL}/drives/{self._drive_id}/items/{item_id}/content"
)
else:
download_url_endpoint = (
f"{self.GRAPH_URL}/sites/{self._site_id}/drive/items/{item_id}/content"
)
# Graph API возвращает 302 redirect на временную ссылку
response = await self._client.get(
download_url_endpoint,
headers=self._auth_header,
follow_redirects=True,
)
response.raise_for_status()
content = response.content
ext = Path(filename).suffix or ".bin"
import tempfile
with tempfile.NamedTemporaryFile(suffix=ext, delete=False) as tmp:
tmp.write(content)
tmp_path = Path(tmp.name)
try:
docs = self._dispatcher.load(tmp_path)
except ValueError:
return [] # неподдерживаемый тип
finally:
tmp_path.unlink(missing_ok=True)
for doc in docs:
doc.metadata.update({
"source": web_url,
"item_id": item_id,
"filename": filename,
"last_modified": item.get("lastModifiedDateTime", ""),
"file_size": item.get("size", 0),
"connector": "sharepoint",
})
return docs
Инкрементальная синхронизация
Полная загрузка 10 000 страниц Confluence занимает несколько часов. Повторять это каждую ночь — расточительно. Инкрементальная синхронизация берёт только изменения с момента последнего запуска.
# sync/pipeline.py
import asyncio
from datetime import datetime, timezone
from pathlib import Path
from core.document import Document
from connectors.base import SyncState
from connectors.notion_connector import NotionConnector
from connectors.confluence_connector import ConfluenceConnector
STATE_DIR = Path(".sync_state")
STATE_DIR.mkdir(exist_ok=True)
async def sync_notion(full: bool = False) -> list[Document]:
"""
Синхронизирует Notion.
full=True → игнорировать состояние, загрузить всё заново.
"""
state_path = STATE_DIR / "notion.json"
state = SyncState() if full else SyncState.load(state_path)
async with NotionConnector() as connector:
docs = await connector.load(state=None if full else state)
# Обновляем состояние
now = datetime.now(timezone.utc).isoformat()
new_state = SyncState(
last_sync_at=now,
processed_ids=state.processed_ids | {d.metadata["page_id"] for d in docs},
)
new_state.save(state_path)
print(f"Notion: synced {len(docs)} documents (since {state.last_sync_at or 'beginning'})")
return docs
async def sync_all(full: bool = False) -> list[Document]:
"""Параллельная синхронизация всех источников."""
results = await asyncio.gather(
sync_notion(full=full),
# sync_confluence(full=full),
# sync_gdrive(full=full),
return_exceptions=True,
)
documents = []
for r in results:
if isinstance(r, Exception):
print(f"Sync error: {r}")
else:
documents.extend(r)
return documents
# Пример запуска
if __name__ == "__main__":
import sys
full_sync = "--full" in sys.argv
docs = asyncio.run(sync_all(full=full_sync))
print(f"Total documents ready for indexing: {len(docs)}")
Краткий обзор: что умеет каждый сервис
lastModified > "2024-01-01"_links.nextmodifiedTime > '…' в query$filter=lastModifiedDateTime gt …Типичные ошибки
filter_properties.python-dotenv для локальной разработки.Шпаргалка
- Notion — API key (
secret_…), 3 req/s, блоки → рекурсия, Search API для листинга. - Confluence — Basic Auth (email:API token в base64), REST v2, body-format=view → HTML → BeautifulSoup.
- Google Drive — Service Account JSON, нативные форматы → export_media(), обычные → get_media().
- SharePoint — MSAL Client Credentials, Graph API,
@odata.nextLinkдля пагинации, 302 redirect на download. - Rate limiting — токен-баккет перед каждым запросом, Retry-After при 429.
- Инкрементальная синхронизация — сохранять last_sync_at и cursor, фильтровать по modified date.
- Секреты — только env vars, никогда в коде.
- Результат — всегда
list[Document]сsource,connector,last_modifiedв metadata.
Практическое задание
-
Notion коннектор.
Создайте интеграцию в Notion, расшарьте на неё 2–3 страницы.
Реализуйте
NotionConnectorи загрузите страницы. Убедитесь, что в metadata естьtitle,page_idиlast_edited. Попробуйте страницу с вложенными блоками — проверьте глубину рекурсии. -
Rate limiter под нагрузкой.
Напишите тест: создайте
RateLimiter(max_rps=5)и выполните 20await limiter.acquire()подряд, замерив общее время. Убедитесь, что 20 вызовов занимают ≥ 4 секунды (20 / 5 rps = 4s). Поэкспериментируйте с параметром. -
Инкрементальная синхронизация.
Запустите загрузку Notion дважды: первый раз — полная загрузка,
второй — инкрементальная (с сохранённым
SyncState). Измените одну страницу в Notion между запусками. Убедитесь, что второй запуск загружает только изменённую страницу, а не все.