HardКейс6 min

Распределённая файловая система

Проектирование распределённой файловой системы: GFS/HDFS архитектура, master/chunk серверы, репликация

Распределённая файловая система (DFS) позволяет хранить файлы на кластере машин с автоматической репликацией и отказоустойчивостью. Проектируем систему, вдохновлённую GFS (Google File System) и HDFS.

Шаг 1: Требования

Функциональные требования

  1. Хранение файлов от нескольких KB до нескольких TB
  2. Операции: create, read, append, delete
  3. Автоматическая репликация (3x по умолчанию)
  4. Поддержка sequential reads и appends (основной паттерн)
  5. Namespace: иерархическая структура директорий
  6. Снапшоты для backup

Нефункциональные требования

  1. Масштабирование до петабайт данных
  2. Высокая пропускная способность (throughput > latency)
  3. Отказоустойчивость: переживает потерю узлов
  4. Consistency: single-writer, multiple-readers

Шаг 2: High-Level архитектура

┌──────────┐         ┌─────────────────────────────────────┐
│  Client  │────────>│        Master Server                │
│          │         │  ┌─────────────┐ ┌───────────────┐  │
└────┬─────┘         │  │ Namespace   │ │ Chunk-to-Node │  │
     │               │  │ Manager     │ │ Mapping       │  │
     │               │  └─────────────┘ └───────────────┘  │
     │               │  ┌─────────────┐ ┌───────────────┐  │
     │               │  │ Replication │ │ Garbage       │  │
     │               │  │ Manager     │ │ Collector     │  │
     │               │  └─────────────┘ └───────────────┘  │
     │               └─────────────────────────────────────┘
     │
     │  Data flow (direct, bypasses Master)
     │
     │         ┌──────────────┐  ┌──────────────┐  ┌──────────────┐
     └────────>│ Chunk Server │  │ Chunk Server │  │ Chunk Server │
               │  ┌────────┐  │  │  ┌────────┐  │  │  ┌────────┐  │
               │  │Chunk 1 │  │  │  │Chunk 1 │  │  │  │Chunk 2 │  │
               │  │Chunk 3 │  │  │  │Chunk 2 │  │  │  │Chunk 3 │  │
               │  └────────┘  │  │  └────────┘  │  │  └────────┘  │
               └──────────────┘  └──────────────┘  └──────────────┘

Ключевой принцип: метаданные идут через Master, данные -- напрямую между клиентом и Chunk-серверами.

Шаг 3: Детальный дизайн

3.1 Master Server

<?php

declare(strict_types=1);

final class MasterServer
{
    private const CHUNK_SIZE = 64 * 1024 * 1024; // 64 MB
    private const REPLICATION_FACTOR = 3;

    public function __construct(
        private readonly NamespaceManager $namespace,
        private readonly ChunkManager $chunks,
        private readonly ChunkServerRegistry $servers,
        private readonly OperationLog $opLog,
    ) {}

    /**
     * Client requests to create a new file
     */
    public function createFile(string $path): FileHandle
    {
        // 1. Validate path
        $this->namespace->validatePath($path);

        // 2. Check if file already exists
        if ($this->namespace->exists($path)) {
            throw new FileExistsException("File already exists: {$path}");
        }

        // 3. Create namespace entry
        $fileId = $this->namespace->create($path, FileType::REGULAR);

        // 4. Log operation for recovery
        $this->opLog->append(new CreateFileOp($fileId, $path));

        return new FileHandle(
            fileId: $fileId,
            path: $path,
        );
    }

    /**
     * Client requests chunk locations for reading
     */
    public function getChunkLocations(string $path, int $offset): ChunkLocationInfo
    {
        $file = $this->namespace->getFile($path);

        // Calculate which chunk contains this offset
        $chunkIndex = intdiv($offset, self::CHUNK_SIZE);
        $chunkHandle = $this->chunks->getChunkHandle($file->id, $chunkIndex);

        if ($chunkHandle === null) {
            throw new ChunkNotFoundException("Chunk not found at index {$chunkIndex}");
        }

        // Get servers that hold replicas of this chunk
        $replicas = $this->chunks->getReplicaLocations($chunkHandle);

        // Sort by proximity to client (or load)
        $sorted = $this->sortByProximity($replicas);

        return new ChunkLocationInfo(
            chunkHandle: $chunkHandle,
            chunkServers: $sorted,
            chunkSize: self::CHUNK_SIZE,
            offsetInChunk: $offset % self::CHUNK_SIZE,
        );
    }

    /**
     * Client requests to write -- Master allocates a new chunk
     */
    public function allocateChunk(string $path): ChunkAllocation
    {
        $file = $this->namespace->getFile($path);
        $chunkIndex = $this->chunks->getNextChunkIndex($file->id);

        // Generate unique chunk handle
        $chunkHandle = bin2hex(random_bytes(8));

        // Select chunk servers (different racks for durability)
        $selectedServers = $this->selectChunkServers(self::REPLICATION_FACTOR);

        // Designate primary (first server, holds the lease)
        $primary = $selectedServers[0];
        $secondaries = array_slice($selectedServers, 1);

        // Grant lease to primary
        $leaseExpiry = time() + 60; // 60 second lease
        $this->chunks->grantLease($chunkHandle, $primary->id, $leaseExpiry);

        // Record chunk mapping
        $this->chunks->createChunk(
            $chunkHandle,
            $file->id,
            $chunkIndex,
            array_map(fn ($s) => $s->id, $selectedServers),
        );

        $this->opLog->append(new AllocateChunkOp(
            $chunkHandle, $file->id, $chunkIndex,
        ));

        return new ChunkAllocation(
            chunkHandle: $chunkHandle,
            primary: $primary,
            secondaries: $secondaries,
            leaseExpiry: $leaseExpiry,
        );
    }

    /**
     * Select chunk servers across different racks
     */
    private function selectChunkServers(int $count): array
    {
        $healthy = $this->servers->getHealthy();

        // Sort by available disk space
        usort($healthy, fn ($a, $b) => $b->availableSpace <=> $a->availableSpace);

        $selected = [];
        $usedRacks = [];

        foreach ($healthy as $server) {
            if (count($selected) >= $count) {
                break;
            }

            // Different rack preference
            if (!in_array($server->rackId, $usedRacks, true)) {
                $selected[] = $server;
                $usedRacks[] = $server->rackId;
            }
        }

        // Fill remaining if not enough racks
        if (count($selected) < $count) {
            foreach ($healthy as $server) {
                if (count($selected) >= $count) break;
                if (!in_array($server, $selected, true)) {
                    $selected[] = $server;
                }
            }
        }

        return $selected;
    }
}

3.2 Chunk Server

<?php

declare(strict_types=1);

final class ChunkServer
{
    private const HEARTBEAT_INTERVAL = 3; // seconds

    public function __construct(
        private readonly string $serverId,
        private readonly string $dataDir,
        private readonly MasterClient $master,
    ) {}

    /**
     * Write chunk data (called by client directly)
     */
    public function writeChunk(
        string $chunkHandle,
        int $offset,
        string $data,
        string $checksum,
    ): bool {
        $chunkPath = $this->getChunkPath($chunkHandle);

        // Verify checksum
        if (md5($data) !== $checksum) {
            throw new ChecksumMismatchException('Data corrupted in transit');
        }

        // Write to local disk
        $fp = fopen($chunkPath, 'cb');
        flock($fp, LOCK_EX);

        fseek($fp, $offset);
        $written = fwrite($fp, $data);

        flock($fp, LOCK_UN);
        fclose($fp);

        // Store checksum for integrity verification
        $this->storeChecksum($chunkHandle, $offset, strlen($data), $checksum);

        return $written === strlen($data);
    }

    /**
     * Read chunk data (called by client directly)
     */
    public function readChunk(
        string $chunkHandle,
        int $offset,
        int $length,
    ): string {
        $chunkPath = $this->getChunkPath($chunkHandle);

        if (!file_exists($chunkPath)) {
            throw new ChunkNotFoundException("Chunk not found: {$chunkHandle}");
        }

        $fp = fopen($chunkPath, 'rb');
        fseek($fp, $offset);
        $data = fread($fp, $length);
        fclose($fp);

        // Verify checksum
        $expectedChecksum = $this->getChecksum($chunkHandle, $offset, $length);
        if ($expectedChecksum !== null && md5($data) !== $expectedChecksum) {
            // Data corrupted, report to master
            $this->master->reportCorruption($this->serverId, $chunkHandle);
            throw new DataCorruptionException("Chunk corrupted: {$chunkHandle}");
        }

        return $data;
    }

    /**
     * Send heartbeat to Master with chunk inventory
     */
    public function sendHeartbeat(): void
    {
        $chunks = $this->listLocalChunks();

        $this->master->heartbeat(
            serverId: $this->serverId,
            chunks: $chunks,
            availableSpace: $this->getAvailableSpace(),
            totalSpace: $this->getTotalSpace(),
        );
    }

    /**
     * Replicate chunk to another server (Master-initiated)
     */
    public function replicateTo(string $chunkHandle, ChunkServerInfo $target): bool
    {
        $data = $this->readChunk($chunkHandle, 0, $this->getChunkSize($chunkHandle));

        $client = new ChunkServerClient($target);
        return $client->writeChunk($chunkHandle, 0, $data, md5($data));
    }

    private function getChunkPath(string $chunkHandle): string
    {
        // Subdirectory based on first 2 chars for filesystem performance
        $subdir = substr($chunkHandle, 0, 2);
        return "{$this->dataDir}/{$subdir}/{$chunkHandle}.chunk";
    }

    private function listLocalChunks(): array
    {
        $chunks = [];
        $dirs = glob("{$this->dataDir}/*", GLOB_ONLYDIR);

        foreach ($dirs as $dir) {
            $files = glob("{$dir}/*.chunk");
            foreach ($files as $file) {
                $handle = basename($file, '.chunk');
                $chunks[] = [
                    'handle' => $handle,
                    'size' => filesize($file),
                ];
            }
        }

        return $chunks;
    }
}

3.3 Namespace Manager

<?php

declare(strict_types=1);

final class NamespaceManager
{
    public function __construct(
        private readonly \PDO $db,
    ) {}

    public function create(string $path, FileType $type): string
    {
        $parentPath = dirname($path);
        $name = basename($path);

        // Ensure parent directory exists
        $parent = $this->getByPath($parentPath);
        if ($parent === null && $parentPath !== '/') {
            throw new DirectoryNotFoundException("Parent not found: {$parentPath}");
        }

        $parentId = $parent?->id;

        $stmt = $this->db->prepare(
            'INSERT INTO namespace (parent_id, name, type, path)
             VALUES (:parent_id, :name, :type, :path)
             RETURNING id'
        );

        $stmt->execute([
            'parent_id' => $parentId,
            'name' => $name,
            'type' => $type->value,
            'path' => $path,
        ]);

        return $stmt->fetchColumn();
    }

    public function listDirectory(string $path): array
    {
        $dir = $this->getByPath($path);

        $stmt = $this->db->prepare(
            'SELECT id, name, type, size, created_at
             FROM namespace
             WHERE parent_id = :parent_id
             ORDER BY name'
        );

        $stmt->execute(['parent_id' => $dir->id]);

        return $stmt->fetchAll(\PDO::FETCH_ASSOC);
    }

    public function delete(string $path): void
    {
        $entry = $this->getByPath($path);

        if ($entry->type === 'directory') {
            // Check if empty
            $children = $this->listDirectory($path);
            if (!empty($children)) {
                throw new DirectoryNotEmptyException("Directory not empty: {$path}");
            }
        }

        $stmt = $this->db->prepare('DELETE FROM namespace WHERE id = ?');
        $stmt->execute([$entry->id]);
    }
}

3.4 Operation Log (WAL)

<?php

declare(strict_types=1);

final class OperationLog
{
    private const CHECKPOINT_INTERVAL = 1000; // operations

    public function __construct(
        private readonly string $logPath,
        private readonly string $checkpointPath,
        private int $opsSinceCheckpoint = 0,
    ) {}

    public function append(Operation $op): void
    {
        $entry = json_encode([
            'type' => $op->getType(),
            'data' => $op->serialize(),
            'timestamp' => microtime(true),
        ]) . "\n";

        file_put_contents($this->logPath, $entry, FILE_APPEND | LOCK_EX);

        $this->opsSinceCheckpoint++;

        if ($this->opsSinceCheckpoint >= self::CHECKPOINT_INTERVAL) {
            $this->createCheckpoint();
        }
    }

    /**
     * Replay log to restore state after Master crash
     */
    public function replay(MasterServer $master): int
    {
        // 1. Load last checkpoint
        $this->loadCheckpoint($master);

        // 2. Replay operations since checkpoint
        $replayed = 0;
        $handle = fopen($this->logPath, 'r');

        while (($line = fgets($handle)) !== false) {
            $entry = json_decode(trim($line), true);
            $op = Operation::deserialize($entry['type'], $entry['data']);
            $op->apply($master);
            $replayed++;
        }

        fclose($handle);
        return $replayed;
    }

    private function createCheckpoint(): void
    {
        // Serialize entire namespace tree and chunk mappings
        // This allows faster recovery
        $this->opsSinceCheckpoint = 0;
    }
}

Шаг 4: Write Flow

1. Client -> Master: "I want to write to /data/file.txt"
2. Master -> Client: "Chunk C1, Primary: S1, Secondaries: S2, S3"
3. Client -> S1, S2, S3: Push data to all (pipelined)
4. Client -> S1 (Primary): "Commit write"
5. S1 -> S2, S3: "Apply write at offset X"
6. S2, S3 -> S1: "ACK"
7. S1 -> Client: "Write successful"

Шаг 5: Масштабирование

Компонент Стратегия
Master Shadow Master (hot standby), Operation Log
Chunk Servers Горизонтальное масштабирование (сотни серверов)
Namespace Шардирование по path prefix
Replication Rack-aware placement для durability
Recovery Automatic re-replication при потере узла

Возможные вопросы интервьюера

  1. Почему chunk size 64 MB, а не 4 KB как в обычных FS?

    • Меньше метаданных на Master
    • Меньше network round-trips
    • Оптимизация для sequential reads (batch processing)
  2. Что происходит при падении Master?

    • Shadow Master берёт на себя роль
    • Replay operation log для восстановления состояния
    • Chunk Servers продолжают обслуживать read запросы
  3. Как обнаружить повреждённые данные?

    • Checksum для каждого 64 KB блока внутри chunk
    • Периодический background scrub
    • При обнаружении -- re-replicate с здоровой копии
  4. Single Master -- bottleneck?

    • Master хранит только метаданные (в памяти)
    • Данные идут напрямую к chunk серверам
    • Для масштабирования: multiple masters с partition namespace