mybook

設計問題: 検索エンジン — 全文検索の仕組み

「検索は魔法じゃない」

「今日は検索エンジン。Googleのような全文検索システムを設計する。」

「Googleって再現できるんですか?」ソウタは目を丸くした。

レイカはコーヒーを一口飲んでから続けた。「もちろん規模は違う。でも仕組みの本質は理解できる。転置インデックス、クロール、ランキング——この3つを設計できれば面接では十分だ。」

「なぜ『転置』インデックスと言うんですか?」

「普通のインデックスは『ドキュメントID → 単語一覧』の対応表。転置インデックスはその逆、『単語 → ドキュメントID一覧』。ユーザーが『Rails』と検索したとき、『Railsという単語が含まれるドキュメントのリスト』を一瞬で返せる。」

「なるほど。図書館で言えば……目次じゃなくて索引の方ですね。」

「正確な理解だ。」レイカが微笑んだ。「では、面接官の視点で一から設計してみよう。」


Step 1: 要件の確認

「まず要件を明確にする。面接ではここを丁寧にやると印象が良い。」

機能要件:
- Webページのクロール(インターネット上のHTMLを収集)
- 全文検索(キーワードに関連するページを返す)
- 検索結果のランキング
- 日本語・英語への対応
- スペルミス訂正とクエリサジェスト
- パーソナライズされた検索結果

非機能要件:
- インデックス対象: 50億ページ(5B)
- QPS(検索): 10,000/sec
- 検索レイテンシ: 500ms以内
- クロール頻度: 新規ページは数日以内、更新ページは定期的
- インデックス更新: 準リアルタイム(数時間以内)
- 可用性: 99.99%

「面接官から『スペルミス訂正は必要か』と問われたとき、要件定義で確認済みなら迷わず答えられる。」レイカが付け加えた。

INFO

要件確認フェーズで「何をスコープに入れるか」を面接官と合意しておくことが重要。途中で範囲が広がると設計が破綻する。「スペルミス訂正やパーソナライゼーションはスコープ外としてよいですか?」と聞けると評価が上がる。


Step 2: 規模の概算

インデックスのストレージ:
  対象ページ数: 5B(50億)
  1ページの平均サイズ(テキスト抽出後): 100KB
  全コンテンツ: 5B × 100KB = 500TB

転置インデックスのサイズ:
  ユニークな単語数(英語): 約100万語
  1単語あたりの平均ドキュメント数: 1,000件
  ポスティングリスト: 100万 × 1,000 × 12B(doc_id 8B + TF 4B)≒ 12TB

クロール帯域:
  1日にクロールするページ: 1億(100M)
  1ページ: 平均500KB
  必要な帯域: 100M × 500KB / 86,400 ≒ 580MB/sec

検索QPS:
  ピーク時: 10,000 QPS
  キャッシュヒット率 90%として: 実際にOpenSearchを叩くのは 1,000 QPS
  OpenSearchノード1台あたりの処理能力: 500 QPS
  必要ノード数: 1,000 / 500 = 2台(レプリカ含めて最低4台)

Step 3: 高レベル設計

Loading diagram...

「コンポーネントの役割を一言で言えるかい?」レイカが問う。

「URL Frontierはクロール対象のURLキュー、Crawlerは実際にHTMLを取得、ParserはHTMLからテキストを抽出、IndexerはOpenSearchへのインデックス構築……」

「正確だ。面接では図を描きながら口頭で説明できると完璧だ。」


Step 4: 詳細設計

クローラーの分散設計

「クロールは単体のプロセスでは絶対に無理だ。1日1億ページ処理するには、クローラーを何台並列で動かす必要があるか計算してみて。」

1台のクローラーの処理能力:
  1リクエストのレイテンシ: 平均500ms(ネットワーク + HTML取得)
  同時並列数: 100スレッド
  1台あたりの処理能力: 100 / 0.5 = 200 RPS = 1,728万 req/day

必要なクローラー台数:
  100M / 1,728万 ≒ 6台(余裕を持って10台)

「Politeness(礼儀正しさ)の問題もある。同じドメインへのリクエストは1秒以上間隔を空けないといけない。」

Loading diagram...

WARNING

クローラーが礼儀のない実装をすると、対象サーバーに過大な負荷をかける。robots.txt の遵守、同じドメインへのアクセス間隔の確保(Crawl-delay)、User-Agent の明示は必須。面接でこれを言えると「実践的な考え方ができる」と評価される。

Golangでのクローラー実装(並列処理 + レートリミット)

// crawler/worker.go
package crawler
 
import (
	"context"
	"log"
	"net/http"
	"net/url"
	"sync"
	"time"
 
	"golang.org/x/time/rate"
)
 
// ドメインごとのレートリミッター管理
type DomainLimiter struct {
	mu       sync.Mutex
	limiters map[string]*rate.Limiter
}
 
func NewDomainLimiter() *DomainLimiter {
	return &DomainLimiter{
		limiters: make(map[string]*rate.Limiter),
	}
}
 
func (dl *DomainLimiter) GetLimiter(domain string) *rate.Limiter {
	dl.mu.Lock()
	defer dl.mu.Unlock()
	if limiter, ok := dl.limiters[domain]; ok {
		return limiter
	}
	// 1ドメインあたり1req/secのレート制限
	limiter := rate.NewLimiter(rate.Every(time.Second), 1)
	dl.limiters[domain] = limiter
	return limiter
}
 
type CrawlerWorker struct {
	client      *http.Client
	domainLimit *DomainLimiter
	robotsCache *RobotsCache
	s3Client    S3Client
	sqsClient   SQSClient
}
 
func NewCrawlerWorker(s3 S3Client, sqs SQSClient) *CrawlerWorker {
	return &CrawlerWorker{
		client:      &http.Client{Timeout: 30 * time.Second},
		domainLimit: NewDomainLimiter(),
		robotsCache: NewRobotsCache(),
		s3Client:    s3,
		sqsClient:   sqs,
	}
}
 
func (w *CrawlerWorker) ProcessURL(ctx context.Context, rawURL string) error {
	parsed, err := url.Parse(rawURL)
	if err != nil {
		return err
	}
	domain := parsed.Host
 
	// robots.txtチェック
	if w.robotsCache.IsDisallowed(domain, rawURL) {
		log.Printf("robots.txt disallows: %s", rawURL)
		return nil
	}
 
	// ドメインごとのレートリミット適用
	limiter := w.domainLimit.GetLimiter(domain)
	if err := limiter.Wait(ctx); err != nil {
		return err
	}
 
	// HTML取得
	resp, err := w.client.Get(rawURL)
	if err != nil {
		return err
	}
	defer resp.Body.Close()
 
	if resp.StatusCode != http.StatusOK {
		return nil
	}
 
	// S3に保存 + 解析ジョブをSQSへ
	key := urlToS3Key(rawURL)
	if err := w.s3Client.PutObject(ctx, key, resp.Body); err != nil {
		return err
	}
	return w.sqsClient.SendMessage(ctx, ParseJob{URL: rawURL, S3Key: key})
}
 
// 並列ワーカープール
func RunWorkerPool(ctx context.Context, concurrency int, urls <-chan string) {
	worker := NewCrawlerWorker(newS3Client(), newSQSClient())
	var wg sync.WaitGroup
 
	for i := 0; i < concurrency; i++ {
		wg.Add(1)
		go func() {
			defer wg.Done()
			for rawURL := range urls {
				if err := worker.ProcessURL(ctx, rawURL); err != nil {
					log.Printf("crawl error %s: %v", rawURL, err)
				}
			}
		}()
	}
	wg.Wait()
}

URL Frontierの優先度付きキュー

// frontier/priority_queue.go
package frontier
 
import (
	"context"
	"encoding/json"
 
	"github.com/aws/aws-sdk-go-v2/service/sqs"
)
 
type URLJob struct {
	URL      string  `json:"url"`
	Priority float64 `json:"priority"` // 0.0〜1.0(PageRankスコア等)
	Depth    int     `json:"depth"`
}
 
type PriorityFrontier struct {
	sqsClient      *sqs.Client
	highQueueURL   string // Priority 0.7以上
	mediumQueueURL string // Priority 0.3〜0.7
	lowQueueURL    string // Priority 0.3未満
}
 
func (f *PriorityFrontier) Enqueue(ctx context.Context, job URLJob) error {
	queueURL := f.selectQueue(job.Priority)
	body, _ := json.Marshal(job)
	_, err := f.sqsClient.SendMessage(ctx, &sqs.SendMessageInput{
		QueueUrl:    &queueURL,
		MessageBody: strPtr(string(body)),
	})
	return err
}
 
func (f *PriorityFrontier) selectQueue(priority float64) string {
	switch {
	case priority >= 0.7:
		return f.highQueueURL
	case priority >= 0.3:
		return f.mediumQueueURL
	default:
		return f.lowQueueURL
	}
}
 
func strPtr(s string) *string { return &s }

Railsでのrobots.txt解析

# app/services/robots_parser.rb
class RobotsParser
  def initialize(domain)
    @rules = parse_robots_txt(domain)
  end
 
  def disallows?(url)
    @rules.any? do |rule|
      rule[:agent] == '*' && url.match?(rule[:pattern])
    end
  end
 
  def crawl_delay
    @rules.find { |r| r[:type] == :crawl_delay }&.dig(:value) || 1
  end
 
  private
 
  def parse_robots_txt(domain)
    response = HTTParty.get("https://#{domain}/robots.txt", timeout: 10)
    return [] unless response.success?
 
    rules = []
    current_agent = nil
    response.body.each_line do |line|
      line = line.strip.downcase
      if line.start_with?('user-agent:')
        current_agent = line.split(':', 2).last.strip
      elsif line.start_with?('disallow:') && current_agent == '*'
        path = line.split(':', 2).last.strip
        rules << { agent: '*', pattern: Regexp.new("^#{Regexp.escape(path)}"), type: :disallow }
      elsif line.start_with?('crawl-delay:')
        delay = line.split(':', 2).last.strip.to_i
        rules << { type: :crawl_delay, value: delay }
      end
    end
    rules
  end
end

転置インデックスの構造と圧縮

「転置インデックスの内部構造を詳しく説明できると差がつく」とレイカが言った。

転置インデックスの構造:

単語         ポスティングリスト
────────────────────────────────────────────
"rails"   → [(doc_id:101, tf:0.05, pos:[3,15,42]),
              (doc_id:307, tf:0.12, pos:[1,8]),
              (doc_id:892, tf:0.03, pos:[20])]

"deploy"  → [(doc_id:101, tf:0.02, pos:[45]),
              (doc_id:512, tf:0.08, pos:[2,19,33])]

ポスティングリストの1エントリ:
  doc_id   : 8バイト(ドキュメントID)
  tf       : 4バイト(TF値 float32)
  pos_count: 2バイト(位置情報の数)
  pos[]    : 4バイト × N(単語の出現位置)
Loading diagram...

ポスティングリストの圧縮(Delta Encoding + VarInt)

// indexer/compression.go
package indexer
 
import "encoding/binary"
 
// Delta Encoding: doc_idの差分を格納してサイズを削減
// 例: [101, 307, 892] → [101, 206, 585](差分)
// さらにVarInt(可変長整数)でエンコードすると70%以上圧縮できる
 
func EncodePostingList(docIDs []uint64) []byte {
	buf := make([]byte, 0, len(docIDs)*4)
	var prev uint64
	for _, id := range docIDs {
		delta := id - prev
		buf = appendVarInt(buf, delta)
		prev = id
	}
	return buf
}
 
func appendVarInt(buf []byte, v uint64) []byte {
	var tmp [binary.MaxVarintLen64]byte
	n := binary.PutUvarint(tmp[:], v)
	return append(buf, tmp[:n]...)
}
 
func DecodePostingList(data []byte) []uint64 {
	var ids []uint64
	var prev uint64
	for len(data) > 0 {
		delta, n := binary.Uvarint(data)
		data = data[n:]
		id := prev + delta
		ids = append(ids, id)
		prev = id
	}
	return ids
}

形態素解析(日本語対応)

「日本語は英語と違って単語の区切りがないから、特別な処理が必要だ」とレイカが続けた。

# app/services/japanese_tokenizer.rb
# Kuromojiを使った日本語形態素解析(natto gem経由でMeCabを利用)
require 'natto'
 
class JapaneseTokenizer
  STOP_WORDS = %w[の は が を に で も と や].freeze
 
  def self.tokenize(text)
    nm = Natto::MeCab.new
    tokens = []
 
    nm.parse(text) do |node|
      next if node.is_eos? || node.is_bos?
 
      feature = node.feature.split(',')
      part_of_speech = feature[0]  # 品詞(名詞、動詞など)
      base_form      = feature[6]  # 基本形(原形)
 
      # 名詞・動詞・形容詞のみ抽出(ストップワード除外)
      next unless %w[名詞 動詞 形容詞].include?(part_of_speech)
      next if STOP_WORDS.include?(node.surface)
      next if base_form == '*'
 
      tokens << base_form.downcase
    end
 
    tokens
  end
 
  # 英語はステミング(語幹化)で対応
  def self.stem_english(word)
    word.gsub(/ing$/, '').gsub(/ed$/, '').gsub(/s$/, '')
  end
 
  # 日本語と英語が混在するテキストを統合処理
  def self.tokenize_mixed(text)
    japanese_parts = text.scan(/[\p{Hiragana}\p{Katakana}\p{Han}]+/)
    english_parts  = text.scan(/[a-zA-Z]+/)
 
    japanese_tokens = japanese_parts.flat_map { |p| tokenize(p) }
    english_tokens  = english_parts.map { |w| stem_english(w.downcase) }
 
    (japanese_tokens + english_tokens).uniq.compact
  end
end

INFO

MeCab / Kuromoji のような形態素解析エンジンは「辞書」を持っており、専門用語(例: Rails、Kubernetes)はユーザー辞書に追加することで精度が上がる。技術系検索エンジンでは専門用語辞書の整備が検索品質に直結する。面接でこの観点を挙げると「ドメイン特化の知識がある」と評価される。


Railsでの全文検索実装(pg_search gem)

# Gemfile
gem 'pg_search'
 
# app/models/document.rb
class Document < ApplicationRecord
  include PgSearch::Model
 
  # PostgreSQLのtsvector型を使った全文検索
  pg_search_scope :search_fulltext,
    against: {
      title:   'A',  # タイトルは重み高
      body:    'B',  # 本文は重み中
      summary: 'C'   # サマリーは重み低
    },
    using: {
      tsearch: {
        dictionary: 'english',
        tsvector_column: 'searchable_content',
        highlight: {
          StartSel: '<mark>',
          StopSel:  '</mark>',
          MaxWords: 50,
          MinWords: 15,
          ShortWord: 3
        }
      }
    }
end
 
# db/migrate/add_search_index_to_documents.rb
class AddSearchIndexToDocuments < ActiveRecord::Migration[7.1]
  def change
    add_column :documents, :searchable_content, :tsvector
    add_index  :documents, :searchable_content, using: :gin
 
    execute <<~SQL
      CREATE OR REPLACE FUNCTION update_searchable_content()
      RETURNS trigger AS $$
      BEGIN
        NEW.searchable_content :=
          setweight(to_tsvector('english', coalesce(NEW.title, '')), 'A') ||
          setweight(to_tsvector('english', coalesce(NEW.body, '')), 'B');
        RETURN NEW;
      END;
      $$ LANGUAGE plpgsql;
 
      CREATE TRIGGER documents_searchable_update
      BEFORE INSERT OR UPDATE ON documents
      FOR EACH ROW EXECUTE FUNCTION update_searchable_content();
    SQL
  end
end
 
# app/controllers/search_controller.rb
class SearchController < ApplicationController
  def index
    query = params[:q].to_s.strip
    return render json: { results: [] } if query.blank?
 
    cache_key = "search:#{Digest::SHA256.hexdigest(query)}:#{params[:page]}"
    results = Rails.cache.fetch(cache_key, expires_in: 5.minutes) do
      Document.search_fulltext(query)
              .with_pg_search_highlight
              .order(pg_search_rank: :desc)
              .page(params[:page]).per(10)
              .map { |d| serialize_result(d) }
    end
 
    render json: { results: results, query: query }
  end
 
  private
 
  def serialize_result(doc)
    {
      id:      doc.id,
      url:     doc.url,
      title:   doc.title,
      snippet: doc.pg_search_highlight,
      score:   doc.pg_search_rank
    }
  end
end

ランキングアルゴリズム

「検索結果の順位付けが最も奥深い部分だ」とレイカが言った。「TF-IDFとBM25とPageRankの違いを説明できるか?」

TF-IDF vs BM25

TF-IDF(古典的手法):
  TF(Term Frequency): 文書内での単語の出現頻度
    TF = 単語の出現回数 / 文書の総単語数

  IDF(Inverse Document Frequency): 単語の希少性
    IDF = log(総文書数 / 単語が出現する文書数)

  TF-IDF = TF × IDF

BM25(Best Match 25 — 現代の標準):
  TF-IDFの弱点を補完:
  ① TFの上限(saturation): 同じ単語が100回出ても差が縮まる
  ② 文書長の正規化: 長い文書ほどTFが低く出る問題を補正

  BM25 = Σ IDF(qi) × [TF × (k1+1)] / [TF + k1×(1 - b + b×|d|/avgdl)]

  k1: TFの飽和パラメータ(通常1.2〜2.0)
  b:  文書長正規化(通常0.75)
  |d|: 文書の長さ
  avgdl: 全文書の平均長

計算例(「Rails」を含む文書Aをスコアリング):
  出現回数: 5回、文書長: 100単語、avgdl: 150単語
  k1=1.5、b=0.75、全文書中の出現: 100万件/50億件

  IDF = log((5B - 1M + 0.5) / (1M + 0.5) + 1) ≒ 8.1
  BM25 = 8.1 × [5 × 2.5] / [5 + 1.5 × (0.25 + 0.75 × 100/150)]
       ≒ 14.7

PageRank(リンクの権威性)

基本的な考え方:
  多くの優良ページからリンクされているページは信頼性が高い

PR(A) = (1-d) + d × Σ(PR(T) / C(T))

  d: ダンピングファクター(通常0.85)
     ランダムサーファーが別のページに飛ぶ確率 = 1-d = 0.15
  T: Aにリンクしているページ
  C(T): Tの外部リンク数

収束するまで繰り返し計算(通常50〜100イテレーション)
Google規模: Apache Spark GraphXで数日かけて分散計算
# app/services/ranking_service.rb
class RankingService
  # 検索結果のスコア計算(複数シグナルの統合)
  def self.calculate_score(doc_id, query_terms, user_context)
    bm25_score      = calculate_bm25(doc_id, query_terms)
    page_rank       = PageRankCache.get(doc_id) || 0.0
    freshness_score = calculate_freshness(doc_id)
    personal_score  = calculate_personalization(doc_id, user_context)
    quality_score   = calculate_quality(doc_id)  # スパム判定、広告比率
 
    bm25_score      * 0.35 +
      page_rank     * 0.30 +
      freshness_score * 0.15 +
      personal_score  * 0.10 +
      quality_score   * 0.10
  end
 
  def self.calculate_bm25(doc_id, query_terms, k1: 1.5, b: 0.75)
    doc    = DocumentMetadata.find(doc_id)
    avg_dl = DocumentMetadata.average(:token_count).to_f
 
    query_terms.sum do |term|
      tf  = doc.term_frequency(term).to_f
      idf = Math.log((TOTAL_DOCS - doc_frequency(term) + 0.5) /
                     (doc_frequency(term) + 0.5) + 1)
      idf * (tf * (k1 + 1)) /
        (tf + k1 * (1 - b + b * doc.token_count / avg_dl))
    end
  end
 
  def self.calculate_freshness(doc_id)
    last_crawled = CrawlMetadata.last_crawled_at(doc_id)
    days_old = (Time.current - last_crawled) / 1.day
    Math.exp(-0.01 * days_old)  # 指数減衰:古いほど低スコア
  end
 
  def self.calculate_personalization(doc_id, user_context)
    return 0.0 unless user_context[:user_id]
    domain = DocumentMetadata.domain(doc_id)
    UserDomainAffinity.score(user_context[:user_id], domain)
  end
end

クエリ処理の詳細

# app/services/query_processor.rb
class QueryProcessor
  def self.process(raw_query, user_id: nil)
    # 1. クエリの正規化
    normalized = normalize(raw_query)
 
    # 2. スペルミス訂正
    corrected, was_corrected = correct_spelling(normalized)
 
    # 3. クエリの解析(AND/OR/NOT検索)
    parsed = parse_query(corrected)
 
    # 4. 検索実行
    results = execute_search(parsed, user_id: user_id)
 
    # 5. スニペット生成
    results_with_snippets = attach_snippets(results, corrected)
 
    {
      results:         results_with_snippets,
      corrected_query: was_corrected ? corrected : nil,
      total:           results.total_count
    }
  end
 
  # 編集距離(Levenshtein)を使ったスペルミス訂正
  def self.correct_spelling(query)
    tokens    = query.split
    corrected = tokens.map do |token|
      candidates = SpellDictionary.nearest_words(token, max_distance: 2)
      candidates.max_by { |w| DictionaryFrequency.get(w) } || token
    end
    was_corrected = corrected != tokens
    [corrected.join(' '), was_corrected]
  end
 
  # スニペット生成(検索語周辺の文脈を抽出してハイライト)
  def self.generate_snippet(doc_content, query_terms, window: 80)
    sentences = doc_content.split(/[。.!?]/)
    best = sentences.max_by { |s| query_terms.count { |t| s.include?(t) } }
    return '' unless best
 
    pos   = query_terms.filter_map { |t| best.index(t) }.min || 0
    start = [pos - window / 2, 0].max
    snippet = best[start, window]
 
    query_terms.reduce(snippet) do |s, term|
      s.gsub(term, "<mark>#{term}</mark>")
    end
  end
end

リアルタイムインデックス更新(Kafkaパイプライン)

Loading diagram...
// indexer/kafka_consumer.go
package indexer
 
import (
	"context"
	"encoding/json"
	"log"
 
	"github.com/segmentio/kafka-go"
)
 
type ParsedDocument struct {
	URL       string   `json:"url"`
	Title     string   `json:"title"`
	Body      string   `json:"body"`
	Tokens    []string `json:"tokens"`
	CrawledAt string   `json:"crawled_at"`
}
 
type IndexerConsumer struct {
	reader    *kafka.Reader
	osClient  OpenSearchClient
	batchSize int
}
 
func NewIndexerConsumer(brokers []string, os OpenSearchClient) *IndexerConsumer {
	return &IndexerConsumer{
		reader: kafka.NewReader(kafka.ReaderConfig{
			Brokers:  brokers,
			Topic:    "parsed-doc",
			GroupID:  "indexer-group",
			MinBytes: 10e3,
			MaxBytes: 10e6,
		}),
		osClient:  os,
		batchSize: 500,
	}
}
 
func (c *IndexerConsumer) Run(ctx context.Context) error {
	batch := make([]ParsedDocument, 0, c.batchSize)
 
	for {
		msg, err := c.reader.ReadMessage(ctx)
		if err != nil {
			return err
		}
 
		var doc ParsedDocument
		if err := json.Unmarshal(msg.Value, &doc); err != nil {
			log.Printf("parse error: %v", err)
			continue
		}
 
		batch = append(batch, doc)
 
		// バッチが溜まったらBulk APIでまとめてインデックス
		if len(batch) >= c.batchSize {
			if err := c.osClient.BulkIndex(ctx, batch); err != nil {
				log.Printf("bulk index error: %v", err)
			} else {
				log.Printf("indexed %d documents", len(batch))
			}
			batch = batch[:0]
		}
	}
}

AWSアーキテクチャ

Loading diagram...
クロールシステム:
  EC2クローラー群:
    インスタンスタイプ: c6i.2xlarge(スポット)
    台数: 10台(オートスケーリング、最大30台)
    コスト削減: スポット価格でオンデマンドの70%削減
  SQS URL Frontier:
    優先度別に3キュー(High / Medium / Low Priority)
    Visibility Timeout: 60秒
    Dead Letter Queue: 5回失敗したURLを隔離
  S3生HTML保存:
    ライフサイクル: 90日後にGlacierへ移行
    ストレージクラス: Intelligent-Tiering
 
解析・インデックスパイプライン:
  Amazon MSK(Kafka):
    トピック: raw-html(パーティション: 24)、parsed-doc(パーティション: 24)
    レプリカ: 3
    保持期間: 7日間
  ECS Parser(Fargate):
    タスク数: 20(QPSに応じてオートスケール)
    CPU: 2vCPU、メモリ: 4GB
  ECS Indexer(Fargate):
    タスク数: 10
    OpenSearch Bulk APIで500件ずつバッチ処理
 
検索インフラ:
  Amazon OpenSearch Service:
    バージョン: OpenSearch 2.x
    マスターノード: 3台(m6g.large.search)
    データノード: 20台(r6g.2xlarge.search、192GB SSD)
    シャード設計: 1,000シャード(5B docs / 5M per shard)
    レプリカ: 1(可用性確保)
  ElastiCache for Redis:
    クラスターモード: 3シャード × 2レプリカ
    用途: 検索結果キャッシュ(Top-10,000クエリ)
    TTL: 5分
  API Gateway + Lambda:
    検索APIのスケールアウト対応
    Provisioned Concurrency: 100(コールドスタート対策)
 
ストレージ:
  DynamoDB: URLメタデータ(クロール日時、URLハッシュ、HTTP Status)
  Aurora PostgreSQL: PageRankスコア、インデックスメタデータ

検索のパーソナライゼーション

# app/services/personalization_service.rb
class PersonalizationService
  HISTORY_WINDOW = 30.days
 
  # ユーザーのクリック履歴からドメイン親和性を更新
  def self.record_click(user_id:, url:, query:, rank:)
    domain = URI.parse(url).host
    # 上位に表示されてもクリックされた場合に重みを増やす
    click_weight = rank <= 3 ? 1.5 : 1.0
 
    Rails.cache.write(
      "affinity:#{user_id}:#{domain}",
      current_affinity(user_id, domain) + click_weight,
      expires_in: HISTORY_WINDOW
    )
 
    SearchLog.create!(
      user_id:     user_id,
      query:       query,
      clicked_url: url,
      rank:        rank,
      searched_at: Time.current
    )
  end
 
  def self.current_affinity(user_id, domain)
    Rails.cache.read("affinity:#{user_id}:#{domain}") || 0.0
  end
 
  # ユーザーの検索履歴からクエリサジェストを生成
  def self.query_suggestions(user_id, prefix)
    recent_queries = SearchLog
      .where(user_id: user_id)
      .where(searched_at: HISTORY_WINDOW.ago..)
      .where("query LIKE ?", "#{prefix}%")
      .group(:query)
      .order(Arel.sql('COUNT(*) DESC'))
      .limit(5)
      .pluck(:query)
 
    # 個人履歴 + 全体の人気クエリをマージ
    popular = PopularQuery.suggestions_for(prefix, limit: 5)
    (recent_queries + popular).uniq.first(8)
  end
end

スペルミス訂正(編集距離 + BKツリー)

// spellcheck/corrector.go
package spellcheck
 
// Levenshtein距離(編集距離)の計算
// 2つの文字列の最小編集操作数(挿入・削除・置換)
func LevenshteinDistance(s, t string) int {
	sr := []rune(s)
	tr := []rune(t)
	m, n := len(sr), len(tr)
 
	dp := make([][]int, m+1)
	for i := range dp {
		dp[i] = make([]int, n+1)
	}
	for i := 0; i <= m; i++ { dp[i][0] = i }
	for j := 0; j <= n; j++ { dp[0][j] = j }
 
	for i := 1; i <= m; i++ {
		for j := 1; j <= n; j++ {
			cost := 1
			if sr[i-1] == tr[j-1] {
				cost = 0
			}
			dp[i][j] = min3(
				dp[i-1][j]+1,
				dp[i][j-1]+1,
				dp[i-1][j-1]+cost,
			)
		}
	}
	return dp[m][n]
}
 
func min3(a, b, c int) int {
	if b < a { a = b }
	if c < a { a = c }
	return a
}
 
// BKツリーで編集距離2以内の候補を高速検索
type SpellCorrector struct {
	bkTree *BKTree
}
 
func (sc *SpellCorrector) Correct(word string) string {
	candidates := sc.bkTree.Search(word, 2)
	if len(candidates) == 0 {
		return word
	}
	// 出現頻度が最も高い候補を選択
	best := candidates[0]
	bestFreq := WordFrequency(best)
	for _, c := range candidates[1:] {
		if f := WordFrequency(c); f > bestFreq {
			best, bestFreq = c, f
		}
	}
	return best
}

監視とSLO

Loading diagram...
# app/services/search_monitor.rb
class SearchMonitor
  NAMESPACE = 'SearchEngine/Metrics'
 
  def self.record_search(query:, latency_ms:, result_count:, cache_hit:)
    cw = Aws::CloudWatch::Client.new
    cw.put_metric_data(
      namespace: NAMESPACE,
      metric_data: [
        { metric_name: 'SearchLatency',  value: latency_ms,              unit: 'Milliseconds' },
        { metric_name: 'ResultCount',    value: result_count,            unit: 'Count' },
        { metric_name: 'CacheHitRate',   value: cache_hit ? 1 : 0,       unit: 'Count' },
        { metric_name: 'ZeroResultRate', value: result_count.zero? ? 1 : 0, unit: 'Count' }
      ]
    )
  end
 
  def self.record_crawl(url:, success:, status_code:, latency_ms:)
    cw = Aws::CloudWatch::Client.new
    cw.put_metric_data(
      namespace: NAMESPACE,
      metric_data: [
        { metric_name: 'CrawlSuccess',  value: success ? 1 : 0,                     unit: 'Count' },
        { metric_name: 'CrawlLatency',  value: latency_ms,                          unit: 'Milliseconds' },
        { metric_name: 'HTTP4xxErrors', value: (400..499).cover?(status_code) ? 1 : 0, unit: 'Count' }
      ]
    )
  end
end
# CloudWatch Alarms(主要なSLO監視)
 
SearchLatencyP99:
  メトリクス: SearchLatency(p99)
  閾値: 500ms超過
  評価期間: 5分
  アクション: SNS → PagerDuty
 
IndexingDelay:
  メトリクス: クロール完了からOpenSearch反映までの遅延
  閾値: 3時間超過
  評価期間: 30分
 
CrawlSuccessRate:
  メトリクス: CrawlSuccess(成功率)
  閾値: 95%未満
  評価期間: 15分
 
ZeroResultRate:
  メトリクス: ZeroResultRate(ゼロ件率)
  閾値: 10%超過
  評価期間: 10分
  用途: インデックス消失・シャード障害の早期検知
 
SLO目標:
  検索レイテンシ p99: 500ms以内
  クロール成功率: 95%以上
  インデックス反映遅延: 3時間以内
  可用性: 99.99%(月間ダウンタイム52分以内)

INFO

ZeroResultRate(検索結果ゼロ件の割合)の監視は見落とされがちだが重要。この数値が急上昇した場合、OpenSearch のインデックスが消えた、シャードが壊れたなどの重大障害を早期検知できる。面接でこの観点を挙げると「運用まで考えている」と評価される。


Step 5: ボトルネック対策

「最後に、ボトルネックと対策を整理しよう」とレイカがホワイトボードに書き始めた。

ボトルネック1: OpenSearchのクエリ遅延
  原因: シャードが多すぎる / ノードが足りない
  対策:
    - ElastiCacheでTop-10,000クエリをキャッシュ(TTL: 5分)
    - ホットシャードの特定と再分散(_cat/shards API)
    - クエリタイムアウト設定(timeout: "200ms")

ボトルネック2: クロール速度の低下
  原因: 特定ドメインでのレートリミット
  対策:
    - DNS解決のキャッシュ(TTL: 1時間)
    - HTTPコネクションプールの再利用(keep-alive)
    - 動的なCrawl-delayの調整

ボトルネック3: インデックス更新の遅延
  原因: Kafkaのラグが増大している
  対策:
    - Indexerのコンシューマー数を増やす(ECSタスク追加)
    - Bulk APIのバッチサイズ調整(500 → 1,000)
    - OpenSearchのリフレッシュ間隔を延ばす(1s → 30s)

ボトルネック4: 検索品質の低下(スパム流入)
  原因: 低品質ページが大量にインデックスされている
  対策:
    - PageRankの低いドメインのクロール優先度を下げる
    - コンテンツ品質スコア(広告比率、テキスト/HTML比)でフィルタ
    - 重複コンテンツの検出(SimHash、Jaccard類似度)

面接官との深掘り会話

「では面接官視点で深掘りするよ」とレイカが眼鏡を外して言った。

面接官「なぜPageRankは計算に時間がかかるのか?」

「イテレーション計算です。全50億ページのPageRankを収束させるには50〜100回の反復が必要で、毎回グラフ全体を走査します。Google規模では Apache Spark GraphX などの分散グラフ処理フレームワークを使って数日かけて計算します。」

面接官「転置インデックスを更新する際の課題は?」

「インデックスの更新中も検索リクエストは来ます。OpenSearch(Elasticsearch)は『セグメント』という不変ファイルにインデックスを分割して管理し、新規ドキュメントは新セグメントに書き込み、バックグラウンドでセグメントをマージします。これにより読み書きの競合を避けています。」

面接官「検索をどのようにスケールアウトするか?」

「読み取り系はOpenSearchのレプリカを増やしてQPSを分散します。書き込み系はKafkaのパーティションとIndexerのコンシューマー数を増やします。クローラーはECSオートスケーリングでスポットインスタンスを追加します。キャッシュ層のElastiCacheも水平スケールが可能です。」

面接官「日本語検索で特有の難しさは?」

「英語と違って単語の区切りがないため、形態素解析が必須です。『東京都知事』は『東京/都/知事』か『東/京都/知事』かで意味が変わります。また、同義語(Synonyms)の辞書整備も重要で、例えば『スマホ』と『スマートフォン』を同一語として扱う設定が必要です。」

ソウタは長い一日の終わりにノートを閉じながら思った。「検索エンジンって、クロール、インデックス、ランキングの3層が全部違う技術で支えられているんですね。」

「そう。そしてそれぞれが独立してスケールできる構造になっているのが設計の核心だ。面接でこの分離を意識した設計を話せると、上位エンジニア向けポジションでも戦える。」

「来週の面接、受けてきます。」

「行ってこい。転置インデックスの話が出たら、BM25とPageRankの違いまで自信を持って話せるはずだ。」