Multi-tenant ETL-пайплайн: Spotify -> Postgres staging -> ClickHouse (DWH) через Apache Airflow.
В аналитике прослушивания считаются "по всем артистам трека" (track.artists), а агрегации считаются в ClickHouse.
- Извлечение
recently-playedиз Spotify API (по cursorafter) для каждого tenant/user. - Идемпотентная загрузка в Postgres staging через дедупликацию по
event_hash(безreplace). - Загрузка новых событий в ClickHouse и построение дневной агрегации
artist_dailyпо всем артистам.
- В Airflow должен быть создан Connection с id
postgre_sql, указывающий на Postgres staging (где будут созданыstg_*таблицы). - Для multi-tenant укажите
SPOTIFY_ACCOUNTS_JSONкак JSON-массив объектов{ tenant_id, spotify_user_id, spotify_token }. - Либо оставьте fallback single-tenant:
TENANT_ID(опционально),SPOTIFY_USER_ID,SPOTIFY_TOKEN. - Для ClickHouse задайте:
CLICKHOUSE_HOST,CLICKHOUSE_PORT(по умолчанию 9000),CLICKHOUSE_USER,CLICKHOUSE_PASSWORD,CLICKHOUSE_DATABASE(по умолчаниюspotify). - Для точной настройки курсора (window):
SPOTIFY_CURSOR_OVERLAP_MINUTESиSPOTIFY_INITIAL_LOOKBACK_HOURS.
Spotify-data-pipeline/
├── .gitignore
├── requirements.txt
└── src/
├── Dags/
│ ├── spotify_staging_pipeline_dag.py
│ └── spotify_clickhouse_loader_dag.py
├── spotify_pipeline/
└── docker-compose.yaml
- Python 3.8+
- PostgreSQL 13+
- Apache Airflow 2.5.1+
- ClickHouse
- Spotify API-токены для tenant/user