Databricks Free Editionで始めるRedditデータパイプライン② データ取得の基本設計
はじめに
前回はパイプライン全体のアーキテクチャとDatabricks Asset Bundlesの設定を解説しました。
今回は、データ取得部分の基本設計を解説します。Inoreader API、PDF生成、pdfplumberという3つのコンポーネントがどのような役割を担い、どのようにデータが流れていくのかを俯瞰します。
実装の詳細なロジックよりも、なぜこの設計にしたのかという思想面に焦点を当てます。
データ取得の全体フロー
┌─────────────────────────────────────────────────────────────────┐
│ データ取得フロー │
├─────────────────────────────────────────────────────────────────┤
│ │
│ Step 1: Inoreader API │
│ ┌─────────────┐ │
│ │ Reddit RSS │──▶ fetch_inoreader.py ──▶ reddit_items_*.csv │
│ │ (メタデータ) │ (OAuth 2.0) (URL, title, etc) │
│ └─────────────┘ │
│ │ │
│ ▼ │
│ Step 2: PDF生成 │
│ ┌─────────────┐ │
│ │ Reddit Page │──▶ scraper.py ──▶ PDF + metadata.json │
│ │ (実ページ) │ (紐付け情報) │
│ └─────────────┘ │
│ │ │
│ ▼ │
│ Step 3: pdfplumber │
│ ┌─────────────┐ │
│ │ PDF │──▶ pdf_to_csv.py ──▶ reddit_results_*.csv │
│ │ (静的ファイル)│ (座標ベース抽出) (upvote, comment, etc) │
│ └─────────────┘ │
│ │ │
│ ▼ │
│ Step 4: Databricks Upload │
│ ┌─────────────┐ │
│ │ results.csv │──▶ databricks fs cp ──▶ UC Volume │
│ └─────────────┘ │
│ │
└─────────────────────────────────────────────────────────────────┘
なぜ4ステップに分けるのか?
各ステップの特性が異なるため、分離することで運用が安定します。
| ステップ | 特性 | 失敗時の影響 |
|---|---|---|
| Step 1: 一覧取得 | 軽量・高速 | 後続すべてが止まる |
| Step 2: PDF生成 | 重い・時間がかかる | 該当URLのみ再実行 |
| Step 3: 抽出 | 重い | 該当フォルダのみ再実行 |
| Step 4: アップロード | シンプル | リトライで解決 |
失敗したステップだけを再実行できるので、全体をやり直す必要がありません。
なぜこの3段(+アップロード)構成なのか
Reddit APIを使わない理由
Reddit公式APIを使う方法もありますが、以下の理由で今回は採用しませんでした。
- レート制限が厳しい — 大量のデータ取得には向かない
- API申請のハードルが高い — 審査プロセスがあり、個人プロジェクトでは通りにくい
採用したアプローチ
代わりに、以下の構成を採用しました。
| やりたいこと | 解決策 |
|---|---|
| 記事一覧を安定して取得したい | Inoreader(RSSアグリゲーター)経由でメタデータを取得 |
| upvote数などの数値を取得したい | ページをPDF化し、そこから数値を抽出 |
| 仕様変更の影響を最小化したい | 「一覧取得」と「ページ取得」を分離し、変更があった側だけ修正 |
つまり、一覧は軽いAPIで効率よく集め、数値が必要なページだけを個別に取得するという段階的なアプローチです。
Step 1: Inoreader API — メタデータ取得
役割
InoreaderはRSSリーダーサービスで、Redditの各subredditをRSSフィードとして購読できます。APIを通じて購読記事のメタデータ(URL、タイトル、公開日時など)を取得します。
⚠️ Inoreader APIの利用条件
Inoreader APIを利用するにはProプラン以上の契約が必要です。無料プランではAPIを使用できません。
詳細はInoreader Developer Portalで確認してください。
設計ポイント
Inoreader
├── redditラベル(フォルダ)
│ ├── r/dataengineering (RSS)
│ ├── r/MachineLearning (RSS)
│ ├── r/python (RSS)
│ └── ...
重複排除の仕組み(seen_urls.json)
# seen_urls.json による重複管理
{
"entries": [
{
"url": "https://reddit.com/r/.../comments/xxx",
"first_seen": "2025-01-01T10:00:00+00:00"
}
]
}
- 一度取得したURLは
seen_urls.jsonに記録 - 90日間のTTL(Time To Live)で古いエントリを自動削除
- これにより、同じ記事を重複して処理することを防止
出力: reddit_items_*.csv
| カラム | 説明 | 例 |
|---|---|---|
| url | Reddit投稿URL | https://reddit.com/r/.../comments/xxx |
| title | 投稿タイトル | How to optimize Spark jobs? |
| categories | カテゴリ(パイプ区切り) | user/123|label/reddit |
| origin_title | subreddit名 | r/dataengineering |
| summary | 投稿サマリ(HTML) | <p>Looking for tips...</p> |
| published_at | 公開日時(ISO8601) | 2025-12-01T10:30:00+00:00 |
| crawlTimeMsec | クロール日時 | 1733050200000 |
Step 2: PDF生成
役割
Inoreaderで取得したURLからページをPDFとして保存します。これにより、upvote数などの数値を後から抽出できます。
なぜPDFなのか
- スクリーンショット(画像)ではOCRが必要になり精度が落ちる
- PDFならテキスト情報がそのまま保持される
- pdfplumberで座標ベースの抽出が可能
出力: PDF + metadata.json
data/pdf/20250101_120000/
├── 001_How_to_optimize_Spark_jobs.pdf
├── 002_Best_practices_for_dbt.pdf
└── metadata.json # PDFとitemsデータの紐付け情報
metadata.json の例:
{
"created_at": "2025-01-01T12:00:00+00:00",
"items": [
{
"url": "https://reddit.com/r/.../comments/xxx",
"title": "How to optimize Spark jobs?",
"pdf_filename": "001_How_to_optimize_Spark_jobs.pdf",
"processed_at": "2025-01-01T12:00:05+00:00"
}
]
}
この metadata.json により、PDF抽出結果と元のメタデータをファイル名ベースでマージできます。Databricks側で複雑なJOINを書かなくても、ローカルで整理してからアップロードできます。
状態管理
state/reddit/
├── auth.json # 認証状態
└── processed_urls.json # 処理済みURL(90日TTL)
- 処理済みURLは
processed_urls.jsonに記録し、重複処理を防止
Step 3: pdfplumber — PDF解析(座標ベース抽出)
役割
PDFからupvote数、トップコメント、コメントupvoteを抽出し、CSVに出力します。
設計ポイント
pdfplumberはPDF内のテキストやレイアウト情報(文字単位の位置など)を扱えるPythonライブラリです。ブラウザで生成したPDFのようにテキスト情報が埋め込まれているPDFに適しています(スキャンした画像PDFには向きません)。
座標ベース抽出のイメージ
┌─────────────────────────────────────────┐
│ [投稿タイトル] │
│ │
│ ▲ 1,234 💬 56 🔗 共有 │ ← upvote数はここ
│ ↑ │
│ この数値を座標ベースで抽出 │
├─────────────────────────────────────────┤
│ [コメント欄] │
│ │
│ user123 • 3日前 │
│ This is a great post... │
│ ▲ 456 返信 │ ← コメントupvote
│ │
└─────────────────────────────────────────┘
抽出ロジック(概要)
- 投稿upvote — 「共有」ボタンの左側にある数値を検出
-
トップコメント —
username • N日前のような開始パターンでコメントを検出 - コメントupvote — 「返信」ボタンの左側にある数値を検出
出力: reddit_results_*.csv
| カラム | 説明 | 例 |
|---|---|---|
| url | Reddit投稿URL | https://reddit.com/r/.../comments/xxx |
| title | 投稿タイトル | How to optimize Spark jobs? |
| upvote | 投稿のupvote数 | 1234 |
| comment | トップコメント本文 | You should try... |
| comment_upvote | コメントのupvote数 | 456 |
| date | 取り込み日(YYYYMMDD) | 20251201 |
| (その他) | itemsデータからマージ | categories, origin_title, etc |
Step 4: Databricks Upload — Unity Catalog Volumeへ配置
役割
生成されたCSVをDatabricksのUnity Catalog Volumeにアップロードします。
databricks fs cp \
"data/csv/results/reddit_results_20250101_120000.csv" \
"dbfs:/Volumes/reddit/sdp/landing/results/"
Volumesはファイルをガバナンス対象にできるUnity Catalogオブジェクトで、/Volumes/<catalog>/<schema>/<volume>/... 形式でアクセスします。
統合実行スクリプト(run_pipeline.sh)
4つのステップを一括実行するシェルスクリプト例です。
#!/bin/bash
# run_pipeline.sh
# Step 1: Inoreader から記事取得
uv run python src/ingestion/inoreader/fetch_inoreader.py --max-items 200
# Step 2: PDF生成
uv run python src/ingestion/reddit/scraper.py --week-ago
# Step 3: PDF から統計情報を抽出
uv run python src/ingestion/reddit/pdf_to_csv.py "$LATEST_PDF_DIR"
# Step 4: Databricks にアップロード
databricks fs cp "$LATEST_RESULTS_FILE" "dbfs:/Volumes/reddit/sdp/landing/results/"
実行オプション例:
# 通常実行(1週間前の記事をPDF化)
./run_pipeline.sh
# 全記事を処理(日付フィルタなし)
./run_pipeline.sh --all
# Databricksアップロードをスキップ
./run_pipeline.sh --no-upload
# CSV変換のみ実行
./run_pipeline.sh --csv-only
# アップロードのみ実行
./run_pipeline.sh --upload-only
処理対象がなければ後続をスキップ:
# 処理対象がなかった場合はStep 3以降をスキップ
if grep -q "処理対象のURLがありません" "$SCRAPER_OUTPUT_FILE"; then
echo "⚠ 処理対象のURLがないため、Step 3以降をスキップします。"
exit 0
fi
データ保持ポリシー(TTL)
各コンポーネントで管理する状態ファイルには90日のTTLを設定します。
| データ | TTL | 管理ファイル |
|---|---|---|
| 取得済みURL(Inoreader) | 90日 | state/inoreader/seen_urls.json |
| 処理済みURL(Scraper) | 90日 | state/reddit/processed_urls.json |
| PDFファイル | 手動削除 | data/pdf/ |
| CSVファイル | 手動削除 | data/csv/ |
TTLを超えた古いURLは削除され、再取得・再処理の対象になります。これにより、長期運用でも状態ファイルが肥大化しません。
今回のまとめ
第2回では、データ取得部分の基本設計を解説しました。
- 4段構成の理由 — 一覧取得・PDF生成・抽出・アップロードを分離し、失敗時の影響を局所化
- Inoreader — RSSアグリゲーターとしてメタデータ取得を担当(Proプラン必要)
- PDF生成 — ページをPDF化し、数値抽出の元データを保存
- pdfplumber — 座標ベースでupvote/コメント情報を抽出
- metadata.json — PDFとitemsデータのマージを簡素化
次回予告
第3回では、Databricks(DLT / Lakeflow Declarative Pipelines)でBronze→Silver→Gold層を構築する流れを解説します。
シリーズ目次:
- 概説 & プロジェクト準備
- データ取得の基本設計(本記事)
- Databricks DLT(Bronze→Silver→Gold)
- ダッシュボード & 運用
Discussion