Scrapyを活用したウェブコンテンツ抽出とEasysearchへのリアルタイムインデックス化

プロジェクト初期化と環境準備

ScrapyはPython向けの非同期基盤に基づく高機能スクレイピングフレームワークです。対象ドメインのDOM構造を効率的に走査し、構造化データを取得するための標準的な雛形生成は以下のCLIコマンドで完了します。

pip install scrapy
scrapy startproject article_indexer
cd article_indexer/spiders
scrapy genspider content_fetcher docs.sampletech.io

Spider実装と抽出ロジックのリファクタリング

従来の単純なテキスト抽出ではなく、アイテム型クラスと安全なセレクター関数を組合せたアーキテクチャに変更しました。これにより、DOM階層の変更や空要素の混入に対して堅牢なデータ格納が可能になります。

import scrapy
from typing import List, Optional
from scrapy.item import Item, Field

class TechArticleItem(Item):
    source_url = Field()
    headline = Field()
    publish_timestamp = Field()
    contributor = Field()
    metadata_tags = Field()
    article_body = Field()

class ContentFetcherSpider(scrapy.Spider):
    name = "content_fetcher"
    allowed_domains = ["docs.sampletech.io"]
    start_urls = ["https://docs.sampletech.io/releases/"]

    def parse(self, response):
        post_hrefs = response.css("div.release-grid article a::attr(href)").getall()
        for relative_path in post_hrefs:
            yield response.follow(relative_path, callback=self.extract_full_record)

    def extract_full_record(self, response):
        record = TechArticleItem()
        record["source_url"] = response.url
        
        # セレクターが存在しない場合のNone返却を明示的に制御
        record["headline"] = self._extract_safe(response, '//h1[@class="release-title"]/text()')
        record["publish_timestamp"] = response.xpath(
            '//meta[@property="article:published_time"]/string(@content)'
        ).get()
        record["contributor"] = self._extract_safe(response, '//span[contains(@class,"author-sig")]/text()')
        
        tag_elements = response.css("ul.tag-cloud span.tag-pill::text").getall()
        record["metadata_tags"] = [t.strip() for t in tag_elements if t.strip()]

        body_paragraphs = response.xpath(
            '//div[@id="main-body"]//p/text() | //div[@id="main-body"]//li/text()'
        ).getall()
        record["article_body"] = "\n".join(p.strip() for p in body_paragraphs if p.strip())

        yield record

    @staticmethod
    def _extract_safe(response: scrapy.http.Response, xpath_expr: str) -> Optional[str]:
        value = response.xpath(xpath_expr).get()
        return value.strip() if value else None

Easysearch連携用パイプライン設定

メモリ上で保持されたアイテムストリームを直接検索エンジンへルーティングするには、専用ミドルウェアパッケージを導入し、プロジェクト設定ファイルを書き換えます。

pip install ScrapyElasticSearch

settings.pyにおいて、バッチ単位での書き込み制御と日付ベースのインデックスローテーションルールを定義します。

ITEM_PIPELINES = {
    "scrapyelasticsearch.scrapyelasticsearch.ElasticSearchPipeline": 10,
}

ELASTICSEARCH_SERVERS = ["http://127.0.0.1:9210"]
ELASTICSEARCH_INDEX = "web_crawl_%Y-%m-%d"
ELASTICSEARCH_TYPE = "_doc"
ELASTICSEARCH_USERNAME = "elastic_operator"
ELASTICSEARCH_PASSWORD = "encrypted_access_token"
# 1バッチあたりの送信件数(メモリ使用量とスループットのバランス調整)
ELASTICSEARCH_BULK_SIZE = 75
# デッドレターキューおよび再試行設定
ELASTICSEARCH_IGNORE_EXCEPTIONS = True

クローリング実行とインデックス検証

プロジェクトルートディレクトリにて、以下のコマンドでバックグランドジョブを起動します。ログ出力先を分離することで、障害発生時のトレースが容易になります。

scrapy crawl content_fetcher --logfile=./crawl_execution_$(date +%F).log

スレッド処理が完了した後、EasysearchコンソールまたはREST API経由でターゲットインデックスのマッピングとドキュメントカウントを確認します。

{
  "took": 12,
  "timed_out": false,
  "_shards": { "total": 1, "successful": 1 },
  "hits": {
    "total": { "value": 84, "relation": "eq" },
    "max_score": null,
    "hits": []
  }
}

取得結果に対し、タイトルフィールドでの前方一致検索やタグフィールドによるAND/OR条件フィルタリングを実行すると、インデックス作成の正常性および用語分割(Analyzer)の適用状況を同時に検証できます。

タグ: Scrapy Elasticsearch Easysearch Python3 DataIngestion

7月26日 16:36 投稿