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
│                                         │
└─────────────────────────────────────────┘

抽出ロジック(概要)

  1. 投稿upvote — 「共有」ボタンの左側にある数値を検出
  2. トップコメントusername • N日前 のような開始パターンでコメントを検出
  3. コメント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層を構築する流れを解説します。


シリーズ目次:

  1. 概説 & プロジェクト準備
  2. データ取得の基本設計(本記事)
  3. Databricks DLT(Bronze→Silver→Gold)
  4. ダッシュボード & 運用

Discussion