設計問題: 検索エンジン — 全文検索の仕組み
「検索は魔法じゃない」
「今日は検索エンジン。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: 高レベル設計
「コンポーネントの役割を一言で言えるかい?」レイカが問う。
「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秒以上間隔を空けないといけない。」
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(単語の出現位置)
ポスティングリストの圧縮(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
endINFO
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パイプライン)
// 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アーキテクチャ
クロールシステム:
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
# 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の違いまで自信を持って話せるはずだ。」