Ответ
В конвейерах данных (ETL/ELT) я сталкивался с разнообразными источниками:
1. Базы данных:
- SQL (PostgreSQL, MySQL): Использовал
psycopg2иSQLAlchemyдля извлечения, часто с инкрементальной загрузкой. - NoSQL (MongoDB): Работал через
pymongo, извлекал вложенные документы.
2. Файловые источники:
- CSV/JSON/Excel: Загружал с помощью
pandas(read_csv,read_json), решая проблемы с кодировкой и разделителями. - Parquet/ORC: Использовал для эффективной работы с большими объёмами в экосистеме Spark (
PySpark).
3. API и потоки данных:
- REST API: Писал скрипты на
requestsиaiohttpдля асинхронного сбора данных, обрабатывал пагинацию и ограничения速率. - Apache Kafka: Потреблял потоки событий с помощью
confluent-kafkaдля обработки в реальном времени.
4. Облачные хранилища:
- Amazon S3 / Google Cloud Storage: Использовал
boto3и соответствующие SDK для загрузки файловых датасетов.
Пример инкрементальной загрузки из PostgreSQL:
import pandas as pd
from sqlalchemy import create_engine
engine = create_engine('postgresql://user:pass@localhost/db')
# Загрузка только новых данных с момента last_update
query = """
SELECT * FROM sales
WHERE updated_at > %(last_run)s
"""
df = pd.read_sql(query, engine, params={'last_run': '2023-11-01'})
Ключевые задачи — обеспечение надёжности, обработка ошибок, валидация схемы и эффективное использование памяти.