[
https://issues.apache.org/jira/browse/HDFS-17973?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
ASF GitHub Bot updated HDFS-17973:
----------------------------------
Labels: pull-request-available (was: )
> HDFS DataNode Write Batching for Efficient Node-Level Scaling
> -------------------------------------------------------------
>
> Key: HDFS-17973
> URL: https://issues.apache.org/jira/browse/HDFS-17973
> Project: Hadoop HDFS
> Issue Type: Improvement
> Components: hdfs
> Reporter: Rajan Dhabalia
> Priority: Major
> Labels: pull-request-available
>
> Below is the same content reformatted for direct copy/paste into Jira wiki
> markup. I have not changed the content.
> h2. Summary
> Introduce an optional in-memory write buffering mechanism in the HDFS
> DataNode to improve per-node I/O efficiency, increase throughput, and reduce
> read/write tail latency.
> The optimization batches small DataNode write packets into larger disk
> writes, uses direct I/O to reduce page-cache pressure, and controls flush
> concurrency to reduce disk contention.
> The feature is configuration-gated and disabled by default.
> h2. Issue Type
> Improvement
> h2. Component
> DataNode
> h3. Motivation
> HDFS blocks are typically large, but the DataNode processes block data as
> much smaller packets. Multiple threads can concurrently process packets
> belonging to different blocks and files.
> While each individual block has a sequential access pattern, packet-level
> interleaving across multiple blocks and files can result in a less sequential
> access pattern at the storage layer.
> This is particularly impactful on HDD-based DataNodes, where seek and
> rotational latency make random I/O significantly more expensive than
> sequential I/O.
> The current architecture can result in:
> * Frequent small disk writes.
> * Increased write amplification and I/O operations.
> * Randomized physical I/O due to packet interleaving.
> * Dirty write pages consuming OS page-cache memory.
> * Eviction of useful read and read-ahead pages.
> * Increased read/write contention.
> * Higher read latency and write tail latency.
> * Lower effective disk and DataNode throughput.
> The goal of this change is to improve {*}vertical efficiency{*}, allowing an
> individual DataNode to achieve higher throughput while providing more
> predictable performance for mixed read/write workloads.
> h1. Current Data Plane
> h3. Block-Level Packetization
> Large HDFS blocks are processed as smaller packets:
> {code:java}
> HDFS Block
> |
> v
> +------+------+------+------+------+
> | P1 | P2 | P3 | P4 | ... |
> +------+------+------+------+------+
> {code}
> Multiple threads can process packets from different blocks/files concurrently:
> {code:java}
> Thread 1 -> File A -> P1 -> P2 -> P3
> Thread 2 -> File B -> P1 -> P2 -> P3
> Thread 3 -> File C -> P1 -> P2 -> P3
> {code}
> This packet-level interleaving can result in non-sequential physical I/O,
> particularly when multiple files are active simultaneously.
> h3. Page Cache
> With the existing buffered write path:
> {code:java}
> DataNode
> |
> v
> OS Page Cache
> |
> v
> Kernel Writeback
> |
> v
> Disk
> {code}
> Writes create dirty pages that compete for memory with:
> Read pages.
> Read-ahead pages.
> Other filesystem cache.
> Heavy write workloads can therefore reduce the effectiveness of the page
> cache for reads.
> h3. HDD Behavior
> HDDs perform best with sequential access. Interleaved requests targeting
> different regions can cause additional head movement and seek overhead.
> As a result, small packet writes combined with concurrent reads can
> significantly reduce effective disk throughput and increase read latency.
> h1. Proposed Data Plane
> The proposed architecture introduces a DataNode-managed write buffer:
> {code:java}
> DataNode
> |
> Small packets
> |
> v
> +-------------------+
> | In-memory buffer |
> +---------+---------+
> |
> Large batched I/O
> |
> v
> O_DIRECT
> |
> v
> Disk
> {code}
> h2. Key Changes
> h3. 1. Large In-Memory Write Buffers
> Accumulate multiple small packets in memory before flushing them to disk.
> This converts:
> {code:java}
> Small writes:
> P1 -> P2 -> P3 -> P4 -> P5 -> ...
> {code}
> into:
> {code:java}
> Large writes:
> +-------------------+
> | P1 P2 P3 ... Pn |
> +-------------------+
> |
> v
> Disk
> {code}
>
> Benefits include:
> - Fewer disk I/O operations.
> - Reduced write amplification.
> - Larger sequential writes.
> - Better utilization of disk bandwidth.
> - Reduced physical I/O fragmentation.
> - The buffer size is configurable.
> h3. 2. Direct I/O for Buffered Writes
> Use
> {code:java}
> O_DIRECT{code}
> when flushing buffered writes to bypass the OS page cache.
> This allows the DataNode to explicitly manage write buffering while
> preserving OS page-cache capacity for reads.
> Benefits:
> - Reduces dirty-page pressure.
> - Prevents DataNode writes from unnecessarily consuming read-cache capacity.
> - Improves read-ahead effectiveness.
> - Reduces read/write memory contention.
> h3. 3. Controlled Flush Concurrency
> Large writes alone are not sufficient. Excessive concurrent flushes can still
> create random I/O and disk contention. The implementation therefore limits
> concurrent flush bytes per volume.
> The goal is to find a balance between:
> - Disk parallelism.
> - Sequential I/O.
> - Disk queue depth.
> - Read latency.
> h3. 4. Last-Replica-Only Buffering
> Provide an option to enable buffering only on the last DataNode in the
> replication pipeline.
> This allows the optimization to be introduced selectively without changing
> the write behavior of every replica.
> h3. 5. Bounded Memory Usage
> The buffering mechanism provides configurable limits for:
> - Total DataNode buffer capacity.
> - Per-block buffer size.
> - Minimum number of volumes.
> - Per-volume concurrent flush capacity.
> This prevents unbounded memory consumption or excessive disk pressure.
> h3. 6. Idle Flush
> Partially filled buffers are flushed after a configurable idle timeout.
> This prevents small writes from being held indefinitely when the workload
> does not generate enough packets to fill the buffer.
> h1. Read Path Optimization
> The write buffer is designed to reduce interference with the read path.
> With DataNode writes bypassing the OS page cache:
> {code:java}
> System Memory
> |
> +------+------+
> | |
> v v
> Read Cache Read-ahead
> | |
> +------+------+
> |
> v
> Disk
> DataNode writes
> |
> v
> Memory Buffer
> |
> v
> O_DIRECT
> |
> v
> Disk
> {code}
> This allows the page cache to be used more effectively for reads and
> read-ahead.
> Expected benefits include:
> - Reduced eviction of useful read pages.
> - Better read-ahead effectiveness.
> - Lower read latency.
> - Reduced read/write contention.
> - More predictable mixed-workload performance.
> h1. Recommended Kernel tuning with this change
> The following kernel settings can complement the DataNode optimization:
> {code:java}
> Read-ahead and larger block-layer requests
> /sys/block/sd${disk}/queue/read_ahead_kb => 4096
> /sys/block/sd${disk}/queue/max_sectors_kb => 4096 {code}
> h3. Rationale
> Read-ahead
> Larger read-ahead allows the kernel to fetch larger sequential ranges and
> improves sequential-read efficiency.
> Maximum request size
> A larger
> {code:java}
> max_sectors_kb{code}
> allows larger requests to be issued to the block layer where supported.
> h1. Configuration
> ||Property||Default||Description||
> |dfs.datanode.write.memory.buffer.enabled|false|Enables DataNode write memory
> buffering.|
> |dfs.datanode.write.memory.buffer.last-replica-only|true|Buffers writes only
> on the last replica.|
> |dfs.datanode.write.memory.buffer.max.capacity.mb|See docs|Maximum total
> in-memory buffer capacity.|
> |dfs.datanode.write.memory.buffer.min.volumes|See docs|Minimum number of
> volumes required.|
> |dfs.datanode.write.buffer.size.bytes|See docs|Per-block write buffer size.|
> |dfs.datanode.write.buffer.idle.flush.timeout.ms|See docs|Idle timeout for
> partially filled buffers.|
> |dfs.datanode.concurrent.flush.mb.per.volume|See docs|Maximum concurrent
> flush bytes per volume.|
> |dfs.datanode.read.ahead.cache.bytes.threshold|See docs|Read-ahead cache
> threshold.|
> h1. Benchmark Results
> Initial benchmarking demonstrates significant improvement in DataNode
> vertical efficiency.
> ||Metric||Improvement||
> |Mixed 80/20 read/write throughput|~30% higher|
> |Read-only throughput|Up to ~45% higher|
> |Large-block P99 write latency|~80% lower|
> |Read latency under concurrent writes|Significant reduction|
> The results indicate that larger buffered writes and controlled flushing can
> improve disk utilization while reducing interference between concurrent reads
> and writes.
> h1. Expected Impact
> For workloads where storage I/O is the primary bottleneck, the expected
> improvements are:
> - 25–40% higher throughput for write-heavy workloads.
> - Up to ~45% higher throughput for read-heavy workloads.
> - ~80% lower large-block write P99 latency.
> - Lower read latency under concurrent writes.
> - Reduced random I/O and disk seek contention.
> - Better utilization of available disk bandwidth.
> - More predictable mixed read/write performance.
> The actual improvement will depend on workload characteristics, including
> read/write ratio, I/O concurrency, packet interleaving, disk utilization, and
> underlying storage characteristics.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]