DEV Community

DASWU
DASWU

Posted on

JuiceFS Enterprise 5.4: From 100 Billion Files to 1 Million Clients

Following the support for 100-billion-file environments introduced in JuiceFS Enterprise Edition 5.3, v5.4 further improves JuiceFS across several dimensions for large-scale deployments. At this scale, even small per-operation overheads can add up to significant resource consumption. Once metadata is distributed across multiple zones, the system also needs to coordinate across nodes while maintaining performance, consistency, and stability.

To address these challenges, JuiceFS Enterprise Edition 5.4 introduces the following new capabilities and improvements:

  • Support for up to 1 million simultaneously mounted clients
  • Fast cloning of directories containing massive files
  • More than 90% RDMA bandwidth utilization
  • Up to 100% higher random-read performance under high concurrency
  • On-demand data synchronization for mirror file systems
  • More flexible caching strategies for a wider range of workloads
  • Multiple operational improvements for large-scale data

Supporting 1 million clients

As AI agents move toward large-scale deployment, the JuiceFS team faced a new challenge: supporting connections and concurrent access from up to 1 million clients on top of the metadata workload of a 100-billion-file file system. The team spent months exploring and optimizing the architecture, and introduced this capability in JuiceFS 5.4.

JuiceFS uses a distributed architecture that partitions metadata across multiple metadata nodes, allowing metadata for hundreds of billions of files to be managed across multiple zones. To allow these nodes to handle a large number of client connections at the same time, v5.4 co-locates metadata nodes with network proxies, enabling clients to connect to any metadata node.

The node that receives a request handles requests for its local partition directly. Requests involving other partitions are forwarded internally to the corresponding nodes. Multiple client sessions share the same physical connections between nodes, reducing the number of connections and the associated maintenance overhead.

Connection reuse is only part of the challenge. Querying large numbers of client sessions, reporting session status, and cleaning up expired sessions can also introduce significant overhead.

In v5.4, indexes are used to locate client sessions instead of scanning unrelated sessions. Paginated queries and sampled reporting reduce the overhead of session status management. Expired sessions are cleaned up in batches to prevent cleanup work from becoming concentrated in a single pass.

With these optimizations, JuiceFS 5.4 successfully handled 1 million simultaneous client connections in stress tests. Previously, when clients connected separately to multiple metadata nodes, the total number of connections could reach tens of millions. After optimization, 10 internal proxy nodes handled 1 million client connections, while connection reuse reduced the number of physical backend connections from the proxies to the metadata leader to about 400. This significantly reduced connection-management pressure on the leader.

Fast cloning for large directories

When creating branches of training data or preparing test environments, users often need an independently modifiable copy of a directory. As directory sizes grow, cloning becomes more complex: in a multi-zone architecture, the metadata for a single directory tree may span multiple metadata zones, so cloning must coordinate subtree operations across nodes while preserving the complete directory hierarchy.

v5.4 introduces cross-zone directory cloning. It recursively processes subtrees while preserving the original cross-zone layout of the directory tree. Cloning only copies metadata. The underlying object data remains shared, so there is no need to copy it. Subsequent writes to the clone do not affect the source files.

In a deployment with 30 metadata zones, cloning 100 million files took 1 minute and 40 seconds, averaging 1 million files per second. More zones can provide greater parallelism, while actual performance also depends on metadata load and how the directory tree is distributed across zones.

Note that cross-zone directory cloning is not atomic. If new data continues to be written to the source directory during cloning, different subtrees in the cloned directory may reflect the source directory at different points in time. If strict consistency is required, pause writes to the source directory before starting the clone.

Performance improvements

More than 90% RDMA bandwidth utilization

Starting with JuiceFS 5.3, clients can use RDMA to communicate with distributed cache nodes, reducing CPU overhead during data transfers and improving transfer efficiency. v5.4 further improves bandwidth utilization and transfer stability.

In a test environment with a 400 Gbps RDMA NIC on the client and an 800 Gbps RDMA NIC on the distributed cache node, a single client achieved 45 GB/s of read throughput. A single distributed cache node delivered 90 GB/s of data throughput, with NIC utilization exceeding 90%.

Higher random-read IOPS under high concurrency

In multi-process training, HPC, and rendering scenarios, high-concurrency random reads can stress not only storage and network resources, but also client lock contention and request-processing overhead.

Over the past six months, the JuiceFS team has continuously optimized high-concurrency random-read performance by simplifying frequently executed code paths and reducing lock contention.

As a result, random-read IOPS has more than doubled overall compared with the previous implementation. With distributed caching, a single client reached up to 189,000 IOPS. With local caching, a single client reached up to 156,000 IOPS.

Test environment: A dual-socket AMD server with 128 cores and 256 threads, 1 TB of DDR4 memory, a 100 Gbps NIC, and a 3.84 TB NVMe SSD used as the cache disk. The system ran Linux 5.14.

Optimizing data access across regions

On-demand synchronization to reduce data replication

In multi-region, multi-provider, and multi-cluster compute environments, a JuiceFS mirror file system can automatically synchronize files and directories from the source. Metadata is continuously synchronized to keep the directory structure and object versions up to date. Previously, object data was synchronized in full in the background by default.

However, a workload in the mirror region may only need a subset of the models, datasets, or historical directories stored at the source. For example, if the source contains 10 PB of data but the mirror accesses only 1% of it, full synchronization would transfer and store a large amount of data that the workload never uses.

v5.4 supports on-demand synchronization of object data for mirror file systems. When enabled, metadata continues to synchronize, while object data is no longer synchronized automatically in the background. When a client reads data, JuiceFS first checks the object storage in the mirror region. If the data has not been synchronized yet, JuiceFS fetches it from the source and stores it in the mirror on demand. This reduces unnecessary data transfer and storage consumption.

If an application requires local read performance on the first access, you can warm up the required data in advance. If a complete copy of the data is required, you can continue to use background synchronization.

Accelerating metadata access with read-only nodes

For workloads that only need to accelerate metadata access, v5.4 introduces a new mode that does not require creating a mirror file system.

Users can deploy read-only metadata service nodes in other regions to accelerate metadata reads for local clients. Clients continue to mount the source file system, while the acceleration nodes can scale dynamically to match access requirements.

More flexible cache management

Improving hot-data cache hit rates

Data from one-time scans is often unlikely to be accessed again. Caching this data can evict hot data that is accessed repeatedly. By default, when a request misses the distributed cache, the cache node fetches the data from object storage and stores it in the cache. v5.4 introduces a new cache-aside mode. When enabled, a cache miss causes the application node to read directly from object storage instead of filling the cache with the requested data.

Preserving cache hit rates when adding multiple cache nodes

When multiple cache nodes are added at once, read requests may be routed to new nodes that do not yet have the required data, triggering additional reads from object storage.

v5.4 introduces the --commission parameter. After a new node is enabled, it retains the node mapping from before scaling. If the data is not available in its local cache, the new node fetches it from the original node instead. This reduces reads from the source during cache scaling.

Reducing the number of files on large cache disks

On large-capacity cache disks, storing each cache block as a separate file can cause the number of files on the cache disk to grow continuously.

The new merge-cache mode combines multiple cache blocks into larger files and stores their indexes separately. This significantly reduces the number of files that the cache disk needs to manage.

This feature is currently available as a public beta.

Warming up specific byte ranges

If a workload reads only part of a large file, warming up the entire file consumes additional cache space and transfers unnecessary data.

v5.4 allows you to specify one or more byte ranges to warm up. JuiceFS only warms up cache blocks that intersect those ranges, aligning cache preparation with the data the workload actually needs to read.

This capability enables more precise warm-up for data formats such as Lance and Parquet.

Simplifying operations for massive files

Restoring files by deletion time

When restoring files from the trash, users may know approximately when an accidental deletion occurred but have difficulty identifying the target files by name or keyword alone.

The restore command now supports the start-time and end-time parameters. In addition to the existing keyword filters, these parameters allow users to restrict the restore operation by file deletion time.

Reducing GC and fsck memory usage

When checking large file systems, GC and fsck need to process large numbers of records, and sorting can consume significant amounts of memory.

v5.4 adds external sorting modes to both commands, using disk storage during sorting to reduce memory requirements. The external sorting mode for GC is available only with dry-run, allowing users to inspect and evaluate the operation without actually deleting data.

Restricting access with token allowlists

v5.4 supports access control through token allowlists. An allowlist can specify multiple first-level subdirectories under the root directory.

Clients using the token can only access and modify content within those subdirectories, allowing different teams or workloads to be restricted to their designated data.

Summary: scale changes everything

JuiceFS 5.3 brought the platform into a new phase of scale. At the same time, the rapid development of AI continues to push the boundaries of data volume and system load.

As file counts, client counts, and request volumes all reach new levels, performance, stability, consistency, and operational complexity are no longer independent concerns. They interact with one another and become increasingly difficult to manage as the system scales.

v5.4 addresses these system-level engineering challenges across large-scale deployments, covering client connectivity, cross-zone directories, data access, mirror synchronization, and cache management. JuiceFS Cloud Service users can now try JuiceFS Enterprise Edition 5.4 online. Users running on-premises deployments can contact the JuiceFS team for upgrade support.

JuiceFS is used across a range of demanding workloads, including LLMs and multimodal models, autonomous driving, embodied AI, and quantitative investment. We’ll continue working with organizations at the forefront of these fields to improve JuiceFS and build infrastructure that can keep pace with their applications.

If you have any feedback on this article or ideas to share, we invite you to participate in the discussions on GitHub and join our community on Discord.

Top comments (0)