YouTube数据管道
Apache Airflow管道,根据指定主题提取YouTube视频数据,检索视频转录,并将数据保存在DuckDB中,以进行高效查询和分析。
特性
- 基于主题的搜索:使用自定义主题搜索YouTube视频
- 全面数据收集:收集视频URL、标题、描述和文字记录
- 转录提取:自动获取可用的视频记录
- 永久存储:将所有数据存储在DuckDB中,以便快速查询和分析
- 简单架构:在没有复杂中间件的情况下直接调用API
技术栈
- 阿帕奇气流:工作流编排和调度
- DuckDB:高性能分析数据库
- YouTube数据API v3:用于视频搜索和元数据
- YouTube转录API:用于转录提取
- python:核心编程语言
项目结构
antigravity-mcp-airflow-duckdb-pipeline-pull-data-from-youtube/
├── dags/
│ └── youtube_pipeline.py # Main Airflow DAG
├── README.md # This file
└── .gitignore # Git ignore rules关键组件
- 气流DAG(
dags/youtube_pipeline.py):
- 编排数据管道工作流 - 包含两个主要任务:搜索和处理视频 - 处理主题输入和数据持久化
- DuckDB数据库:
- 存储视频元数据和转录本 - 针对分析查询进行了优化 - 基于文件的数据库,便于部署
DAG架构和工作流
DAG结构
youtube_data_pipeline (DAG)
├── search_videos (@task)
│ └── Searches YouTube for videos
└── process_videos (@task)
└── Retrieves transcripts and stores in DuckDBDAG的工作原理
- 初始化:
- 设置与DuckDB的数据库连接 - 创造 videos 如果表不存在 - 配置API身份验证
- 搜索阶段 (
search_task):
- 接受主题输入(默认:“机器学习”) - 调用YouTube搜索API - 返回视频元数据列表(ID、标题、描述、URL)
- 处理阶段 (
process_task):
- 从搜索任务接收视频列表 - 对于每个视频: - 使用YouTube转录API获取转录 - 将完整数据存储在DuckDB中 - 优雅地处理错误(例如,没有转录的视频)
- 数据持久层:
- 所有存储在DuckDB中的数据都有模式:
CREATE TABLE videos (
video_id VARCHAR PRIMARY KEY,
title VARCHAR,
description VARCHAR,
url VARCHAR,
transcript VARCHAR,
created_at TIMESTAMP
);进度安排和执行
- 日程:可配置(默认值:每日)
- 重试:1次重试,失败后延迟5分钟
- 依赖项:顺序执行(搜索→ 过程)
使用DuckDB进行数据分析
为什么选择DuckDB?
- 演出:对结构化数据进行快速分析查询
- 简洁:基于文件的数据库,不需要服务器
- SQL支持:用于复杂查询的标准SQL
- Python集成:用于数据操作的原生Python API
示例分析查询
-- Count videos by topic
SELECT COUNT(*) as video_count
FROM videos
WHERE description ILIKE '%machine learning%';
-- Find videos with longest transcripts
SELECT title, LENGTH(transcript) as transcript_length
FROM videos
ORDER BY transcript_length DESC
LIMIT 10;
-- Search transcripts for specific keywords
SELECT title, url
FROM videos
WHERE transcript ILIKE '%neural network%';
-- Analyze content trends
SELECT
DATE_TRUNC('day', created_at) as date,
COUNT(*) as videos_added
FROM videos
GROUP BY DATE_TRUNC('day', created_at)
ORDER BY date DESC;Python分析示例
import duckdb
# Connect to database
con = duckdb.connect('youtube_data.db')
# Load data into pandas
df = con.execute("SELECT * FROM videos").fetchdf()
# Analyze transcript lengths
df['transcript_length'] = df['transcript'].str.len()
print(df.groupby('topic')['transcript_length'].describe())先决条件
- Python 3.8+
- Apache Airflow 2.x
- YouTube数据API v3密钥(从 谷歌云控制台)
安装
- 克隆存储库:
git clone https://github.com/denis911/antigravity-mcp-airflow-duckdb-pipeline-pull-data-from-youtube.git
cd antigravity-mcp-airflow-duckdb-pipeline-pull-data-from-youtube- 安装Python依赖项:
pip install apache-airflow duckdb requests youtube-transcript-api- 将YouTube API密钥设置为气流变量:
# Via Airflow CLI
airflow variables set youtube_api_key "your_youtube_api_key_here"
# Via Airflow UI: Admin → Variables → Create Variable
# Key: youtube_api_key
# Value: your_youtube_api_key_here- 初始化气流(如果尚未完成):
airflow db init本地设置和运行
1.环境设置
创建虚拟环境并安装依赖项:
# Create virtual environment
python -m venv airflow_env
source airflow_env/bin/activate # On Windows: airflow_env\Scripts\activate
# Install dependencies
pip install apache-airflow duckdb requests youtube-transcript-api2.气流初始化
# Set AIRFLOW_HOME (optional, defaults to ~/airflow)
# Unix/macOS
export AIRFLOW_HOME=./airflow_home
# Windows (PowerShell)
$env:AIRFLOW_HOME="./airflow_home"
# Initialize Airflow database
airflow db init
# Create admin user
airflow users create \
--username admin \
--firstname Admin \
--lastname User \
--role Admin \
--email admin@example.com3.配置YouTube API
从获取API密钥 谷歌云控制台:
- 创建新项目或选择现有项目
- 启用YouTube数据API v3
- 创建凭据(API密钥)
- 设置气流变量:
# Via Airflow CLI
airflow variables set youtube_api_key "your_actual_api_key_here"
# Or via Airflow UI: Admin → Variables → Create Variable
# Key: youtube_api_key
# Value: your_actual_api_key_here4.设置DAG目录
# Copy DAG to Airflow dags folder
cp dags/youtube_pipeline.py $AIRFLOW_HOME/dags/
# Or set AIRFLOW__CORE__DAGS_FOLDER to current directory
# Unix/macOS
export AIRFLOW__CORE__DAGS_FOLDER=./dags
# Windows (PowerShell)
$env:AIRFLOW__CORE__DAGS_FOLDER="./dags"5.运行气流
# Start scheduler (runs DAGs)
airflow scheduler
# In another terminal, start webserver
airflow webserver --port 80806.访问气流UI
打开http://localhost:8080在浏览器中,使用以下命令登录:
- 用户名:admin
- 密码:(无论您在创建用户时设置了什么)
7.运行管道
- 在Airflow UI中,查找
youtube_data_pipeline有向无环图 - 单击切换以启用它
- 点击“触发DAG”立即运行
- 监控UI中的进度
8.检查结果
# Connect to DuckDB and query results
python -c "
import duckdb
con = duckdb.connect('./data/youtube_data.db')
result = con.execute('SELECT COUNT(*) as total_videos FROM videos').fetchone()
print(f'Total videos collected: {result[0]}')
con.close()
"用法
- 将DAG文件放置在Airflow dags文件夹中
- 更新DAG中的主题(当前设置为“机器学习”)
- 启动气流:
airflow scheduler &
airflow webserver &- 从气流UI或CLI触发DAG
配置
气流变量
youtube_api_key:YouTube API访问必需(通过Airflow UI或CLI设置)duckdb_path:DuckDB数据库的路径(默认值:/opt/aflow/data/youtube_data.db)
DAG参数
topic:搜索主题(默认:“机器学习”)max_results:要处理的视频数量(默认值:5)schedule:运行管道的频率(默认值:每天)
监控与维护
气流UI
- 监控DAG运行和任务状态
- 查看调试日志
- 手动触发测试运行
数据库维护
-- Check data volume
SELECT COUNT(*) FROM videos;
-- Remove old data if needed
DELETE FROM videos WHERE created_at < '2023-01-01';
-- Optimize database
VACUUM;API配额
- 监控YouTube API配额使用情况
- 优雅地处理速率限制
- 考虑升级API配额以供生产使用
贡献
- 分叉存储库
- 创建要素分支
- 进行更改
- 添加新功能的测试
- 提交拉取请求
故障排除
常见问题
- API密钥错误:确保
youtube_api_key气流变量设置正确(管理员→ 气流UI中的变量) - 成绩单不可用:有些视频没有文字记录
- 数据库连接:检查DuckDB文件权限和路径
- 气流未启动:验证Python路径和依赖关系
日志
检查UI或日志文件中的气流日志,了解详细的错误消息。
未来的增强功能
- 多主题支持:在单个DAG运行中处理多个主题
- 增量更新:仅获取自上次运行以来的新视频
- 高级分析:成绩单情感分析
- 仪表板集成:与Tableau等BI工具连接
- 通知系统:管道故障或新内容警报
免责声明
该项目用于教育和研究目的。请遵守YouTube的服务条款和API使用限制。
