<?php
declare(strict_types=1);
final class AutocompleteService
{
public function __construct(
private readonly \Redis $redis,
) {}
/**
* Build autocomplete index using sorted sets
*/
public function index(string $phrase, float $score = 1.0): void
{
$normalized = mb_strtolower(trim($phrase));
// Add all prefixes
for ($i = 1; $i <= mb_strlen($normalized); $i++) {
$prefix = mb_substr($normalized, 0, $i);
$this->redis->zIncrBy("autocomplete:{$prefix}", $score, $normalized);
}
// Keep only top N suggestions per prefix
$this->redis->zRemRangeByRank("autocomplete:{$normalized[0]}", 0, -101);
}
/**
* Get suggestions for a prefix
*/
public function suggest(string $prefix, int $limit = 10): array
{
$normalized = mb_strtolower(trim($prefix));
// Get top scored suggestions
$results = $this->redis->zRevRange(
"autocomplete:{$normalized}",
0,
$limit - 1,
true, // with scores
);
return array_map(
fn (string $phrase, float $score) => [
'text' => $phrase,
'score' => $score,
],
array_keys($results),
array_values($results),
);
}
/**
* Trie-based approach for more efficient memory usage
*/
public function suggestWithTrie(string $prefix, int $limit = 10): array
{
$key = "trie:" . mb_strtolower(trim($prefix));
// ZRANGEBYLEX for prefix matching
$results = $this->redis->zRangeByLex(
'autocomplete:trie',
"[{$prefix}",
"[{$prefix}\xff",
0,
$limit,
);
return $results;
}
}
package search
import (
"context"
"fmt"
"strings"
"github.com/redis/go-redis/v9"
)
// Suggestion represents an autocomplete suggestion.
type Suggestion struct {
Text string `json:"text"`
Score float64 `json:"score"`
}
// AutocompleteService provides prefix-based autocomplete via Redis sorted sets.
type AutocompleteService struct {
rdb *redis.Client
}
func NewAutocompleteService(rdb *redis.Client) *AutocompleteService {
return &AutocompleteService{rdb: rdb}
}
// Index adds a phrase to the autocomplete index with all prefixes.
func (s *AutocompleteService) Index(ctx context.Context, phrase string, score float64) error {
normalized := strings.ToLower(strings.TrimSpace(phrase))
runes := []rune(normalized)
pipe := s.rdb.Pipeline()
for i := 1; i <= len(runes); i++ {
prefix := string(runes[:i])
key := fmt.Sprintf("autocomplete:%s", prefix)
pipe.ZIncrBy(ctx, key, score, normalized)
}
// Keep only top 100 suggestions per first-char prefix
firstChar := string(runes[0])
pipe.ZRemRangeByRank(ctx, fmt.Sprintf("autocomplete:%s", firstChar), 0, -101)
_, err := pipe.Exec(ctx)
return err
}
// Suggest returns top suggestions for a prefix.
func (s *AutocompleteService) Suggest(ctx context.Context, prefix string, limit int) ([]Suggestion, error) {
normalized := strings.ToLower(strings.TrimSpace(prefix))
key := fmt.Sprintf("autocomplete:%s", normalized)
results, err := s.rdb.ZRevRangeWithScores(ctx, key, 0, int64(limit-1)).Result()
if err != nil {
return nil, fmt.Errorf("autocomplete suggest: %w", err)
}
suggestions := make([]Suggestion, 0, len(results))
for _, z := range results {
suggestions = append(suggestions, Suggestion{
Text: z.Member.(string),
Score: z.Score,
})
}
return suggestions, nil
}
// SuggestWithTrie uses ZRANGEBYLEX for prefix matching (trie-based).
func (s *AutocompleteService) SuggestWithTrie(ctx context.Context, prefix string, limit int) ([]string, error) {
normalized := strings.ToLower(strings.TrimSpace(prefix))
results, err := s.rdb.ZRangeByLex(ctx, "autocomplete:trie", &redis.ZRangeBy{
Min: "[" + normalized,
Max: "[" + normalized + "\xff",
Offset: 0,
Count: int64(limit),
}).Result()
if err != nil {
return nil, fmt.Errorf("trie suggest: %w", err)
}
return results, nil
}
using System.Globalization;
using StackExchange.Redis;
namespace Search;
// Suggestion represents an autocomplete suggestion.
public sealed record Suggestion(string Text, double Score);
// AutocompleteService provides prefix-based autocomplete via Redis sorted sets.
public sealed class AutocompleteService(IDatabase redis)
{
// Index adds a phrase to the autocomplete index under all of its prefixes.
public async Task IndexAsync(string phrase, double score = 1.0)
{
var normalized = phrase.Trim().ToLowerInvariant();
if (normalized.Length == 0)
{
return;
}
// StringInfo walks grapheme clusters, so surrogate pairs stay intact
var elements = new List<string>();
var enumerator = StringInfo.GetTextElementEnumerator(normalized);
while (enumerator.MoveNext())
{
elements.Add((string)enumerator.Current);
}
var batch = redis.CreateBatch();
var tasks = new List<Task>(elements.Count + 1);
for (var i = 1; i <= elements.Count; i++)
{
var prefix = string.Concat(elements.Take(i));
tasks.Add(batch.SortedSetIncrementAsync($"autocomplete:{prefix}", normalized, score));
}
// Keep only top 100 suggestions per first-char prefix
tasks.Add(batch.SortedSetRemoveRangeByRankAsync($"autocomplete:{elements[0]}", 0, -101));
batch.Execute();
await Task.WhenAll(tasks);
}
// SuggestAsync returns the top scored suggestions for a prefix.
public async Task<IReadOnlyList<Suggestion>> SuggestAsync(string prefix, int limit = 10)
{
var normalized = prefix.Trim().ToLowerInvariant();
var entries = await redis.SortedSetRangeByRankWithScoresAsync(
$"autocomplete:{normalized}",
start: 0,
stop: limit - 1,
order: Order.Descending);
return entries
.Select(entry => new Suggestion(entry.Element.ToString(), entry.Score))
.ToList();
}
// SuggestWithTrieAsync uses ZRANGEBYLEX for prefix matching (trie-based).
public async Task<IReadOnlyList<string>> SuggestWithTrieAsync(string prefix, int limit = 10)
{
var normalized = prefix.Trim().ToLowerInvariant();
var values = await redis.SortedSetRangeByValueAsync(
"autocomplete:trie",
min: normalized,
max: normalized + "\xff",
exclude: Exclude.None,
skip: 0,
take: limit);
return values.Select(value => value.ToString()).ToList();
}
}
from __future__ import annotations
from dataclasses import dataclass
from redis.asyncio import Redis
# Suggestion represents an autocomplete suggestion.
@dataclass(frozen=True, slots=True)
class Suggestion:
text: str
score: float
class AutocompleteService:
"""Prefix-based autocomplete backed by Redis sorted sets."""
def __init__(self, redis: Redis) -> None:
self._redis = redis
async def index(self, phrase: str, score: float = 1.0) -> None:
"""Add a phrase to the autocomplete index under all of its prefixes."""
normalized = phrase.strip().lower()
if not normalized:
return
async with self._redis.pipeline(transaction=False) as pipe:
for i in range(1, len(normalized) + 1):
pipe.zincrby(f"autocomplete:{normalized[:i]}", score, normalized)
# Keep only top 100 suggestions per first-char prefix
pipe.zremrangebyrank(f"autocomplete:{normalized[0]}", 0, -101)
await pipe.execute()
async def suggest(self, prefix: str, limit: int = 10) -> list[Suggestion]:
"""Return the top scored suggestions for a prefix."""
normalized = prefix.strip().lower()
entries = await self._redis.zrevrange(
f"autocomplete:{normalized}", 0, limit - 1, withscores=True
)
return [
Suggestion(text=member.decode() if isinstance(member, bytes) else member, score=score)
for member, score in entries
]
async def suggest_with_trie(self, prefix: str, limit: int = 10) -> list[str]:
"""Use ZRANGEBYLEX for prefix matching (trie-based)."""
normalized = prefix.strip().lower()
results = await self._redis.zrangebylex(
"autocomplete:trie",
min=f"[{normalized}",
max=f"[{normalized}\xff",
start=0,
num=limit,
)
return [item.decode() if isinstance(item, bytes) else item for item in results]
## Шаг 5: Indexing Pipeline
<?php
declare(strict_types=1);
final class IndexingPipeline
{
public function __construct(
private readonly DocumentSource $source,
private readonly TextAnalyzer $analyzer,
private readonly IndexWriter $indexWriter,
private readonly AutocompleteService $autocomplete,
) {}
/**
* Process a batch of documents for indexing
*/
public function processBatch(array $documentIds): IndexingResult
{
$indexed = 0;
$errors = 0;
foreach ($documentIds as $docId) {
try {
$document = $this->source->getDocument($docId);
// 1. Analyze text
$analyzed = $this->analyzer->analyze($document->content);
// 2. Build index entry
$entry = new IndexEntry(
docId: $document->id,
tokens: $analyzed->tokens,
termFrequencies: $analyzed->termFrequencies,
metadata: [
'title' => $document->title,
'category' => $document->category,
'created_at' => $document->createdAt,
],
);
// 3. Write to index
$this->indexWriter->write($entry);
// 4. Update autocomplete
$this->autocomplete->index($document->title, 1.0);
$indexed++;
} catch (\Throwable $e) {
$errors++;
}
}
return new IndexingResult($indexed, $errors);
}
}
final class TextAnalyzer
{
public function analyze(string $text): AnalyzedText
{
// Pipeline: lowercase -> tokenize -> remove stops -> stem
$text = mb_strtolower($text);
$tokens = $this->tokenize($text);
$tokens = $this->removeStopWords($tokens);
$tokens = array_map(fn (string $t) => $this->stem($t), $tokens);
return new AnalyzedText(
tokens: $tokens,
termFrequencies: array_count_values($tokens),
);
}
private function tokenize(string $text): array
{
return preg_split('/[\s\p{P}]+/u', $text, -1, PREG_SPLIT_NO_EMPTY);
}
private function removeStopWords(array $tokens): array
{
$stops = ['the', 'a', 'is', 'are', 'and', 'or', 'but', 'in', 'on', 'at', 'to', 'for'];
return array_values(array_filter($tokens, fn (string $t) => !in_array($t, $stops, true)));
}
private function stem(string $word): string
{
// Use Snowball stemmer in production
return $word;
}
}
package search
import (
"context"
"fmt"
"regexp"
"strings"
)
// IndexingResult holds the outcome of a batch indexing operation.
type IndexingResult struct {
Indexed int
Errors int
}
// IndexEntry represents a single document to be written to the index.
type IndexEntry struct {
DocID int
Tokens []string
TermFrequencies map[string]int
Metadata map[string]string
}
// AnalyzedText holds the result of text analysis.
type AnalyzedText struct {
Tokens []string
TermFrequencies map[string]int
}
// IndexingPipeline processes documents for indexing.
type IndexingPipeline struct {
source DocumentSource
analyzer *TextAnalyzer
indexWriter IndexWriter
autocomplete *AutocompleteService
}
func NewIndexingPipeline(
source DocumentSource,
analyzer *TextAnalyzer,
writer IndexWriter,
ac *AutocompleteService,
) *IndexingPipeline {
return &IndexingPipeline{
source: source,
analyzer: analyzer,
indexWriter: writer,
autocomplete: ac,
}
}
// ProcessBatch indexes a batch of documents by their IDs.
func (p *IndexingPipeline) ProcessBatch(ctx context.Context, documentIDs []int) IndexingResult {
var result IndexingResult
for _, docID := range documentIDs {
doc, err := p.source.GetDocument(ctx, docID)
if err != nil {
result.Errors++
continue
}
// 1. Analyze text
analyzed := p.analyzer.Analyze(doc.Content)
// 2. Build index entry
entry := IndexEntry{
DocID: doc.ID,
Tokens: analyzed.Tokens,
TermFrequencies: analyzed.TermFrequencies,
Metadata: map[string]string{
"title": doc.Title,
"category": doc.Category,
"created_at": doc.CreatedAt,
},
}
// 3. Write to index
if err := p.indexWriter.Write(ctx, entry); err != nil {
result.Errors++
continue
}
// 4. Update autocomplete
if err := p.autocomplete.Index(ctx, doc.Title, 1.0); err != nil {
// Non-critical, just log
fmt.Printf("autocomplete index error for doc %d: %v\n", docID, err)
}
result.Indexed++
}
return result
}
var tokenizeRe = regexp.MustCompile(`[\s\p{P}]+`)
var analyzerStopWords = map[string]bool{
"the": true, "a": true, "is": true, "are": true,
"and": true, "or": true, "but": true, "in": true,
"on": true, "at": true, "to": true, "for": true,
}
// TextAnalyzer processes text through a pipeline: lowercase, tokenize, stop words, stem.
type TextAnalyzer struct{}
func (a *TextAnalyzer) Analyze(text string) AnalyzedText {
text = strings.ToLower(text)
parts := tokenizeRe.Split(text, -1)
tokens := make([]string, 0, len(parts))
for _, t := range parts {
if t == "" || analyzerStopWords[t] {
continue
}
// Use Snowball stemmer in production
tokens = append(tokens, t)
}
freqs := make(map[string]int, len(tokens))
for _, t := range tokens {
freqs[t]++
}
return AnalyzedText{
Tokens: tokens,
TermFrequencies: freqs,
}
}
using System.Text.RegularExpressions;
using Microsoft.Extensions.Logging;
namespace Search;
// IndexingResult holds the outcome of a batch indexing operation.
public sealed record IndexingResult(int Indexed, int Errors);
// IndexEntry represents a single document to be written to the index.
public sealed record IndexEntry(
int DocId,
IReadOnlyList<string> Tokens,
IReadOnlyDictionary<string, int> TermFrequencies,
IReadOnlyDictionary<string, string> Metadata);
// AnalyzedText holds the result of text analysis.
public sealed record AnalyzedText(
IReadOnlyList<string> Tokens,
IReadOnlyDictionary<string, int> TermFrequencies);
// IndexingPipeline processes documents for indexing.
public sealed class IndexingPipeline(
IDocumentSource source,
TextAnalyzer analyzer,
IIndexWriter indexWriter,
AutocompleteService autocomplete,
ILogger<IndexingPipeline> logger)
{
// ProcessBatchAsync indexes a batch of documents by their IDs.
public async Task<IndexingResult> ProcessBatchAsync(
IReadOnlyList<int> documentIds,
CancellationToken cancellationToken = default)
{
var indexed = 0;
var errors = 0;
foreach (var docId in documentIds)
{
try
{
var document = await source.GetDocumentAsync(docId, cancellationToken);
// 1. Analyze text
var analyzed = analyzer.Analyze(document.Content);
// 2. Build index entry
var entry = new IndexEntry(
DocId: document.Id,
Tokens: analyzed.Tokens,
TermFrequencies: analyzed.TermFrequencies,
Metadata: new Dictionary<string, string>
{
["title"] = document.Title,
["category"] = document.Category,
["created_at"] = document.CreatedAt.ToString("O"),
});
// 3. Write to index
await indexWriter.WriteAsync(entry, cancellationToken);
// 4. Update autocomplete (non-critical, failure must not drop the document)
try
{
await autocomplete.IndexAsync(document.Title);
}
catch (Exception ex)
{
logger.LogWarning(ex, "Autocomplete index failed for document {DocId}", docId);
}
indexed++;
}
catch (Exception ex)
{
logger.LogError(ex, "Indexing failed for document {DocId}", docId);
errors++;
}
}
return new IndexingResult(indexed, errors);
}
}
// TextAnalyzer processes text: lowercase, tokenize, remove stop words, stem.
public sealed partial class TextAnalyzer
{
private static readonly HashSet<string> StopWords =
[
"the", "a", "is", "are", "and", "or",
"but", "in", "on", "at", "to", "for",
];
[GeneratedRegex(@"[\s\p{P}]+")]
private static partial Regex TokenizeRegex();
public AnalyzedText Analyze(string text)
{
var tokens = TokenizeRegex()
.Split(text.ToLowerInvariant())
.Where(t => t.Length > 0 && !StopWords.Contains(t))
.Select(Stem)
.ToList();
var frequencies = tokens
.GroupBy(t => t)
.ToDictionary(g => g.Key, g => g.Count());
return new AnalyzedText(tokens, frequencies);
}
// Use a Snowball stemmer in production
private static string Stem(string word) => word;
}
from __future__ import annotations
import logging
import re
from collections import Counter
from dataclasses import dataclass
logger = logging.getLogger(__name__)
# IndexingResult holds the outcome of a batch indexing operation.
@dataclass(frozen=True, slots=True)
class IndexingResult:
indexed: int
errors: int
# IndexEntry represents a single document to be written to the index.
@dataclass(frozen=True, slots=True)
class IndexEntry:
doc_id: int
tokens: list[str]
term_frequencies: dict[str, int]
metadata: dict[str, str]
# AnalyzedText holds the result of text analysis.
@dataclass(frozen=True, slots=True)
class AnalyzedText:
tokens: list[str]
term_frequencies: dict[str, int]
TOKENIZE_RE = re.compile(r"[\s\W_]+", re.UNICODE)
ANALYZER_STOP_WORDS = frozenset(
{"the", "a", "is", "are", "and", "or", "but", "in", "on", "at", "to", "for"}
)
class TextAnalyzer:
"""Pipeline: lowercase -> tokenize -> remove stop words -> stem."""
def analyze(self, text: str) -> AnalyzedText:
tokens = [
self._stem(t)
for t in TOKENIZE_RE.split(text.lower())
if t and t not in ANALYZER_STOP_WORDS
]
return AnalyzedText(tokens=tokens, term_frequencies=dict(Counter(tokens)))
@staticmethod
def _stem(word: str) -> str:
# Use a Snowball stemmer (nltk / PyStemmer) in production
return word
class IndexingPipeline:
"""Processes documents for indexing."""
def __init__(
self,
source: DocumentSource,
analyzer: TextAnalyzer,
index_writer: IndexWriter,
autocomplete: AutocompleteService,
) -> None:
self._source = source
self._analyzer = analyzer
self._index_writer = index_writer
self._autocomplete = autocomplete
async def process_batch(self, document_ids: list[int]) -> IndexingResult:
indexed = 0
errors = 0
for doc_id in document_ids:
try:
document = await self._source.get_document(doc_id)
# 1. Analyze text
analyzed = self._analyzer.analyze(document.content)
# 2. Build index entry
entry = IndexEntry(
doc_id=document.id,
tokens=analyzed.tokens,
term_frequencies=analyzed.term_frequencies,
metadata={
"title": document.title,
"category": document.category,
"created_at": document.created_at.isoformat(),
},
)
# 3. Write to index
await self._index_writer.write(entry)
except Exception:
logger.exception("indexing failed for document %d", doc_id)
errors += 1
continue
# 4. Update autocomplete (non-critical, must not drop the document)
try:
await self._autocomplete.index(document.title, 1.0)
except Exception:
logger.warning("autocomplete index failed for document %d", doc_id, exc_info=True)
indexed += 1
return IndexingResult(indexed=indexed, errors=errors)
## Шаг 6: Шардирование индекса
Стратегия
Описание
Плюсы
Минусы
Document-based
Документы распределяются по шардам
Простая индексация
Scatter-gather для каждого запроса
Term-based
Термины распределяются по шардам
Эффективный поиск одного термина
Сложная индексация, hotspots
Рекомендуется document-based шардирование (как в Elasticsearch).
Возможные вопросы интервьюера
TF-IDF vs BM25?
BM25 лучше: saturation для TF, нормализация длины документа
BM25 -- стандарт в современных поисковых системах
Как обрабатывать fuzzy search?
Edit distance (Levenshtein)
N-gram индекс
Phonetic encoding (Soundex, Metaphone)
Как масштабировать до миллиарда документов?
Шардирование индекса (100+ шардов)
Реплики для каждого шарда
Scatter-gather pattern для поиска
Как обновлять индекс без downtime?
Near-real-time indexing (segment-based, как в Lucene)