From dc007393cb4cd0440987c917f6bb5c8be5f2f10d Mon Sep 17 00:00:00 2001 From: Tobias Macey Date: Wed, 1 Jul 2026 09:01:20 -0400 Subject: [PATCH 1/3] feat(ingest): add YouTube dlt source to ol_dlt Migrates MIT Learn's YouTube catalog ETL to the data platform as a pure-dlt source in src/ol_dlt, so channel/playlist/video/transcript extraction is observable and backfillable in Dagster rather than opaque in the Django app. Reads per-channel configs from mitodl/open-video-data and calls the YouTube Data API v3 over plain requests (consistent with the other ol_dlt sources) plus youtube-transcript-api for transcripts, landing raw__youtube__api__{channels,playlists,videos,transcripts}. Wired into the data_loading code location as @dlt_assets with a daily schedule. Co-Authored-By: Claude Opus 4.8 --- .../data_loading/defs/ingestion/assets.py | 7 + .../data_loading/defs/ingestion/schedules.py | 10 + dg_projects/data_loading/uv.lock | 15 + src/ol_dlt/ol_dlt/sources/youtube/README.md | 50 +++ src/ol_dlt/ol_dlt/sources/youtube/__init__.py | 370 ++++++++++++++++++ src/ol_dlt/ol_dlt/sources/youtube/__main__.py | 10 + src/ol_dlt/pyproject.toml | 1 + src/ol_dlt/tests/sources/test_youtube.py | 223 +++++++++++ src/ol_dlt/uv.lock | 24 ++ 9 files changed, 710 insertions(+) create mode 100644 src/ol_dlt/ol_dlt/sources/youtube/README.md create mode 100644 src/ol_dlt/ol_dlt/sources/youtube/__init__.py create mode 100644 src/ol_dlt/ol_dlt/sources/youtube/__main__.py create mode 100644 src/ol_dlt/tests/sources/test_youtube.py diff --git a/dg_projects/data_loading/data_loading/defs/ingestion/assets.py b/dg_projects/data_loading/data_loading/defs/ingestion/assets.py index 050c95b88..c0729cb7c 100644 --- a/dg_projects/data_loading/data_loading/defs/ingestion/assets.py +++ b/dg_projects/data_loading/data_loading/defs/ingestion/assets.py @@ -19,6 +19,7 @@ mitpe, oll, podcast_rss, + youtube, ) from ol_orchestrate.lib.constants import EDXORG_DB_TABLES @@ -73,6 +74,11 @@ def _assets( source=podcast_rss.build_source(), pipeline=podcast_rss.podcast_rss_pipeline, ) +youtube_assets = build_ingest_assets( + name="youtube_ingest", + source=youtube.build_source(), + pipeline=youtube.youtube_pipeline, +) # --- edxorg_s3: custom upstream deps + one op per table --------------------- @@ -131,6 +137,7 @@ def _asset( mit_climate_assets, mit_edx_programs_assets, podcast_rss_assets, + youtube_assets, *edxorg_s3_table_assets, ], ) diff --git a/dg_projects/data_loading/data_loading/defs/ingestion/schedules.py b/dg_projects/data_loading/data_loading/defs/ingestion/schedules.py index 12bc9cd66..04bda7270 100644 --- a/dg_projects/data_loading/data_loading/defs/ingestion/schedules.py +++ b/dg_projects/data_loading/data_loading/defs/ingestion/schedules.py @@ -48,6 +48,15 @@ execution_timezone="Etc/UTC", ) +# The four raw__youtube__api__* tables are materialized by a single @dlt_assets +# run, so schedule the whole youtube source group rather than one table. +youtube_ingest_schedule = dg.ScheduleDefinition( + name="youtube_ingest_daily_schedule", + target=dg.AssetSelection.groups("youtube"), + cron_schedule="15 4 * * *", + execution_timezone="Etc/UTC", +) + defs = dg.Definitions( schedules=[ oll_ingest_schedule, @@ -55,5 +64,6 @@ mit_climate_ingest_schedule, mit_edx_programs_ingest_schedule, podcast_rss_ingest_schedule, + youtube_ingest_schedule, ], ) diff --git a/dg_projects/data_loading/uv.lock b/dg_projects/data_loading/uv.lock index eb642a5d5..992b9e4c4 100644 --- a/dg_projects/data_loading/uv.lock +++ b/dg_projects/data_loading/uv.lock @@ -1892,6 +1892,7 @@ dependencies = [ { name = "duckdb" }, { name = "pyiceberg", extra = ["glue"] }, { name = "requests" }, + { name = "youtube-transcript-api" }, ] [package.metadata] @@ -1901,6 +1902,7 @@ requires-dist = [ { name = "duckdb", specifier = ">=1.0,!=1.5.0,!=1.5.1" }, { name = "pyiceberg", extras = ["glue"], specifier = ">=0.9" }, { name = "requests", specifier = ">=2.32" }, + { name = "youtube-transcript-api", specifier = ">=1.0,<2" }, ] [package.metadata.requires-dev] @@ -3586,6 +3588,19 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/93/6f/7403e6ae864a0a7f1cdd8814d39690062766e141339127f2b3469201ff6f/yaspin-3.4.0-py3-none-any.whl", hash = "sha256:2a40572a38d39846d0df0a421733459481b7da17789f7a2618c3181bb0a82819", size = 21822, upload-time = "2025-12-06T12:33:50.633Z" }, ] +[[package]] +name = "youtube-transcript-api" +version = "1.2.4" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "defusedxml" }, + { name = "requests" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/60/43/4104185a2eaa839daa693b30e15c37e7e58795e8e09ec414f22b3db54bec/youtube_transcript_api-1.2.4.tar.gz", hash = "sha256:b72d0e96a335df599d67cee51d49e143cff4f45b84bcafc202ff51291603ddcd", size = 469839, upload-time = "2026-01-29T09:09:17.088Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/be/95/129ea37efd6cd6ed00f62baae6543345c677810b8a3bf0026756e1d3cf3c/youtube_transcript_api-1.2.4-py3-none-any.whl", hash = "sha256:03878759356da5caf5edac77431780b91448fb3d8c21d4496015bdc8a7bc43ff", size = 485227, upload-time = "2026-01-29T09:09:15.427Z" }, +] + [[package]] name = "zstandard" version = "0.25.0" diff --git a/src/ol_dlt/ol_dlt/sources/youtube/README.md b/src/ol_dlt/ol_dlt/sources/youtube/README.md new file mode 100644 index 000000000..22b9bf786 --- /dev/null +++ b/src/ol_dlt/ol_dlt/sources/youtube/README.md @@ -0,0 +1,50 @@ +# YouTube source + +Loads MIT Learn YouTube data from the YouTube Data API v3 into the raw-data +layer, replacing the `learning_resources/etl/youtube.py` ETL in the MIT Learn +application. Wrapped as Dagster assets by the `data_loading` code location +(`data_loading/defs/ingestion/assets.py`). + +**Config source:** [mitodl/open-video-data](https://github.com/mitodl/open-video-data) +(`youtube/*.yaml`, one file per channel). Each config provides a `channel_id` +and an optional `playlists` list (`{id, ignore?}`). A config with no `playlists`, +or one containing the `all` wildcard id, ingests every playlist on the channel. + +## Data flow + +``` +GitHub (mitodl/open-video-data/youtube/*.yaml) + → channel_id + playlist configs per channel + → YouTube Data API v3 (channels / playlists / playlistItems / videos) + → raw__youtube__api__channels (one row per configured channel) + → raw__youtube__api__playlists (one row per ingested playlist) + → raw__youtube__api__videos (one row per video) + → youtube-transcript-api + → raw__youtube__api__transcripts (one row per video with a transcript) +``` + +All tables use `write_disposition="merge"` keyed on the entity id so reruns +update in place. + +## Implementation notes + +- The Data API is called directly over HTTP with `requests` (API-key auth via + the `key` query param), not the `google-api-python-client` SDK, to stay + consistent with the other ol_dlt sources. +- Transcripts are **not** part of the Data API; they are fetched separately with + `youtube-transcript-api` (coded against its 1.x `fetch()` instance API). + Videos whose transcripts are disabled or unavailable are skipped (logged), + not raised. + +## Configuration + +- **Secret:** `YOUTUBE_DEVELOPER_KEY` — YouTube Data API v3 key, resolved lazily + at run time (`config.resolve_secret` / `config.require_secrets`). +- **Destination:** built by `ol_dlt.config.pipeline_for("youtube")` from the + active `DLT_PROFILE` (no per-source `.dlt/config.toml` block). + +## Standalone run + +```bash +DLT_PROFILE=dev YOUTUBE_DEVELOPER_KEY=... python -m ol_dlt.sources.youtube +``` diff --git a/src/ol_dlt/ol_dlt/sources/youtube/__init__.py b/src/ol_dlt/ol_dlt/sources/youtube/__init__.py new file mode 100644 index 000000000..2d9b25b01 --- /dev/null +++ b/src/ol_dlt/ol_dlt/sources/youtube/__init__.py @@ -0,0 +1,370 @@ +"""MIT Learn YouTube ingestion via dlt. + +Fetches YouTube channel configuration from the mitodl/open-video-data GitHub +repository, then loads raw channel, playlist, video, and transcript data from +the YouTube Data API v3. + +The Data API is called directly over HTTP with ``requests`` (API-key auth via +the ``key`` query param), not the ``google-api-python-client`` SDK, to stay +consistent with the other ol_dlt sources. Transcripts are not part of the Data +API, so they are retrieved separately with ``youtube-transcript-api``. + +Data flow: + GitHub (mitodl/open-video-data/youtube/*.yaml) + -> channel_id + playlist configs per channel + -> raw__youtube__api__channels (one row per configured channel) + -> raw__youtube__api__playlists (one row per playlist) + -> raw__youtube__api__videos (one row per video) + -> raw__youtube__api__transcripts (one row per video with a transcript) + +Run standalone: + DLT_PROFILE=dev YOUTUBE_DEVELOPER_KEY=... python -m ol_dlt.sources.youtube +""" + +import base64 +import logging +from collections.abc import Generator, Iterable +from typing import Any + +import dlt +import requests +import yaml + +from ol_dlt import config + +logger = logging.getLogger(__name__) + +GITHUB_API_BASE = "https://api.github.com" +YOUTUBE_API_BASE = "https://www.googleapis.com/youtube/v3" +YOUTUBE_MAX_RESULTS = 50 +# Sentinel playlist id meaning "ingest every playlist on the channel". +WILDCARD_PLAYLIST_ID = "all" + +_CONFIG_FILE_REPO_DEFAULT = "mitodl/open-video-data" +_CONFIG_FILE_FOLDER_DEFAULT = "youtube" + + +def _github_headers(token: str | None) -> dict[str, str]: + headers = {"Accept": "application/vnd.github.v3+json"} + if token: + headers["Authorization"] = f"Bearer {token}" + return headers + + +def _fetch_channel_configs( + repo: str, + folder: str, + branch: str, + token: str | None, +) -> list[dict[str, Any]]: + """Fetch and parse all YouTube channel YAML config files from a GitHub repo. + + Each config must contain a ``channel_id`` to be included. The ``playlists`` + key (a list of ``{id, ignore?}`` dicts) is optional; a config with no + ``playlists`` is treated as a wildcard (ingest all of the channel's + playlists). + """ + headers = _github_headers(token) + listing_url = f"{GITHUB_API_BASE}/repos/{repo}/contents/{folder}?ref={branch}" + resp = requests.get(listing_url, headers=headers, timeout=30) + resp.raise_for_status() + + configs = [] + for file_meta in resp.json(): + if not file_meta["name"].endswith((".yaml", ".yml")): + continue + + file_resp = requests.get(file_meta["url"], headers=headers, timeout=30) + file_resp.raise_for_status() + raw_content = base64.b64decode(file_resp.json()["content"]).decode("utf-8") + + try: + channel_config = yaml.safe_load(raw_content) + except yaml.YAMLError: + logger.exception("Failed to parse YAML config: %s", file_meta["name"]) + continue + + if not isinstance(channel_config, dict): + logger.warning("Skipping non-dict youtube config: %s", file_meta["name"]) + continue + + if "channel_id" not in channel_config: + logger.warning( + "Skipping config missing required key (channel_id): %s", + file_meta["name"], + ) + continue + + configs.append(channel_config) + + logger.info("Loaded %d youtube configs from %s/%s", len(configs), repo, folder) + return configs + + +def _resolve_api_key(api_key: str | None) -> str: + """Resolve the YouTube Data API key lazily, failing loudly if it is absent.""" + return config.require_secrets( + YOUTUBE_DEVELOPER_KEY=config.resolve_secret(api_key, "YOUTUBE_DEVELOPER_KEY") + )["YOUTUBE_DEVELOPER_KEY"] + + +def _yt_paged_items( + endpoint: str, + params: dict[str, Any], + api_key: str, +) -> Generator[dict[str, Any]]: + """Yield every ``items`` entry from a paginated Data API v3 endpoint. + + Follows ``nextPageToken`` until the API stops returning one. + """ + page_params = {**params, "key": api_key, "maxResults": YOUTUBE_MAX_RESULTS} + page_token: str | None = None + while True: + if page_token: + page_params["pageToken"] = page_token + resp = requests.get( + f"{YOUTUBE_API_BASE}/{endpoint}", params=page_params, timeout=30 + ) + resp.raise_for_status() + data = resp.json() + yield from data.get("items", []) + page_token = data.get("nextPageToken") + if not page_token: + break + + +def _playlist_ids_for_config( + channel_config: dict[str, Any], api_key: str +) -> Generator[str]: + """Yield the playlist ids to ingest for a single channel config. + + A config with no ``playlists`` list, or one containing the ``"all"`` + wildcard, expands to every playlist on the channel. Playlists flagged + ``ignore: true`` are skipped. + """ + channel_id = channel_config["channel_id"] + playlist_configs = channel_config.get("playlists") or [] + configs_by_id = { + pc["id"]: pc for pc in playlist_configs if isinstance(pc, dict) and "id" in pc + } + + if not playlist_configs or WILDCARD_PLAYLIST_ID in configs_by_id: + for playlist in _yt_paged_items( + "playlists", {"part": "id", "channelId": channel_id}, api_key + ): + playlist_id = playlist["id"] + if configs_by_id.get(playlist_id, {}).get("ignore", False): + continue + yield playlist_id + return + + for playlist_id, playlist_config in configs_by_id.items(): + if playlist_config.get("ignore", False): + continue + yield playlist_id + + +def _video_ids_for_playlist(playlist_id: str, api_key: str) -> Generator[str]: + """Yield the video ids contained in a playlist.""" + for item in _yt_paged_items( + "playlistItems", {"part": "contentDetails", "playlistId": playlist_id}, api_key + ): + video_id = item.get("contentDetails", {}).get("videoId") + if video_id: + yield video_id + + +def _batched(items: Iterable[str], size: int) -> Generator[list[str]]: + """Yield successive lists of at most ``size`` items.""" + batch: list[str] = [] + for item in items: + batch.append(item) + if len(batch) >= size: + yield batch + batch = [] + if batch: + yield batch + + +def _channel_video_ids(configs: list[dict[str, Any]], api_key: str) -> Generator[str]: + """Yield the deduplicated video ids across every configured channel.""" + seen: set[str] = set() + for channel_config in configs: + for playlist_id in _playlist_ids_for_config(channel_config, api_key): + for video_id in _video_ids_for_playlist(playlist_id, api_key): + if video_id not in seen: + seen.add(video_id) + yield video_id + + +@dlt.source(name="youtube") +def youtube_source( # noqa: C901 + api_key: str | None = None, + github_access_token: str | None = None, + github_repo: str = _CONFIG_FILE_REPO_DEFAULT, + github_folder: str = _CONFIG_FILE_FOLDER_DEFAULT, + github_branch: str = "master", +) -> Generator[Any]: + """Load MIT Learn YouTube data from the Data API v3 into four raw tables. + + Channel configs are read from YAML files in the mitodl/open-video-data + GitHub repository (one file per channel). For each channel the source + yields records into: + + raw__youtube__api__channels - one record per configured channel + raw__youtube__api__playlists - one record per ingested playlist + raw__youtube__api__videos - one record per video across all playlists + raw__youtube__api__transcripts - one record per video that has a transcript + + Credentials are resolved lazily at execution time (not at import) so the + module loads cleanly when secrets are absent in local development. + + Args: + api_key: YouTube Data API v3 key. Resolved from YOUTUBE_DEVELOPER_KEY + if not provided. + github_access_token: Optional GitHub token to raise the config-fetch + rate limit; the public repo works unauthenticated. + github_repo: GitHub repository containing the channel YAML configs. + github_folder: Folder within the repo containing YAML config files. + github_branch: Git branch to read configs from. + """ + table_format = config.active_table_format() + + def _configs() -> list[dict[str, Any]]: + return _fetch_channel_configs( + repo=github_repo, + folder=github_folder, + branch=github_branch, + token=github_access_token, + ) + + @dlt.resource( + name="raw__youtube__api__channels", + write_disposition="merge", + primary_key="channel_id", + table_format=table_format, + ) + def youtube_channels() -> Generator[dict[str, Any]]: + """Yield one record per configured YouTube channel.""" + key = _resolve_api_key(api_key) + for channel_config in _configs(): + channel_id = channel_config["channel_id"] + items = list( + _yt_paged_items( + "channels", + {"part": "snippet,contentDetails,statistics", "id": channel_id}, + key, + ) + ) + if not items: + logger.warning("No channel data returned for channel_id=%s", channel_id) + continue + yield { + "channel_id": channel_id, + "offered_by": channel_config.get("offered_by"), + "etl_source": "youtube", + **items[0], + } + + @dlt.resource( + name="raw__youtube__api__playlists", + write_disposition="merge", + primary_key="playlist_id", + table_format=table_format, + ) + def youtube_playlists() -> Generator[dict[str, Any]]: + """Yield one record per ingested playlist across all channels.""" + key = _resolve_api_key(api_key) + for channel_config in _configs(): + channel_id = channel_config["channel_id"] + for playlist_id in _playlist_ids_for_config(channel_config, key): + items = list( + _yt_paged_items( + "playlists", + {"part": "snippet,contentDetails", "id": playlist_id}, + key, + ) + ) + if not items: + logger.warning("No playlist data for playlist_id=%s", playlist_id) + continue + yield { + "playlist_id": playlist_id, + "channel_id": channel_id, + "etl_source": "youtube", + **items[0], + } + + @dlt.resource( + name="raw__youtube__api__videos", + write_disposition="merge", + primary_key="video_id", + table_format=table_format, + ) + def youtube_videos() -> Generator[dict[str, Any]]: + """Yield one record per video across all configured playlists.""" + key = _resolve_api_key(api_key) + for batch in _batched(_channel_video_ids(_configs(), key), YOUTUBE_MAX_RESULTS): + for video in _yt_paged_items( + "videos", + {"part": "snippet,contentDetails,statistics", "id": ",".join(batch)}, + key, + ): + yield { + "video_id": video["id"], + "etl_source": "youtube", + **video, + } + + @dlt.resource( + name="raw__youtube__api__transcripts", + write_disposition="merge", + primary_key="video_id", + table_format=table_format, + ) + def youtube_transcripts() -> Generator[dict[str, Any]]: + """Yield one record per video that has a transcript. + + Transcripts are fetched with ``youtube-transcript-api`` (imported lazily + so the module loads without the dependency present). Videos with + transcripts disabled or unavailable are skipped, not raised. + """ + from youtube_transcript_api import ( + NoTranscriptFound, + TranscriptsDisabled, + VideoUnavailable, + YouTubeTranscriptApi, + ) + from youtube_transcript_api.formatters import TextFormatter + + key = _resolve_api_key(api_key) + ytt_api = YouTubeTranscriptApi() + formatter = TextFormatter() + for video_id in _channel_video_ids(_configs(), key): + try: + fetched = ytt_api.fetch(video_id) + except (NoTranscriptFound, TranscriptsDisabled, VideoUnavailable): + logger.debug("No transcript available for video_id=%s", video_id) + continue + except Exception: + logger.exception("Failed to fetch transcript for video_id=%s", video_id) + continue + yield { + "video_id": video_id, + "etl_source": "youtube", + "transcript": formatter.format_transcript(fetched), + "segments": fetched.to_raw_data(), + } + + yield youtube_channels + yield youtube_playlists + yield youtube_videos + yield youtube_transcripts + + +youtube_pipeline = config.pipeline_for("youtube") + + +def build_source() -> Any: # noqa: ANN401 + """Instantiate the source (uniform entrypoint for the Dagster wrapper).""" + return youtube_source() diff --git a/src/ol_dlt/ol_dlt/sources/youtube/__main__.py b/src/ol_dlt/ol_dlt/sources/youtube/__main__.py new file mode 100644 index 000000000..b3f8f2e95 --- /dev/null +++ b/src/ol_dlt/ol_dlt/sources/youtube/__main__.py @@ -0,0 +1,10 @@ +"""Standalone smoke run: ``DLT_PROFILE=dev python -m ol_dlt.sources.youtube``.""" + +import logging + +from ol_dlt.sources.youtube import build_source, youtube_pipeline + +logging.basicConfig(level=logging.INFO) +logging.getLogger(__name__).info( + "Pipeline completed: %s", youtube_pipeline.run(build_source()) +) diff --git a/src/ol_dlt/pyproject.toml b/src/ol_dlt/pyproject.toml index 619d4967c..2f1ad535a 100644 --- a/src/ol_dlt/pyproject.toml +++ b/src/ol_dlt/pyproject.toml @@ -11,6 +11,7 @@ dependencies = [ "requests>=2.32", "duckdb>=1.0,!=1.5.0,!=1.5.1", "defusedxml>=0.7", + "youtube-transcript-api>=1.0,<2", ] [dependency-groups] diff --git a/src/ol_dlt/tests/sources/test_youtube.py b/src/ol_dlt/tests/sources/test_youtube.py new file mode 100644 index 000000000..621753df7 --- /dev/null +++ b/src/ol_dlt/tests/sources/test_youtube.py @@ -0,0 +1,223 @@ +"""Unit + materialization tests for the YouTube source.""" + +import base64 +import sys +import types +from pathlib import Path +from typing import Any + +import pytest + +from ol_dlt import config +from ol_dlt.sources import youtube +from tests.conftest import FakeResponse + + +def _queue_get(monkeypatch, payloads): + """Patch youtube.requests.get to return queued JSON payloads in order.""" + queue = [FakeResponse(json_data=p) for p in payloads] + calls = [] + + def _fake_get(url, params=None, **_kwargs): + calls.append({"url": url, "params": params or {}}) + return queue.pop(0) + + monkeypatch.setattr(youtube.requests, "get", _fake_get) + return calls + + +def test_resolve_api_key_prefers_argument(monkeypatch): + monkeypatch.delenv("YOUTUBE_DEVELOPER_KEY", raising=False) + assert youtube._resolve_api_key("explicit-key") == "explicit-key" + + +def test_resolve_api_key_falls_back_to_env(monkeypatch): + monkeypatch.setenv("YOUTUBE_DEVELOPER_KEY", "env-key") + assert youtube._resolve_api_key(None) == "env-key" + + +def test_resolve_api_key_raises_when_missing(monkeypatch): + monkeypatch.delenv("YOUTUBE_DEVELOPER_KEY", raising=False) + with pytest.raises(ValueError, match="YOUTUBE_DEVELOPER_KEY"): + youtube._resolve_api_key(None) + + +def test_github_headers_includes_token_when_present(): + assert youtube._github_headers(None) == {"Accept": "application/vnd.github.v3+json"} + assert youtube._github_headers("tok")["Authorization"] == "Bearer tok" + + +def test_batched_splits_into_chunks(): + assert list(youtube._batched(range(5), 2)) == [[0, 1], [2, 3], [4]] + assert list(youtube._batched([], 2)) == [] + + +def test_yt_paged_items_follows_next_page_token(monkeypatch): + calls = _queue_get( + monkeypatch, + [ + {"items": [{"id": "a"}], "nextPageToken": "PAGE2"}, + {"items": [{"id": "b"}]}, + ], + ) + items = list(youtube._yt_paged_items("videos", {"part": "id"}, "k")) + assert [i["id"] for i in items] == ["a", "b"] + assert calls[0]["params"]["key"] == "k" + assert calls[0]["params"]["maxResults"] == youtube.YOUTUBE_MAX_RESULTS + assert calls[1]["params"]["pageToken"] == "PAGE2" + + +def test_playlist_ids_explicit_configs_skip_ignored(monkeypatch): + def _boom(*_args, **_kwargs): + pytest.fail("should not call the API for explicit playlists") + + monkeypatch.setattr(youtube.requests, "get", _boom) + channel_config = { + "channel_id": "chan", + "playlists": [{"id": "keep"}, {"id": "drop", "ignore": True}], + } + assert list(youtube._playlist_ids_for_config(channel_config, "k")) == ["keep"] + + +def test_playlist_ids_wildcard_lists_channel_playlists(monkeypatch): + _queue_get(monkeypatch, [{"items": [{"id": "p1"}, {"id": "p2"}]}]) + channel_config = {"channel_id": "chan", "playlists": [{"id": "all"}]} + assert list(youtube._playlist_ids_for_config(channel_config, "k")) == ["p1", "p2"] + + +def test_playlist_ids_empty_config_is_wildcard(monkeypatch): + _queue_get(monkeypatch, [{"items": [{"id": "p1"}]}]) + assert list(youtube._playlist_ids_for_config({"channel_id": "c"}, "k")) == ["p1"] + + +def test_video_ids_for_playlist_reads_content_details(monkeypatch): + _queue_get( + monkeypatch, + [ + { + "items": [ + {"contentDetails": {"videoId": "v1"}}, + {"contentDetails": {}}, + {"contentDetails": {"videoId": "v2"}}, + ] + } + ], + ) + assert list(youtube._video_ids_for_playlist("pl", "k")) == ["v1", "v2"] + + +def test_channel_video_ids_deduplicates(monkeypatch): + _queue_get( + monkeypatch, + [ + {"items": [{"id": "pl1"}, {"id": "pl2"}]}, + { + "items": [ + {"contentDetails": {"videoId": "v1"}}, + {"contentDetails": {"videoId": "v2"}}, + ] + }, + { + "items": [ + {"contentDetails": {"videoId": "v2"}}, + {"contentDetails": {"videoId": "v3"}}, + ] + }, + ], + ) + assert list(youtube._channel_video_ids([{"channel_id": "c"}], "k")) == [ + "v1", + "v2", + "v3", + ] + + +def _install_fake_transcript_api(monkeypatch): + """Inject a fake youtube-transcript-api matching the 1.x fetch() interface.""" + + class _Fetched: + def to_raw_data(self): + return [{"text": "hello", "start": 0.0, "duration": 1.0}] + + class _Api: + def fetch(self, _video_id): + return _Fetched() + + fake: Any = types.ModuleType("youtube_transcript_api") + fake.YouTubeTranscriptApi = _Api + fake.NoTranscriptFound = type("NoTranscriptFound", (Exception,), {}) + fake.TranscriptsDisabled = type("TranscriptsDisabled", (Exception,), {}) + fake.VideoUnavailable = type("VideoUnavailable", (Exception,), {}) + + fmt_mod: Any = types.ModuleType("youtube_transcript_api.formatters") + + class _TextFormatter: + def format_transcript(self, _fetched): + return "hello" + + fmt_mod.TextFormatter = _TextFormatter + monkeypatch.setitem(sys.modules, "youtube_transcript_api", fake) + monkeypatch.setitem(sys.modules, "youtube_transcript_api.formatters", fmt_mod) + + +def test_transcripts_resource_yields_formatted_text(monkeypatch): + _install_fake_transcript_api(monkeypatch) + monkeypatch.setattr(youtube, "_resolve_api_key", lambda _key: "k") + monkeypatch.setattr(youtube, "_fetch_channel_configs", lambda **_kw: [{"c": 1}]) + monkeypatch.setattr( + youtube, "_channel_video_ids", lambda _configs, _key: iter(["vid1"]) + ) + + source = youtube.youtube_source(api_key="k") + rows = list(source.resources["raw__youtube__api__transcripts"]) + assert rows == [ + { + "video_id": "vid1", + "etl_source": "youtube", + "transcript": "hello", + "segments": [{"text": "hello", "start": 0.0, "duration": 1.0}], + } + ] + + +def _fake_channel_get(url, params=None, **_kwargs): + """Serve GitHub config listing + a single channel from the Data API.""" + if "contents" in url: + return FakeResponse( + json_data=[{"name": "c.yaml", "url": "https://api/file/c.yaml"}] + ) + if "api/file" in url: + yaml_bytes = b"channel_id: CHAN\noffered_by: ocw\n" + return FakeResponse( + json_data={"content": base64.b64encode(yaml_bytes).decode("ascii")} + ) + if "youtube/v3/channels" in url: + return FakeResponse( + json_data={"items": [{"id": "CHAN", "snippet": {"title": "A channel"}}]} + ) + return FakeResponse(json_data={"items": []}) + + +def test_channels_resource_builds_record(monkeypatch): + monkeypatch.setattr(youtube.requests, "get", _fake_channel_get) + source = youtube.youtube_source(api_key="k") + rows = list(source.resources["raw__youtube__api__channels"]) + assert len(rows) == 1 + assert rows[0]["channel_id"] == "CHAN" + assert rows[0]["offered_by"] == "ocw" + assert rows[0]["etl_source"] == "youtube" + assert rows[0]["snippet"]["title"] == "A channel" + + +@pytest.mark.integration +def test_youtube_channels_materialization(test_profile: Path, monkeypatch): + monkeypatch.setattr(youtube.requests, "get", _fake_channel_get) + pipeline = config.pipeline_for("youtube") + source = youtube.youtube_source(api_key="k").with_resources( + "raw__youtube__api__channels" + ) + info = pipeline.run(source) + assert not info.has_failed_jobs + + dataset = pipeline.dataset() + assert dataset["raw__youtube__api__channels"].arrow().num_rows == 1 diff --git a/src/ol_dlt/uv.lock b/src/ol_dlt/uv.lock index 1aa149068..dd759bf99 100644 --- a/src/ol_dlt/uv.lock +++ b/src/ol_dlt/uv.lock @@ -333,6 +333,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/07/6c/aa3f2f849e01cb6a001cd8554a88d4c77c5c1a31c95bdf1cf9301e6d9ef4/defusedxml-0.7.1-py2.py3-none-any.whl", hash = "sha256:a352e7e428770286cc899e2542b6cdaedb2b4953ff269a210103ec58f6198a61", size = 25604, upload-time = "2021-03-08T10:59:24.45Z" }, ] +[[package]] +name = "defusedxml" +version = "0.7.1" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/0f/d5/c66da9b79e5bdb124974bfe172b4daf3c984ebd9c2a06e2b8a4dc7331c72/defusedxml-0.7.1.tar.gz", hash = "sha256:1bb3032db185915b62d7c6209c5a8792be6a32ab2fedacc84e01b52c51aa3e69", size = 75520, upload-time = "2021-03-08T10:59:26.269Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/07/6c/aa3f2f849e01cb6a001cd8554a88d4c77c5c1a31c95bdf1cf9301e6d9ef4/defusedxml-0.7.1-py2.py3-none-any.whl", hash = "sha256:a352e7e428770286cc899e2542b6cdaedb2b4953ff269a210103ec58f6198a61", size = 25604, upload-time = "2021-03-08T10:59:24.45Z" }, +] + [[package]] name = "dlt" version = "1.28.1" @@ -778,6 +787,7 @@ dependencies = [ { name = "duckdb" }, { name = "pyiceberg", extra = ["glue"] }, { name = "requests" }, + { name = "youtube-transcript-api" }, ] [package.dev-dependencies] @@ -794,6 +804,7 @@ requires-dist = [ { name = "duckdb", specifier = ">=1.0,!=1.5.0,!=1.5.1" }, { name = "pyiceberg", extras = ["glue"], specifier = ">=0.9" }, { name = "requests", specifier = ">=2.32" }, + { name = "youtube-transcript-api", specifier = ">=1.0,<2" }, ] [package.metadata.requires-dev] @@ -1714,6 +1725,19 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/fd/4d/4b880086bd0d3e034d25647be1d830afc3e3f610e98c4ab3490af6b1b6d5/yarl-1.24.2-py3-none-any.whl", hash = "sha256:2783d9226db8797636cd6896e4de81feed252d1db72265686c9558d97a4d94b9", size = 53576, upload-time = "2026-05-19T21:31:03.909Z" }, ] +[[package]] +name = "youtube-transcript-api" +version = "1.2.4" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "defusedxml" }, + { name = "requests" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/60/43/4104185a2eaa839daa693b30e15c37e7e58795e8e09ec414f22b3db54bec/youtube_transcript_api-1.2.4.tar.gz", hash = "sha256:b72d0e96a335df599d67cee51d49e143cff4f45b84bcafc202ff51291603ddcd", size = 469839, upload-time = "2026-01-29T09:09:17.088Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/be/95/129ea37efd6cd6ed00f62baae6543345c677810b8a3bf0026756e1d3cf3c/youtube_transcript_api-1.2.4-py3-none-any.whl", hash = "sha256:03878759356da5caf5edac77431780b91448fb3d8c21d4496015bdc8a7bc43ff", size = 485227, upload-time = "2026-01-29T09:09:15.427Z" }, +] + [[package]] name = "zstandard" version = "0.25.0" From 5e13c2afd2de5ee2727fecd336aa92aa9994dcb9 Mon Sep 17 00:00:00 2001 From: Tobias Macey Date: Wed, 1 Jul 2026 09:14:48 -0400 Subject: [PATCH 2/3] fix(youtube): parse list-format configs, default to main branch MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Address PR review: the mitodl/open-video-data youtube/*.yaml files (channels.yaml, shorts.yaml) are YAML *lists* of channel dicts on the `main` branch, not one dict per file on `master` — the previous parser loaded zero configs against the real repo. Parse both list and single-dict files, default github_branch to "main", and accept bare-string playlist entries so `playlists: [all]` still triggers the wildcard instead of silently ingesting nothing. Update the README format docs and tests to match the real list-based config shape. Co-Authored-By: Claude Opus 4.8 --- src/ol_dlt/ol_dlt/sources/youtube/README.md | 16 ++++-- src/ol_dlt/ol_dlt/sources/youtube/__init__.py | 49 +++++++++++-------- src/ol_dlt/tests/sources/test_youtube.py | 45 +++++++++++++++-- 3 files changed, 83 insertions(+), 27 deletions(-) diff --git a/src/ol_dlt/ol_dlt/sources/youtube/README.md b/src/ol_dlt/ol_dlt/sources/youtube/README.md index 22b9bf786..70aa0a76f 100644 --- a/src/ol_dlt/ol_dlt/sources/youtube/README.md +++ b/src/ol_dlt/ol_dlt/sources/youtube/README.md @@ -6,9 +6,19 @@ application. Wrapped as Dagster assets by the `data_loading` code location (`data_loading/defs/ingestion/assets.py`). **Config source:** [mitodl/open-video-data](https://github.com/mitodl/open-video-data) -(`youtube/*.yaml`, one file per channel). Each config provides a `channel_id` -and an optional `playlists` list (`{id, ignore?}`). A config with no `playlists`, -or one containing the `all` wildcard id, ingests every playlist on the channel. +(`youtube/*.yaml` — e.g. `channels.yaml`, `shorts.yaml` — on the `main` branch). +Each file is a **YAML list** of channel config dicts; add or edit a channel by +adding a `- channel_id: ...` entry. Each entry provides a `channel_id` and an +optional `playlists` list (entries may be `{id, ignore?}` dicts or bare id +strings). A channel with no `playlists`, or one containing the `all` wildcard id, +ingests every playlist on the channel. + +```yaml +- channel_id: UCEBb1b_L6zDS3xTUrIALZOw + offered_by: ocw + playlists: + - id: all +``` ## Data flow diff --git a/src/ol_dlt/ol_dlt/sources/youtube/__init__.py b/src/ol_dlt/ol_dlt/sources/youtube/__init__.py index 2d9b25b01..d296f9858 100644 --- a/src/ol_dlt/ol_dlt/sources/youtube/__init__.py +++ b/src/ol_dlt/ol_dlt/sources/youtube/__init__.py @@ -59,10 +59,13 @@ def _fetch_channel_configs( ) -> list[dict[str, Any]]: """Fetch and parse all YouTube channel YAML config files from a GitHub repo. - Each config must contain a ``channel_id`` to be included. The ``playlists`` - key (a list of ``{id, ignore?}`` dicts) is optional; a config with no - ``playlists`` is treated as a wildcard (ingest all of the channel's - playlists). + The mitodl/open-video-data ``youtube/`` folder holds a few YAML files + (``channels.yaml``, ``shorts.yaml``) each containing a *list* of channel + config dicts (``- channel_id: ...``). A single top-level dict is also + accepted for robustness. Each channel dict must contain a ``channel_id``; + its ``playlists`` key (a list of ``{id, ignore?}`` dicts or bare id + strings) is optional and a channel with no ``playlists`` is treated as a + wildcard (ingest all of the channel's playlists). """ headers = _github_headers(token) listing_url = f"{GITHUB_API_BASE}/repos/{repo}/contents/{folder}?ref={branch}" @@ -79,23 +82,21 @@ def _fetch_channel_configs( raw_content = base64.b64decode(file_resp.json()["content"]).decode("utf-8") try: - channel_config = yaml.safe_load(raw_content) + parsed = yaml.safe_load(raw_content) except yaml.YAMLError: logger.exception("Failed to parse YAML config: %s", file_meta["name"]) continue - if not isinstance(channel_config, dict): - logger.warning("Skipping non-dict youtube config: %s", file_meta["name"]) - continue - - if "channel_id" not in channel_config: - logger.warning( - "Skipping config missing required key (channel_id): %s", - file_meta["name"], - ) - continue - - configs.append(channel_config) + # A file is either a single channel dict or a list of channel dicts. + entries = parsed if isinstance(parsed, list) else [parsed] + for entry in entries: + if not isinstance(entry, dict) or "channel_id" not in entry: + logger.warning( + "Skipping youtube config entry without channel_id in %s", + file_meta["name"], + ) + continue + configs.append(entry) logger.info("Loaded %d youtube configs from %s/%s", len(configs), repo, folder) return configs @@ -144,9 +145,15 @@ def _playlist_ids_for_config( """ channel_id = channel_config["channel_id"] playlist_configs = channel_config.get("playlists") or [] - configs_by_id = { - pc["id"]: pc for pc in playlist_configs if isinstance(pc, dict) and "id" in pc - } + # Each entry is either {"id": ..., "ignore"?: ...} or a bare id string; + # accept both so a config like `playlists: [all]` still triggers the + # wildcard instead of silently ingesting nothing. + configs_by_id: dict[str, dict[str, Any]] = {} + for pc in playlist_configs: + if isinstance(pc, str): + configs_by_id[pc] = {"id": pc} + elif isinstance(pc, dict) and "id" in pc: + configs_by_id[pc["id"]] = pc if not playlist_configs or WILDCARD_PLAYLIST_ID in configs_by_id: for playlist in _yt_paged_items( @@ -203,7 +210,7 @@ def youtube_source( # noqa: C901 github_access_token: str | None = None, github_repo: str = _CONFIG_FILE_REPO_DEFAULT, github_folder: str = _CONFIG_FILE_FOLDER_DEFAULT, - github_branch: str = "master", + github_branch: str = "main", ) -> Generator[Any]: """Load MIT Learn YouTube data from the Data API v3 into four raw tables. diff --git a/src/ol_dlt/tests/sources/test_youtube.py b/src/ol_dlt/tests/sources/test_youtube.py index 621753df7..b852445cc 100644 --- a/src/ol_dlt/tests/sources/test_youtube.py +++ b/src/ol_dlt/tests/sources/test_youtube.py @@ -180,16 +180,23 @@ def test_transcripts_resource_yields_formatted_text(monkeypatch): ] +# Matches the real mitodl/open-video-data format: a YAML *list* of channel dicts. +_CHANNELS_YAML = ( + b"---\n- channel_id: CHAN\n offered_by: ocw\n playlists:\n - id: all\n" +) + + def _fake_channel_get(url, params=None, **_kwargs): """Serve GitHub config listing + a single channel from the Data API.""" if "contents" in url: return FakeResponse( - json_data=[{"name": "c.yaml", "url": "https://api/file/c.yaml"}] + json_data=[ + {"name": "channels.yaml", "url": "https://api/file/channels.yaml"} + ] ) if "api/file" in url: - yaml_bytes = b"channel_id: CHAN\noffered_by: ocw\n" return FakeResponse( - json_data={"content": base64.b64encode(yaml_bytes).decode("ascii")} + json_data={"content": base64.b64encode(_CHANNELS_YAML).decode("ascii")} ) if "youtube/v3/channels" in url: return FakeResponse( @@ -198,6 +205,38 @@ def _fake_channel_get(url, params=None, **_kwargs): return FakeResponse(json_data={"items": []}) +def test_fetch_channel_configs_parses_yaml_list(monkeypatch): + """A YAML file that is a list of channel dicts loads every valid entry.""" + yaml_bytes = ( + b"---\n" + b"- channel_id: CHAN_A\n offered_by: ocw\n" + b"- channel_id: CHAN_B\n" + b"- offered_by: no_channel_id\n" # invalid: skipped + ) + + def _fake_get(url, params=None, **_kwargs): + if "contents" in url: + return FakeResponse( + json_data=[{"name": "channels.yaml", "url": "https://api/f.yaml"}] + ) + return FakeResponse( + json_data={"content": base64.b64encode(yaml_bytes).decode("ascii")} + ) + + monkeypatch.setattr(youtube.requests, "get", _fake_get) + configs = youtube._fetch_channel_configs( + repo="mitodl/open-video-data", folder="youtube", branch="main", token=None + ) + assert [c["channel_id"] for c in configs] == ["CHAN_A", "CHAN_B"] + + +def test_playlist_ids_bare_string_wildcard(monkeypatch): + """A bare-string ``all`` entry (not ``{id: all}``) still triggers the wildcard.""" + _queue_get(monkeypatch, [{"items": [{"id": "p1"}, {"id": "p2"}]}]) + channel_config = {"channel_id": "chan", "playlists": ["all"]} + assert list(youtube._playlist_ids_for_config(channel_config, "k")) == ["p1", "p2"] + + def test_channels_resource_builds_record(monkeypatch): monkeypatch.setattr(youtube.requests, "get", _fake_channel_get) source = youtube.youtube_source(api_key="k") From 12325621911a8c6aa3fe549260b39f50264e0461 Mon Sep 17 00:00:00 2001 From: Tobias Macey Date: Fri, 10 Jul 2026 21:01:25 -0400 Subject: [PATCH 3/3] fix(ol_dlt): align youtube source with #2424 hardening pattern Use dlt.sources.helpers.requests (retry/backoff) instead of plain requests, and apply JSON_API_SCHEMA_CONTRACT to all four resources, matching mit_climate/mit_edx_programs/mitpe. Also dedupes a duplicate defusedxml package block left in uv.lock by the earlier rebase's textual auto-merge. --- src/ol_dlt/ol_dlt/sources/youtube/__init__.py | 6 +++++- src/ol_dlt/uv.lock | 9 --------- 2 files changed, 5 insertions(+), 10 deletions(-) diff --git a/src/ol_dlt/ol_dlt/sources/youtube/__init__.py b/src/ol_dlt/ol_dlt/sources/youtube/__init__.py index d296f9858..02aa75f36 100644 --- a/src/ol_dlt/ol_dlt/sources/youtube/__init__.py +++ b/src/ol_dlt/ol_dlt/sources/youtube/__init__.py @@ -27,8 +27,8 @@ from typing import Any import dlt -import requests import yaml +from dlt.sources.helpers import requests from ol_dlt import config @@ -250,6 +250,7 @@ def _configs() -> list[dict[str, Any]]: write_disposition="merge", primary_key="channel_id", table_format=table_format, + schema_contract=config.JSON_API_SCHEMA_CONTRACT, ) def youtube_channels() -> Generator[dict[str, Any]]: """Yield one record per configured YouTube channel.""" @@ -278,6 +279,7 @@ def youtube_channels() -> Generator[dict[str, Any]]: write_disposition="merge", primary_key="playlist_id", table_format=table_format, + schema_contract=config.JSON_API_SCHEMA_CONTRACT, ) def youtube_playlists() -> Generator[dict[str, Any]]: """Yield one record per ingested playlist across all channels.""" @@ -307,6 +309,7 @@ def youtube_playlists() -> Generator[dict[str, Any]]: write_disposition="merge", primary_key="video_id", table_format=table_format, + schema_contract=config.JSON_API_SCHEMA_CONTRACT, ) def youtube_videos() -> Generator[dict[str, Any]]: """Yield one record per video across all configured playlists.""" @@ -328,6 +331,7 @@ def youtube_videos() -> Generator[dict[str, Any]]: write_disposition="merge", primary_key="video_id", table_format=table_format, + schema_contract=config.JSON_API_SCHEMA_CONTRACT, ) def youtube_transcripts() -> Generator[dict[str, Any]]: """Yield one record per video that has a transcript. diff --git a/src/ol_dlt/uv.lock b/src/ol_dlt/uv.lock index dd759bf99..709ad0af5 100644 --- a/src/ol_dlt/uv.lock +++ b/src/ol_dlt/uv.lock @@ -333,15 +333,6 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/07/6c/aa3f2f849e01cb6a001cd8554a88d4c77c5c1a31c95bdf1cf9301e6d9ef4/defusedxml-0.7.1-py2.py3-none-any.whl", hash = "sha256:a352e7e428770286cc899e2542b6cdaedb2b4953ff269a210103ec58f6198a61", size = 25604, upload-time = "2021-03-08T10:59:24.45Z" }, ] -[[package]] -name = "defusedxml" -version = "0.7.1" -source = { registry = "https://pypi.org/simple" } -sdist = { url = "https://files.pythonhosted.org/packages/0f/d5/c66da9b79e5bdb124974bfe172b4daf3c984ebd9c2a06e2b8a4dc7331c72/defusedxml-0.7.1.tar.gz", hash = "sha256:1bb3032db185915b62d7c6209c5a8792be6a32ab2fedacc84e01b52c51aa3e69", size = 75520, upload-time = "2021-03-08T10:59:26.269Z" } -wheels = [ - { url = "https://files.pythonhosted.org/packages/07/6c/aa3f2f849e01cb6a001cd8554a88d4c77c5c1a31c95bdf1cf9301e6d9ef4/defusedxml-0.7.1-py2.py3-none-any.whl", hash = "sha256:a352e7e428770286cc899e2542b6cdaedb2b4953ff269a210103ec58f6198a61", size = 25604, upload-time = "2021-03-08T10:59:24.45Z" }, -] - [[package]] name = "dlt" version = "1.28.1"