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
- Replication approaches
- 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
- How to organize network communication?
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
- write:
- "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).
- Chain Replication with Apportioned Queries
- 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 (
kubeletor yarn's NodeManager) - runs job tasks, sends heartbeats, tracks task status
- responsible for resource and performance isolation
- daemon (
- 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.)
- centralized
- 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
- Task executor
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
- example: Kubernetes
- 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
DataFramesthe 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