ddia/ch11-batch.md - DDIA

Sources:

  • 3FS
    • https://maknee.github.io/blog/2025/3FS-Performance-Journal-1/
    • Notes: https://github.com/deepseek-ai/3FS/blob/ee9a5cee0a85c64f4797bf380257350ca1becd36/docs/design_notes.md
  • Ionia: High-Performance Replication for Modern Disk-based KV Stores
    • Replication approaches
      • https://www.usenix.org/system/files/fast24-xu.pdf
  • DeepSeek HAI
    • https://arxiv.org/pdf/2408.14158v1

Filesystem

  • Interface through VFS and FUSE
  • FUSE allows to mount different storages conveniently: S3-style, SSHed server, zip folder
  • FUSE problems:
    • memory-copy overhead (userspace -> kernel -> userspace)
    • primitive multi-threading support

Distributed filesystems

  • Larger block sizes - 128MB (HDFS), 4MB (JuiceFS)
    • => smaller block writes are possible
  • System idea
    • Each node runs a daemon (type depends on node type): data node, metadata node, etc.
    • Several node types: client, data nodes (store data), metadata nodes (store metadata), sometimes management mode
    • Data is usually broken up into chunks
      • "The design goal of chunk storage system is to achieve the highest bandwidth possible even when there are storage medium failures"
      • Smaller block sizes are actually better + chunks must be distributed among nodes
        • DeepSeek about these problems: https://www.high-flyer.cn/blog/3fs-4/
    • Support for node cache stored in memory.
    • Support for replication
  • Problems
    • How to organize network communication?
      • Specifically: moving in network is expensive
      • MapReduce hierarchy + GFS. store locally to map
      • Solution: Fat-Trees https://courses.csail.mit.edu/6.896/spring04/handouts/papers/fat_trees.pdf
    • How to avoid copies?
      • Solution: RDMA
      • 3FS experience: https://www.high-flyer.cn/blog/3fs-2/
      • experimental RDMA read size 64KB

3FS

  • Distributed filesystem
  • Hardware
3FS Storage Node Hardware: In Fire-Flyer 2 AI-HPC,
we deployed 180 storage nodes, as shown in Table IV, each
node contains 16 PCIe 4.0 NVMe SSDs and 2 Mellanox CX6
200Gbps InfiniBand HCAs. With totally 360 * 200Gbps out-
bound InfiniBand HCAs, the system can total provide 9TB/s
outbound bandwidh, and we actually achieved total read
throughput of 8TB/s. The total 2880 NVMe SSDs provide
over 20PiB storage space with an mirror data redundancy
  • Nodes
    • 3 types of nodes: storage, meta (stores data locations and metadata), management (information about nodes)
  • Replication
    • Each storage node
    • Chained replication
      • have a chain of nodes. write goes into HEAD, reads always go to TAIL
      • gives consistency and linearizability for free
      • tradeoff: write latency suitable for read-heavy workloads.
    • 3FS uses CRAQ
      • Chain Replication with Apportioned Queries
        • main difference from CR : serve reads from any chain node BUT preserve strong consistency
      • Idea:
        • write:
          • accept new write in HEAD
          • propagate new value to TAIL, while writing the "dirty" version on every intermediate
          • when TAIL ACKs the value, propagate ACK to all nodes. the "dirty" value becomes "committed"
        • read:
          • if the latest version is "committed" => serve directly
          • if the latest version is "dirty" => ask TAIL for the latest committed version
            • this means that every intermediate node stores a reference to TAIL
      • "The meta service selects an offset in the chain table and a stripe size k for each file. The file chunks are assigned to the next k chains starting at the offset"
      • Implementation details:
        • each SSD participates in several chains (each SSD is storage target).
  • TODO

Object storages

  • Optimized for large files
  • Simpler interface: no things like atomic renames, file locking, directory hierarchy

Distributed job orchestration

  • Components
    • Task executor
      • daemon (kubelet or yarn's NodeManager)
      • runs job tasks, sends heartbeats, tracks task status
      • responsible for resource and performance isolation
    • Resource manager
      • centralized
        • but can be scaled with coordination services (kubernetes uses etcd, YARN uses ZooKeeper)
      • stores cluster state as metadata about each node (hardware, network, node/task status, etc.)
    • Scheduler
      • centralized
      • processes job start requests. example of workflow
        • client sends "start this job with 10 tasks on this Docker image"
        • scheduler forms a request
        • scheduler gets the state of resource manager
        • scheduler determines what nodes to run on
        • scheduler informs task executors that they need to run task X
      • some scheduling heuristics need to be application-specific. in Kubernetes it's operators, in YARN ApplicationMasters

Types of job schedulers

  • Resource allocation scheduler
    • example: Kubernetes
      • all requests enter a queue
      • uses heuristics and rules on how to allocate limited resources to incoming jobs
    • example: DeepSeek HAI platform
  • Workflow scheduler (dataflow)
    • example: Airflow
    • idea: have a DAG

Use cases for DFS

  • GFS in MapReduce
  • 3FS in HAI training
    • need to perform checkpoints during model training
    • tasks that we run are Python scripts that MUST adhere to cluster resource management mechanisms:
      • accepting interruption signal from the cluster manager
      • save checkpoint (model weights, optimizer params, etc.)
      • notify the cluster of interruptions
      • recover from the checkpoint and run
    • when checkpoint interrupt is received
      • params and optimization states are divided into chunks and written to 3FS using batch write API (~10GiB/s throughput)
      • asynchronously transfer data from GPU to CPU
    • checkpointing every 5 minutes (so max. 5 minutes of training can be lost)

Shuffling

  • Often operations within a framework require that processing happens over items with the same key.
  • Therefore, we need to move the data with the same key to same machines.
  • Typically, output of "mapping" task --- (key,value) pair --- is getting written to location partitioned on "reducer" it's assigned to (so on key).

Data organization

  • data mesh = each org. unit takes ownership over its data
  • data fabric = architectural approach (set of tools) for unifying access to an organization's data across many separate sources --- databases, data lakes, warehouses, streaming systems, SaaS applications.

DataFrames: Distributed computation

  • In local DataFrames the evaluation is eager. In Spark, you write your flow, then it builds a query, optimizes it, and only then executes

TODO:

  • Fat tree vs Dragonfly https://www.youtube.com/watch?v=cLSn7Q0QXG4
  • Metaflow Netflix. AI workflow orch : https://www.youtube.com/watch?v=JCbOI_1ZA5E&t=1032s