rdhabalia opened a new pull request, #8718:
URL: https://github.com/apache/hadoop/pull/8718
### Description of PR
A block written to a DataNode remains active until the client sends the final
packet and the block is finalized.
When a client crashes, loses network connectivity, or hangs before sending
the
last packet, the DataNode transfer thread can remain blocked indefinitely
while
waiting for input. Over time, these stalled transfers can accumulate and
cause:
- Transfer-thread starvation
- Resource leakage
- Degraded DataNode responsiveness
Today, there is no built-in inactivity detection to reclaim resources from
stalled transfers. The system relies on client behavior or eventual lease
recovery, which may occur long after the actual failure.
## Solution
Introduce a configurable, inactivity-based timeout for block transfers on the
DataNode.
When enabled, the DataNode tracks packet-arrival activity for each ongoing
block write using a monotonic clock. A single scheduled check per transfer
runs every `timeout / 2` and compares the time since the last received packet
against `timeout / 2`.
### Transfer Activity Detection
- If a packet arrived within `timeout / 2`:
- The stream is considered active.
- The check is rescheduled.
- This provides an automatic reset for healthy streams.
- Otherwise:
- The transfer is considered stalled.
- The transfer is aborted.
Polling every `timeout / 2` bounds the detection latency to between
`timeout / 2` and `timeout` after the client stops sending.
### Safe Transfer Abort
The abort runs on the shared scheduler thread and deliberately performs **no
disk I/O**.
The on-disk streams (`out` / `checksumOut`) are written by the receive thread
**without holding the `BlockReceiver` monitor**. Flushing these streams from
the scheduler thread could race with those writes and potentially corrupt the
replica.
Instead, the scheduler:
1. Closes only the client input stream.
2. This unblocks the receive thread's blocked socket read.
3. The receive thread unwinds and performs its own single-threaded cleanup.
Data already acknowledged to the client is already durable on disk. The
replica remains in the **RBW (Replica Being Written)** state, allowing the
NameNode's existing lease-recovery / block-synchronization path to finalize
the block without data loss.
## Scheduler Design
The timeout checks use a shared `ScheduledThreadPoolExecutor` with **2 daemon
threads**.
The scheduler:
- Is created lazily at DataNode startup only when the timeout is configured.
- Adds **zero threads and zero overhead when the feature is disabled**.
- Uses `remove-on-cancel` so per-block checks cancelled when a transfer
completes do not accumulate in the delay queue.
- Does not execute delayed tasks after shutdown.
The last-packet timestamp is updated only when the timeout feature is
enabled,
avoiding unnecessary atomic writes on the hot path when the feature is
disabled.
## Client Socket Timeout Considerations
Because a stall is declared after no packet is received for `timeout / 2`,
the DataNode timeout must be configured comfortably larger than the client
socket read timeout.
A healthy idle `hflush` / `hsync` stream still sends heartbeat packets
roughly
every half of the socket timeout.
The DataNode logs a warning if the configured transfer timeout is not
sufficiently larger than the client socket timeout.
## Configuration
| Configuration | Default | Description |
|---|---:|---|
| `dfs.datanode.last.packet.receive.timeout.ms` | `0` | Last-packet
inactivity timeout in milliseconds. `0` disables the feature. |
A typical enabled value is:
```text
dfs.datanode.last.packet.receive.timeout.ms=600000
````
This corresponds to a **10-minute timeout**.
The change is fully backward compatible and disabled by default.
## Testing
### `TestBlockReceiverLastPacketTimeout`
Verifies the DataNode-side scheduler wiring:
* The shared scheduler is created only when the timeout is enabled.
* `remove-on-cancel` is configured correctly.
* Delayed tasks do not execute after scheduler shutdown.
### `TestBlockReceiverTransferTimeout`
Runs an end-to-end test using `MiniDFSCluster` and verifies that a genuinely
stuck client is deterministically aborted.
The test uses a single DataNode with DataNode replacement disabled so that a
broken pipeline cannot be silently recovered.
It also verifies that healthy streams are never falsely aborted, including:
* Full transfers
* Slow-but-steady writers
* Writers whose packet gaps remain just below the timeout threshold, even
when the total transfer time exceeds the timeout
* Disabled timeout configuration
* Multiple concurrent transfers
### How was this patch tested?
Tested by newly added unit test
### For code changes:
- [ ] Does the title of this PR start with the corresponding JIRA issue id
(e.g. 'HADOOP-17799. Your PR title ...')?
- [ ] Object storage: Have the integration tests been executed and the
endpoint
declared according to the connector-specific documentation? *Note:
Automated CI
testing doesn't cover all cases so manual testing with cloud storage
is still
required.*
- [ ] If adding new dependencies to the code, are these dependencies
licensed in a way that is compatible for inclusion under [ASF
2.0](http://www.apache.org/legal/resolved.html#category-a)?
- [ ] If applicable, have you updated the `LICENSE`, `LICENSE-binary`,
`NOTICE-binary` files?
### AI Tooling
If an AI tool was used:
- [ ] The PR includes the phrase "Contains content generated by <tool>"
where <tool> is the name of the AI tool used.
- [ ] My use of AI contributions follows the ASF legal policy
https://www.apache.org/legal/generative-tooling.html
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]