This is an automated email from the ASF dual-hosted git repository.
Gargi-jais11 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ozone.git
The following commit(s) were added to refs/heads/master by this push:
new 66447d3663a HDDS-15690. Support Node-ID for DiskBalancer Commands
(#10755).
66447d3663a is described below
commit 66447d3663a73f4bc80e335d67c6679e6b3b4ecf
Author: Gargi Jaiswal <[email protected]>
AuthorDate: Mon Jul 20 16:41:48 2026 +0530
HDDS-15690. Support Node-ID for DiskBalancer Commands (#10755).
---
hadoop-hdds/docs/content/design/diskbalancer.md | 12 +-
hadoop-hdds/docs/content/feature/DiskBalancer.md | 37 ++-
.../docs/content/feature/DiskBalancer.zh.md | 34 +-
.../datanode/AbstractDiskBalancerSubCommand.java | 162 +++++++++-
.../hdds/scm/cli/datanode/DatanodeParameters.java | 3 +-
.../scm/cli/datanode/DiskBalancerCommands.java | 31 +-
.../cli/datanode/DiskBalancerCommonOptions.java | 12 +
.../cli/datanode/DiskBalancerReportSubcommand.java | 35 ++-
.../cli/datanode/DiskBalancerStatusSubcommand.java | 46 +--
.../cli/datanode/DiskBalancerSubCommandUtil.java | 139 ++++++++-
.../datanode/TestDiskBalancerSubCommandUtil.java | 135 ++++++++
.../cli/datanode/TestDiskBalancerSubCommands.java | 345 +++++++++++++++++++--
12 files changed, 855 insertions(+), 136 deletions(-)
diff --git a/hadoop-hdds/docs/content/design/diskbalancer.md
b/hadoop-hdds/docs/content/design/diskbalancer.md
index 05ed92e77be..299ff689950 100644
--- a/hadoop-hdds/docs/content/design/diskbalancer.md
+++ b/hadoop-hdds/docs/content/design/diskbalancer.md
@@ -70,6 +70,7 @@ Administrators use the `ozone admin datanode diskbalancer`
CLI to manage and mon
- Query DiskBalancer status and volume density reports
* Each datanode performs its own **authentication** (via RPC) and
**authorization** checks (using `OzoneAdmins` based on `ozone.administrators`
configuration).
* For batch operations, clients can use the `--in-service-datanodes` flag to
automatically query SCM for all IN_SERVICE and HEALTHY datanodes and execute
commands on all of them.
+* When `--node-id` is used, the CLI resolves the datanode UUID to the
CLIENT_RPC address through SCM before issuing the datanode RPC.
**DN - DiskBalancer Service:**
@@ -81,8 +82,10 @@ A daemon thread, the **Scheduler**, runs periodically on
each Datanode.
from the most over-utilized disk (source) to the least utilized disk
(destination).
3. The scheduler dispatches these move tasks to a pool of **Worker** threads
for parallel execution.
-**Note:** SCM is used **only** for datanode discovery when using the
`--in-service-datanodes` flag. SCM provides a list of IN_SERVICE and HEALTHY
datanodes for batch operations but
-does **not** participate in DiskBalancer control operations
(start/stop/update/status/report). All DiskBalancer operations are performed
directly between client and datanode.
+**Note:** SCM is used for datanode discovery when using the
`--in-service-datanodes` flag and to resolve
+`--node-id` UUIDs to CLIENT_RPC addresses. SCM provides a list of IN_SERVICE
and HEALTHY datanodes for
+batch operations and node metadata for UUID lookup, but does **not**
participate in DiskBalancer control operations
+(start/stop/update/status/report). All DiskBalancer operations are performed
directly between client and datanode.
## Container Move Process
@@ -149,10 +152,11 @@ The DiskBalancer CLI provides five main commands that
communicate directly with
5. **report** - Retrieves volume density report showing imbalance analysis.
The CLI supports:
-- **Direct datanode addressing**: Commands can target specific datanodes by
hostname or IP address
+- **Direct datanode addressing**: Commands can target specific datanodes by
hostname or IP address as positional arguments
+- **UUID targeting**: `--node-id` resolves datanode UUIDs to CLIENT_RPC
addresses through SCM; resolution failures are reported per node
- **Batch operations**: The `--in-service-datanodes` flag queries SCM for all
IN_SERVICE and HEALTHY datanodes and executes commands on all of them
- **Flexible input**: Datanode addresses can be provided as positional
arguments or read from stdin
-- **Output formats**: Results can be displayed in human-readable format or
JSON for programmatic access
+- **Output formats**: Results can be displayed in human-readable format or
JSON for programmatic access; hostname targets show `hostname (ip:port)`,
`--node-id` targets show the UUID
### Operational State Awareness
diff --git a/hadoop-hdds/docs/content/feature/DiskBalancer.md
b/hadoop-hdds/docs/content/feature/DiskBalancer.md
index b88f1710028..d584408dfef 100644
--- a/hadoop-hdds/docs/content/feature/DiskBalancer.md
+++ b/hadoop-hdds/docs/content/feature/DiskBalancer.md
@@ -126,41 +126,42 @@ The DiskBalancer is managed through the `ozone admin
datanode diskbalancer` comm
**Start DiskBalancer:**
```bash
-ozone admin datanode diskbalancer start [<datanode-address> ...] [OPTIONS]
[--in-service-datanodes]
+ozone admin datanode diskbalancer start [<datanode-address or -id> ...]
[OPTIONS] [--in-service-datanodes]
```
**Stop DiskBalancer:**
```bash
-ozone admin datanode diskbalancer stop [<datanode-address> ...]
[--in-service-datanodes]
+ozone admin datanode diskbalancer stop [<datanode-address or -id> ...]
[--in-service-datanodes]
```
**Update Configuration:**
```bash
-ozone admin datanode diskbalancer update [<datanode-address> ...] [OPTIONS]
[--in-service-datanodes]
+ozone admin datanode diskbalancer update [<datanode-address or -id> ...]
[OPTIONS] [--in-service-datanodes]
```
**Get Status:**
```bash
-ozone admin datanode diskbalancer status [<datanode-address> ...]
[--in-service-datanodes] [--json]
+ozone admin datanode diskbalancer status [<datanode-address or -id> ...]
[--in-service-datanodes] [--json]
```
**Get Report:**
```bash
-ozone admin datanode diskbalancer report [<datanode-address> ...]
[--in-service-datanodes] [--json]
+ozone admin datanode diskbalancer report [<datanode-address or -id> ...]
[--in-service-datanodes] [--json]
```
### Command Options
-| Option | Description
| Example |
-|-------------------------------------|-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|------------------------------------------------|
-| `<datanode-address>` | One or more datanode addresses as
positional arguments. Addresses can be:<br>- Hostname (e.g., `DN-1`) - uses
default CLIENT_RPC port (19864)<br>- Hostname with port (e.g.,
`DN-1:19864`)<br>- IP address (e.g., `192.168.1.10`)<br>- IP address with port
(e.g., `192.168.1.10:19864`)<br>- Stdin (`-`) - reads datanode addresses from
standard input, one per line | `DN-1`<br>`DN-1:19864`<br>`192.168.1.10`<br>`-` |
-| `--in-service-datanodes` | It queries SCM for all IN_SERVICE and
HEALTHY datanodes and executes the command on all of them.
| `--in-service-datanodes` |
-| `--json` | Format output as JSON.
| `--json` |
-| `-t/--threshold-percentage` | Volume density threshold percentage
(default: 10.0). Used with `start` and `update` commands.
| `-t 5`<br>`--threshold-percentage 5.0` |
-| `-b/--bandwidth-in-mb` | Maximum disk bandwidth in MB/s
(default: 10). Used with `start` and `update` commands.
| `-b 20`<br>`--bandwidth-in-mb 50` |
-| `-p/--parallel-thread` | Number of parallel threads (default:
5). Used with `start` and `update` commands.
| `-p 5`<br>`--parallel-thread 10` |
-| `-s/--stop-after-disk-even` | Stop automatically after disks are
balanced (default: true). Used with `start` and `update` commands.
| `-s false`<br>`--stop-after-disk-even true` |
-| `-c/--container-states` | Comma-separated container lifecycle
state names that may be moved between disks. Used with `start` and `update`
commands.
| `-c
CLOSED,QUASI_CLOSED`<br>`--container-states CLOSED` |
+| Option | Description
|
Example |
+|-----------------------------|-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|------------------------------------------------|
+| `<datanode-address>` | One or more datanode addresses as positional
arguments. Each can be:<br>- Hostname (e.g., `DN-1`) - uses default CLIENT_RPC
port (19864)<br>- Hostname with port (e.g., `DN-1:19864`)<br>- IP address
(e.g., `192.168.1.10`)<br>- IP address with port (e.g.,
`192.168.1.10:19864`)<br>- Stdin (`-`) - reads addresses from standard input,
one per line | `DN-1`<br>`DN-1:19864`<br>`192.168.1.10`<br>`-` |
+| `--node-id` | Datanode UUID to target. Requires SCM to
resolve the UUID to a CLIENT_RPC address.
| `--node-id
a3b63511-bdf8-4fa1-8ab6-d19c0e806f84` |
+| `--in-service-datanodes` | It queries SCM for all IN_SERVICE and HEALTHY
datanodes and executes the command on all of them.
|
`--in-service-datanodes` |
+| `--json` | Format output as JSON.
|
`--json` |
+| `-t/--threshold-percentage` | Volume density threshold percentage (default:
10.0). Used with `start` and `update` commands.
| `-t
5`<br>`--threshold-percentage 5.0` |
+| `-b/--bandwidth-in-mb` | Maximum disk bandwidth in MB/s (default: 10).
Used with `start` and `update` commands.
| `-b
20`<br>`--bandwidth-in-mb 50` |
+| `-p/--parallel-thread` | Number of parallel threads (default: 5). Used
with `start` and `update` commands.
| `-p
5`<br>`--parallel-thread 10` |
+| `-s/--stop-after-disk-even` | Stop automatically after disks are balanced
(default: true). Used with `start` and `update` commands.
| `-s
false`<br>`--stop-after-disk-even true` |
+| `-c/--container-states` | Comma-separated container lifecycle state
names that may be moved between disks. Used with `start` and `update` commands.
| `-c CLOSED,QUASI_CLOSED`<br>`--container-states
CLOSED` |
### Examples
@@ -169,6 +170,9 @@ ozone admin datanode diskbalancer report
[<datanode-address> ...] [--in-service-
# Start DiskBalancer on multiple datanodes
ozone admin datanode diskbalancer start DN-1 DN-2 DN-3
+# Start DiskBalancer using a datanode UUID
+ozone admin datanode diskbalancer start --node-id
a3b63511-bdf8-4fa1-8ab6-d19c0e806f84
+
# Start DiskBalancer on all IN_SERVICE and HEALTHY datanodes
ozone admin datanode diskbalancer start --in-service-datanodes
@@ -216,6 +220,9 @@ ozone admin datanode diskbalancer update DN-1 -b 50 --json
# Get status from multiple datanodes
ozone admin datanode diskbalancer status DN-1 DN-2 DN-3
+# Get status using a datanode UUID
+ozone admin datanode diskbalancer status --node-id
a3b63511-bdf8-4fa1-8ab6-d19c0e806f84
+
# Get status from all IN_SERVICE and HEALTHY datanodes
ozone admin datanode diskbalancer status --in-service-datanodes
diff --git a/hadoop-hdds/docs/content/feature/DiskBalancer.zh.md
b/hadoop-hdds/docs/content/feature/DiskBalancer.zh.md
index 585d1535d29..fe7305dc04b 100644
--- a/hadoop-hdds/docs/content/feature/DiskBalancer.zh.md
+++ b/hadoop-hdds/docs/content/feature/DiskBalancer.zh.md
@@ -121,41 +121,42 @@ DiskBalancer 通过 `ozone admin datanode diskbalancer`
命令进行管理。
### 命令语法
**启动 DiskBalancer:**
```bash
-ozone admin datanode diskbalancer start [<datanode-address> ...] [OPTIONS]
[--in-service-datanodes]
+ozone admin datanode diskbalancer start [<datanode-address or -id> ...]
[OPTIONS] [--in-service-datanodes]
```
**停止 DiskBalancer:**
```bash
-ozone admin datanode diskbalancer stop [<datanode-address> ...]
[--in-service-datanodes]
+ozone admin datanode diskbalancer stop [<datanode-address or -id> ...]
[--in-service-datanodes]
```
**更新配置:**
```bash
-ozone admin datanode diskbalancer update [<datanode-address> ...] [OPTIONS]
[--in-service-datanodes]
+ozone admin datanode diskbalancer update [<datanode-address or -id> ...]
[OPTIONS] [--in-service-datanodes]
```
**获取状态:**
```bash
-ozone admin datanode diskbalancer status [<datanode-address> ...]
[--in-service-datanodes] [--json]
+ozone admin datanode diskbalancer status [<datanode-address or -id> ...]
[--in-service-datanodes] [--json]
```
**获取报告:**
```bash
-ozone admin datanode diskbalancer report [<datanode-address> ...]
[--in-service-datanodes] [--json]
+ozone admin datanode diskbalancer report [<datanode-address or -id> ...]
[--in-service-datanodes] [--json]
```
### 命令选项
-| Option | Description
| Example |
-|-------------------------------------|------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|------------------------------------------------|
-| `<datanode-address>` | 一个或多个数据节点地址作为位置参数。地址可以是:<br>-
主机名(例如,`DN-1`)- 使用默认的 CLIENT_RPC 端口 (19864)<br>- 带端口的主机名(例如,`DN-1:19864`)<br>-
IP 地址(例如,`192.168.1.10`)<br>- 带端口的 IP 地址(例如,`192.168.1.10:19864`)<br>- 标准输入
(`-`) - 从标准输入读取数据节点地址,每行一个 | `DN-1`<br>`DN-1:19864`<br>`192.168.1.10`<br>`-` |
-| `--in-service-datanodes` | 它向 SCM 查询所有 IN_SERVICE 且 HEALTHY
的数据节点,并在所有这些数据节点上执行该命令。
| `--in-service-datanodes` |
-| `--json` | 输出格式设置为JSON。
| `--json` |
-| `-t/--threshold-percentage` | 磁盘使用率阈值百分比(默认值:10.0)。与 `start` 和
`update` 命令一起使用。
| `-t 5`<br>`--threshold-percentage 5.0` |
-| `-b/--bandwidth-in-mb` | 最大磁盘带宽,单位为 MB/s(默认值:10)。与 `start` 和
`update` 命令一起使用。
| `-b 20`<br>`--bandwidth-in-mb 50` |
-| `-p/--parallel-thread` | 并行线程数(默认值:5)。与 `start` 和 `update`
命令一起使用。
| `-p 5`<br>`--parallel-thread 10` |
-| `-s/--stop-after-disk-even` | 磁盘平衡完成后自动停止(默认值:true)。与 `start` 和
`update` 命令一起使用。
| `-s false`<br>`--stop-after-disk-even true` |
-| `-c/--container-states` | 以逗号分隔的容器生命周期状态名称,表示可在磁盘之间移动的状态。配合
`start` 和 `update` 命令使用。
| `-c
CLOSED,QUASI_CLOSED`<br>`--container-states CLOSED` |
+| Option | Description
| Example |
+|-----------------------------|------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|------------------------------------------------|
+| `<datanode-address>` | 一个或多个数据节点地址作为位置参数。每个可以是:<br>- 主机名(例如 `DN-1`)-
使用默认 CLIENT_RPC 端口 (19864)<br>- 带端口的主机名(例如 `DN-1:19864`)<br>- IP 地址(例如
`192.168.1.10`)<br>- 带端口的 IP 地址(例如 `192.168.1.10:19864`)<br>- 标准输入 (`-`) -
从标准输入读取地址,每行一个 | `DN-1`<br>`DN-1:19864`<br>`192.168.1.10`<br>`-` |
+| `--node-id` | 数据节点 UUID。需要通过 SCM 解析为 CLIENT_RPC 地址。
| `--node-id
a3b63511-bdf8-4fa1-8ab6-d19c0e806f84` |
+| `--in-service-datanodes` | 它向 SCM 查询所有 IN_SERVICE 且 HEALTHY
的数据节点,并在所有这些数据节点上执行该命令。
| `--in-service-datanodes` |
+| `--json` | 输出格式设置为JSON。
| `--json` |
+| `-t/--threshold-percentage` | 磁盘使用率阈值百分比(默认值:10.0)。与 `start` 和 `update`
命令一起使用。
| `-t 5`<br>`--threshold-percentage 5.0` |
+| `-b/--bandwidth-in-mb` | 最大磁盘带宽,单位为 MB/s(默认值:10)。与 `start` 和 `update`
命令一起使用。
| `-b 20`<br>`--bandwidth-in-mb 50` |
+| `-p/--parallel-thread` | 并行线程数(默认值:5)。与 `start` 和 `update` 命令一起使用。
| `-p 5`<br>`--parallel-thread 10` |
+| `-s/--stop-after-disk-even` | 磁盘平衡完成后自动停止(默认值:true)。与 `start` 和 `update`
命令一起使用。
| `-s false`<br>`--stop-after-disk-even true` |
+| `-c/--container-states` | 以逗号分隔的容器生命周期状态名称,表示可在磁盘之间移动的状态。配合 `start` 和
`update` 命令使用。
| `-c
CLOSED,QUASI_CLOSED`<br>`--container-states CLOSED` |
### 示例
**启动 DiskBalancer:**
@@ -164,6 +165,9 @@ ozone admin datanode diskbalancer report
[<datanode-address> ...] [--in-service-
# 在多个数据节点上启动 DiskBalancer
ozone admin datanode diskbalancer start DN-1 DN-2 DN-3
+# 使用数据节点 UUID 启动 DiskBalancer
+ozone admin datanode diskbalancer start --node-id
a3b63511-bdf8-4fa1-8ab6-d19c0e806f84
+
# 在所有 IN_SERVICE 且 HEALTHY 的数据节点上启动 DiskBalancer
ozone admin datanode diskbalancer start --in-service-datanodes
diff --git
a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/AbstractDiskBalancerSubCommand.java
b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/AbstractDiskBalancerSubCommand.java
index 266795fcb08..38c3ba1d03b 100644
---
a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/AbstractDiskBalancerSubCommand.java
+++
b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/AbstractDiskBalancerSubCommand.java
@@ -26,6 +26,7 @@
import java.util.stream.Collectors;
import org.apache.hadoop.hdds.HddsConfigKeys;
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
+import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
import org.apache.hadoop.hdds.scm.cli.ContainerOperationClient;
import org.apache.hadoop.hdds.scm.client.ScmClient;
import org.apache.hadoop.hdds.server.JsonUtils;
@@ -43,11 +44,34 @@ public abstract class AbstractDiskBalancerSubCommand
implements Callable<Void> {
// Track if we're in batch mode to run commands on all in-service datanodes
private boolean isBatchMode = false;
- // Pre-fetched datanode address names for batch mode (address -> "hostname
(ip:port)"); null in non-batch
+ // Pre-fetched datanode address names (address -> display name); null when
not resolved via SCM
private Map<String, String> datanodeDisplayNames = null;
+ // Datanode identifiers that failed UUID-to-address resolution before RPC
execution
+ private final List<ResolutionFailure> resolutionFailures = new ArrayList<>();
+
+ // Normalized hostname/address args and --node-id UUIDs for the current
invocation
+ private List<String> explicitAddressArgs = new ArrayList<>();
+ private List<String> explicitNodeIds = new ArrayList<>();
+
+ private static final class ResolutionFailure {
+ private final String datanode;
+ private final String errorMsg;
+
+ private ResolutionFailure(String datanode, String errorMsg) {
+ this.datanode = datanode;
+ this.errorMsg = errorMsg;
+ }
+ }
+
@Override
public Void call() throws Exception {
+ resolutionFailures.clear();
+ datanodeDisplayNames = null;
+ explicitAddressArgs = new ArrayList<>();
+ explicitNodeIds = new ArrayList<>();
+ resetCommandState();
+
// Check if DiskBalancer is enabled in configuration
OzoneConfiguration conf = new OzoneConfiguration();
if
(!conf.getBoolean(HddsConfigKeys.HDDS_DATANODE_DISK_BALANCER_ENABLED_KEY,
@@ -57,10 +81,27 @@ public Void call() throws Exception {
return null;
}
- // Validate that either datanode addresses or --in-service-datanodes is
specified
- if ((options.getDatanodes() == null || options.getDatanodes().isEmpty())
+
explicitNodeIds.addAll(DiskBalancerSubCommandUtil.normalizeNodeIds(options.getNodeIds()));
+ for (String datanodeArg : options.getDatanodes()) {
+ if (DiskBalancerSubCommandUtil.isDatanodeUuid(datanodeArg)) {
+ if (explicitNodeIds.isEmpty()) {
+ System.err.println("Error: Datanode UUID must be specified with
--node-id, not as a "
+ + "positional argument. For multiple UUIDs use a comma-separated
list, for example "
+ + "--node-id uuid1,uuid2 or --node-id \"uuid1, uuid2\".");
+ return null;
+ }
+ explicitNodeIds.add(datanodeArg);
+ } else {
+ explicitAddressArgs.add(datanodeArg);
+ }
+ }
+
+ // Validate that either datanode addresses, --node-id, or
--in-service-datanodes is specified
+ if (explicitAddressArgs.isEmpty()
+ && explicitNodeIds.isEmpty()
&& !options.isInServiceDatanodes()) {
- System.err.println("Error: Either datanode address(es) or
--in-service-datanodes must be specified.");
+ System.err.println("Error: Either datanode address(es), --node-id, or
--in-service-datanodes "
+ + "must be specified.");
return null;
}
@@ -73,7 +114,10 @@ public Void call() throws Exception {
// Get the list of datanodes to execute on
List<String> targetDatanodes = getTargetDatanodes();
- if (targetDatanodes == null || targetDatanodes.isEmpty()) {
+ if (targetDatanodes == null) {
+ targetDatanodes = new ArrayList<>();
+ }
+ if (targetDatanodes.isEmpty() && resolutionFailures.isEmpty()) {
System.err.println("Error: No datanodes found to execute command on.");
return null;
}
@@ -90,6 +134,16 @@ public Void call() throws Exception {
List<String> successNodes = new ArrayList<>();
List<String> failedNodes = new ArrayList<>();
List<Object> jsonResults = new ArrayList<>();
+
+ for (ResolutionFailure resolutionFailure : resolutionFailures) {
+ failedNodes.add(resolutionFailure.datanode);
+ if (options.isJson()) {
+ jsonResults.add(createErrorResult(resolutionFailure.datanode,
resolutionFailure.errorMsg));
+ } else {
+ System.err.printf("Error on node [%s]: %s%n",
+ formatDatanodeDisplayName(resolutionFailure.datanode),
resolutionFailure.errorMsg);
+ }
+ }
// Execute commands and collect results
for (String dn : deduplicatedDatanodes) {
@@ -114,7 +168,7 @@ public Void call() throws Exception {
jsonResults.add(errorResult);
} else {
// Print error messages in non-JSON mode
- System.err.printf("Error on node [%s]: %s%n", dn, errorMsg);
+ System.err.printf("Error on node [%s]: %s%n",
formatDatanodeDisplayName(dn), errorMsg);
}
}
}
@@ -148,15 +202,70 @@ protected DiskBalancerCommonOptions getOptions() {
/**
* Get the list of target datanodes to execute the command on.
- * Either from positional arguments or by querying SCM for in-service
datanodes.
+ * Either from positional arguments, --node-id, or by querying SCM for
in-service datanodes.
*/
private List<String> getTargetDatanodes() {
if (options.isInServiceDatanodes()) {
return getAllInServiceDatanodes();
- } else {
- datanodeDisplayNames = null; // Non-batch: use user input as-is, no SCM
for formatting
- return options.getDatanodes();
}
+ return resolveExplicitDatanodeTargets(explicitAddressArgs,
explicitNodeIds);
+ }
+
+ /**
+ * Resolves hostname/host:port arguments and --node-id UUIDs to CLIENT_RPC
addresses.
+ * Hostname and host:port arguments are passed through unchanged without
contacting SCM.
+ */
+ private List<String> resolveExplicitDatanodeTargets(
+ List<String> addressArgs, List<String> nodeIdArgs) {
+ List<String> resolvedAddresses = new ArrayList<>(addressArgs);
+ if (nodeIdArgs.isEmpty()) {
+ datanodeDisplayNames = null;
+ return resolvedAddresses;
+ }
+
+ Map<String, String> displayNames = new LinkedHashMap<>();
+ ScmClient scmClient;
+ try {
+ scmClient = new ContainerOperationClient(new OzoneConfiguration());
+ } catch (IOException e) {
+ String msg = e.getMessage();
+ System.err.printf("Error resolving datanode address(es).%n%s%n", msg);
+ nodeIdArgs.forEach(nodeId -> addResolutionFailure(nodeId, msg));
+ datanodeDisplayNames = null;
+ return addressArgs;
+ }
+
+ try {
+ for (String nodeId : nodeIdArgs) {
+ try {
+ DiskBalancerSubCommandUtil.DatanodeTarget target =
+
DiskBalancerSubCommandUtil.resolveDatanodeTargetByUuid(scmClient, nodeId);
+ resolvedAddresses.add(target.getClientRpcAddress());
+ displayNames.put(target.getClientRpcAddress(),
target.getDisplayName());
+ } catch (IOException e) {
+ addResolutionFailure(nodeId, e.getMessage());
+ }
+ }
+ } finally {
+ try {
+ scmClient.close();
+ } catch (IOException e) {
+ System.err.printf("Error closing SCM client after resolving datanode
address(es).%n%s%n",
+ e.getMessage());
+ }
+ }
+
+ datanodeDisplayNames = displayNames.isEmpty() ? null : displayNames;
+ return resolvedAddresses;
+ }
+
+ private void addResolutionFailure(String datanode, String errorMsg) {
+ for (ResolutionFailure existing : resolutionFailures) {
+ if (existing.datanode.equals(datanode)) {
+ return;
+ }
+ }
+ resolutionFailures.add(new ResolutionFailure(datanode, errorMsg));
}
/**
@@ -172,6 +281,13 @@ private List<String> getAllInServiceDatanodes() {
}
}
+ /**
+ * Reset subcommand-specific state before each invocation.
+ */
+ protected void resetCommandState() {
+ // Default: no subcommand-specific state
+ }
+
/**
* Validate command parameters before execution.
*
@@ -248,12 +364,12 @@ private Map<String, Object> createErrorResult(String
datanode, String errorMsg)
}
/**
- * Format a datanode address for display.
- * In batch mode, uses pre-fetched display names from the SCM query.
- * In non-batch mode, returns the user's input as-is (no SCM call).
+ * Format a datanode address for display using pre-fetched SCM metadata when
available.
+ * Batch mode and --node-id populate {@link #datanodeDisplayNames}.
+ * Hostname and host:port arguments without SCM resolution are returned
as-is.
*
- * @param address the datanode address in "ip:port" format
- * @return formatted string "hostname (ip:port)" in batch mode, or address
as-is in non-batch
+ * @param address the datanode CLIENT_RPC address or unresolved user input
+ * @return UUID when resolved via --node-id, hostname (ip:port) in batch
mode, or address as-is
*/
protected String formatDatanodeDisplayName(String address) {
if (datanodeDisplayNames != null) {
@@ -261,5 +377,21 @@ protected String formatDatanodeDisplayName(String address)
{
}
return address;
}
+
+ /**
+ * Format a datanode for status/report output.
+ * Uses SCM-enriched display names when available; otherwise shows hostname
(ip:port) without UUID.
+ *
+ * @param address the datanode CLIENT_RPC address used for command execution
+ * @param nodeProto datanode details from the DiskBalancer RPC response
+ * @return formatted datanode identifier for output
+ */
+ protected String formatDatanodeDisplayName(
+ String address, HddsProtos.DatanodeDetailsProto nodeProto) {
+ if (datanodeDisplayNames != null &&
datanodeDisplayNames.containsKey(address)) {
+ return datanodeDisplayNames.get(address);
+ }
+ return DiskBalancerSubCommandUtil.getDatanodeHostAndIp(nodeProto);
+ }
}
diff --git
a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DatanodeParameters.java
b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DatanodeParameters.java
index 3406e340ae2..c239426cbe3 100644
---
a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DatanodeParameters.java
+++
b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DatanodeParameters.java
@@ -43,7 +43,8 @@ public class DatanodeParameters extends ItemsFromStdin {
" # From file having list of dns to balance",
" ozone admin datanode diskbalancer report - < datanode-lists.txt",
"Port is optional and defaults to 19864 (CLIENT_RPC port).",
- "Address examples: 'DN-1', 'DN-1:19864', '192.168.1.10'."
+ "Address examples: 'DN-1', 'DN-1:19864', '192.168.1.10'.",
+ "Use --node-id to target a datanode by UUID (requires SCM)."
},
arity = "0..*",
paramLabel = "<datanode address>")
diff --git
a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerCommands.java
b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerCommands.java
index 88c48244dda..7cf2e56befe 100644
---
a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerCommands.java
+++
b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerCommands.java
@@ -39,8 +39,15 @@
* in IN_SERVICE operational state, excluding
non-HEALTHY,
* DECOMMISSIONING, DECOMMISSIONED, and nodes
in maintenance states.
*
+ * Datanode identifiers:
+ * Positional arguments accept hostname, host:port, or IP address.
+ * Use --node-id to target a datanode by UUID (requires SCM).
+ * --in-service-datanodes queries SCM for all HEALTHY IN_SERVICE
datanodes.
+ *
* To start:
- * ozone admin datanode diskbalancer start {@literal <host[:port]>}
[{@literal <host[:port]>} ...]
+ * ozone admin datanode diskbalancer start {@literal <datanode-address>}
+ * [{@literal <datanode-address>} ...]
+ * [ --node-id {@literal <uuid>} ...]
* [ -t/--threshold-percentage {@literal <threshold>}]
* [ -b/--bandwidth-in-mb {@literal <bandwidthInMB>}]
* [ -p/--parallel-thread {@literal <parallelThread>}]
@@ -55,6 +62,9 @@
* ozone admin datanode diskbalancer start 192.168.1.10:19864
* Start balancer with explicit port specification
*
+ * ozone admin datanode diskbalancer start --node-id
a3b63511-bdf8-4fa1-8ab6-d19c0e806f84
+ * Start balancer using a datanode UUID (resolved via SCM)
+ *
* ozone admin datanode diskbalancer start DN-1 DN-2 DN-3
* Start balancer on multiple datanodes (using default port)
*
@@ -80,7 +90,9 @@
* Start balancer on all IN_SERVICE and HEALTHY datanodes and output
results in JSON format
*
* To stop:
- * ozone admin datanode diskbalancer stop {@literal <host[:port]>}
[{@literal <host[:port]>} ...]
+ * ozone admin datanode diskbalancer stop {@literal <datanode-address>}
+ * [{@literal <datanode-address>} ...]
+ * [ --node-id {@literal <uuid>} ...]
* [ --json ]
* [ --in-service-datanodes ]
*
@@ -98,7 +110,9 @@
* Stop diskbalancer on DN-1 and output result in JSON format
*
* To update:
- * ozone admin datanode diskbalancer update {@literal <host[:port]>}
[{@literal <host[:port]>} ...]
+ * ozone admin datanode diskbalancer update {@literal <datanode-address>}
+ * [{@literal <datanode-address>} ...]
+ * [ --node-id {@literal <uuid>} ...]
* [ -t/--threshold-percentage {@literal <threshold>}]
* [ -b/--bandwidth-in-mb {@literal <bandwidthInMB>}]
* [ -p/--parallel-thread {@literal <parallelThread>}]
@@ -117,7 +131,9 @@
* Update diskbalancer threshold to 10% on DN-1 and output result in
JSON format
*
* To get report:
- * ozone admin datanode diskbalancer report {@literal <host[:port]>}
[{@literal <host[:port]>} ...]
+ * ozone admin datanode diskbalancer report {@literal <datanode-address>}
+ * [{@literal <datanode-address>} ...]
+ * [ --node-id {@literal <uuid>} ...]
* [ --json ]
* [ --in-service-datanodes ]
*
@@ -135,7 +151,9 @@
* Retrieve volume density report from DN-1 in JSON format
*
* To get status:
- * ozone admin datanode diskbalancer status {@literal <host[:port]>}
[{@literal <host[:port]>} ...]
+ * ozone admin datanode diskbalancer status {@literal <datanode-address>}
+ * [{@literal <datanode-address>} ...]
+ * [ --node-id {@literal <uuid>} ...]
* [ --json ]
* [ --in-service-datanodes ]
*
@@ -143,6 +161,9 @@
* ozone admin datanode diskbalancer status DN-1
* Return the diskbalancer status on DN-1
*
+ * ozone admin datanode diskbalancer status --node-id
a3b63511-bdf8-4fa1-8ab6-d19c0e806f84
+ * Return the diskbalancer status using a datanode UUID
+ *
* ozone admin datanode diskbalancer status DN-1 DN-2 DN-3
* Return the diskbalancer status on multiple datanodes
*
diff --git
a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerCommonOptions.java
b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerCommonOptions.java
index 60da4a0fa6b..8ec25c4839d 100644
---
a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerCommonOptions.java
+++
b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerCommonOptions.java
@@ -35,6 +35,14 @@ public class DiskBalancerCommonOptions {
required = false)
private boolean inServiceDatanodes;
+ @CommandLine.Option(names = {"--node-id"},
+ description = "Datanode UUID(s). Requires SCM to resolve each UUID to a
CLIENT_RPC address. "
+ + "Pass a comma-separated list (for example, --node-id uuid1,uuid2).
"
+ + "When SCM is unavailable, use hostname or host:port positional
arguments instead.",
+ paramLabel = "<uuid>",
+ split = ",\\s*")
+ private List<String> nodeIds;
+
@CommandLine.Option(names = {"--json"},
description = "Format output as JSON",
defaultValue = "false")
@@ -50,6 +58,10 @@ public boolean isInServiceDatanodes() {
return inServiceDatanodes;
}
+ public List<String> getNodeIds() {
+ return nodeIds != null ? nodeIds : Collections.emptyList();
+ }
+
public boolean isJson() {
return json;
}
diff --git
a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerReportSubcommand.java
b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerReportSubcommand.java
index 16a91d45e28..39e806310e2 100644
---
a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerReportSubcommand.java
+++
b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerReportSubcommand.java
@@ -20,6 +20,7 @@
import static java.util.stream.Collectors.toList;
import java.io.IOException;
+import java.util.AbstractMap;
import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
@@ -49,6 +50,11 @@ public class DiskBalancerReportSubcommand extends
AbstractDiskBalancerSubCommand
private static final String PERCENT_FORMAT = "%.2f%%";
+ @Override
+ protected void resetCommandState() {
+ reports.clear();
+ }
+
@Override
protected Object executeCommand(String hostName) throws IOException {
DiskBalancerProtocol diskBalancerProxy = DiskBalancerSubCommandUtil
@@ -58,7 +64,7 @@ protected Object executeCommand(String hostName) throws
IOException {
// Only create JSON result object if JSON mode is enabled
if (getOptions().isJson()) {
- return toJson(report);
+ return toJson(hostName, report);
}
// For non-JSON mode, store the proto for later consolidation
@@ -89,20 +95,27 @@ protected void displayResults(List<String> successNodes,
List<String> failedNode
List<DatanodeDiskBalancerInfoProto> reportList = successNodes.stream()
.map(reports::get)
.collect(toList());
- System.out.println(generateReport(reportList));
+ System.out.println(generateReport(successNodes, reportList));
}
}
- private String generateReport(List<DatanodeDiskBalancerInfoProto> protos) {
- protos.sort((a, b) ->
- Double.compare(b.getCurrentVolumeDensitySum(),
a.getCurrentVolumeDensitySum()));
+ private String generateReport(
+ List<String> successNodes, List<DatanodeDiskBalancerInfoProto> protos) {
+ List<Map.Entry<String, DatanodeDiskBalancerInfoProto>> entries = new
ArrayList<>();
+ for (int i = 0; i < protos.size(); i++) {
+ entries.add(new AbstractMap.SimpleEntry<>(successNodes.get(i),
protos.get(i)));
+ }
+ entries.sort((a, b) -> Double.compare(
+ b.getValue().getCurrentVolumeDensitySum(),
+ a.getValue().getCurrentVolumeDensitySum()));
StringBuilder formatBuilder = new StringBuilder("Report result:%n");
List<String> contentList = new ArrayList<>();
- for (int i = 0; i < protos.size(); i++) {
- DatanodeDiskBalancerInfoProto p = protos.get(i);
- String dn = DiskBalancerSubCommandUtil.getDatanodeHostAndIp(p.getNode());
+ for (int i = 0; i < entries.size(); i++) {
+ Map.Entry<String, DatanodeDiskBalancerInfoProto> entry = entries.get(i);
+ DatanodeDiskBalancerInfoProto p = entry.getValue();
+ String dn = formatDatanodeDisplayName(entry.getKey(), p.getNode());
StringBuilder header = new StringBuilder();
header.append("Datanode: ").append(dn).append(System.lineSeparator())
@@ -156,7 +169,7 @@ private String
generateReport(List<DatanodeDiskBalancerInfoProto> protos) {
formatBuilder.append("%n");
}
- if (i < protos.size() - 1) {
+ if (i < entries.size() - 1) {
formatBuilder.append("-------%n%n");
}
}
@@ -198,9 +211,9 @@ private static String formatPercent(double ratio) {
* @param report the DiskBalancer report proto
* @return JSON result map
*/
- private Map<String, Object> toJson(DatanodeDiskBalancerInfoProto report) {
+ private Map<String, Object> toJson(String hostName,
DatanodeDiskBalancerInfoProto report) {
Map<String, Object> result = new LinkedHashMap<>();
- result.put("datanode",
DiskBalancerSubCommandUtil.getDatanodeHostAndIp(report.getNode()));
+ result.put("datanode", formatDatanodeDisplayName(hostName,
report.getNode()));
result.put("action", "report");
result.put("status", "success");
result.put("volumeDensity",
formatPercent(report.getCurrentVolumeDensitySum()));
diff --git
a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerStatusSubcommand.java
b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerStatusSubcommand.java
index 291ff2d2d3c..1c133f5932d 100644
---
a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerStatusSubcommand.java
+++
b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerStatusSubcommand.java
@@ -24,7 +24,6 @@
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
-import java.util.Objects;
import org.apache.hadoop.hdds.cli.HddsVersionProvider;
import org.apache.hadoop.hdds.protocol.DiskBalancerProtocol;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
@@ -45,6 +44,11 @@ public class DiskBalancerStatusSubcommand extends
AbstractDiskBalancerSubCommand
private final Map<String, DatanodeDiskBalancerInfoProto> statuses =
new LinkedHashMap<>();
+ @Override
+ protected void resetCommandState() {
+ statuses.clear();
+ }
+
@Override
protected Object executeCommand(String hostName) throws IOException {
DiskBalancerProtocol diskBalancerProxy = DiskBalancerSubCommandUtil
@@ -54,7 +58,7 @@ protected Object executeCommand(String hostName) throws
IOException {
// Only create JSON result object if JSON mode is enabled
if (getOptions().isJson()) {
- return createStatusResult(status);
+ return createStatusResult(hostName, status);
}
// For non-JSON mode, store the proto for later consolidation
@@ -82,18 +86,23 @@ protected void displayResults(List<String> successNodes,
List<String> failedNode
// Display consolidated status for successful nodes
if (!successNodes.isEmpty() && !statuses.isEmpty()) {
- List<DatanodeDiskBalancerInfoProto> statusList =
- successNodes.stream()
- .map(statuses::get)
- .filter(Objects::nonNull)
- .collect(toList());
- System.out.println(generateStatus(statusList));
+ List<DatanodeDiskBalancerInfoProto> statusList = new ArrayList<>();
+ List<String> displayNames = new ArrayList<>();
+ for (String successNode : successNodes) {
+ DatanodeDiskBalancerInfoProto proto = statuses.get(successNode);
+ if (proto != null) {
+ statusList.add(proto);
+ displayNames.add(formatDatanodeDisplayName(successNode,
proto.getNode()));
+ }
+ }
+ System.out.println(generateStatus(statusList, displayNames));
}
}
- private String generateStatus(List<DatanodeDiskBalancerInfoProto> protos) {
+ private String generateStatus(
+ List<DatanodeDiskBalancerInfoProto> protos, List<String>
datanodeDisplayNames) {
StringBuilder formatBuilder = new StringBuilder("Status result:%n" +
- "%-60s %-12s %-15s %-15s %-12s %-20s %-40s %-12s %-12s %-15s %-18s
%-20s%n");
+ "%-60s %-10s %-15s %-15s %-10s %-18s %-30s %-12s %-12s %-15s %-18s
%-20s%n");
List<String> contentList = new ArrayList<>();
contentList.add("Datanode");
@@ -109,16 +118,14 @@ private String
generateStatus(List<DatanodeDiskBalancerInfoProto> protos) {
contentList.add("EstBytesToMove(MB)");
contentList.add("EstTimeLeft(min)");
- for (HddsProtos.DatanodeDiskBalancerInfoProto proto : protos) {
- formatBuilder.append("%-60s %-12s %-15s %-15s %-12s %-20s %-40s %-12s
%-12s %-15s %-18s %-20s%n");
+ for (int i = 0; i < protos.size(); i++) {
+ HddsProtos.DatanodeDiskBalancerInfoProto proto = protos.get(i);
+ formatBuilder.append("%-60s %-10s %-15s %-15s %-10s %-18s %-30s %-12s
%-12s %-15s %-18s %-20s%n");
long estimatedTimeLeft = calculateEstimatedTimeLeft(proto);
long bytesMovedMB = (long) Math.ceil(proto.getBytesMoved() / (1024.0 *
1024.0));
long bytesToMoveMB = (long) Math.ceil(proto.getBytesToMove() / (1024.0 *
1024.0));
- // Format datanode string with hostname and IP address
- String formattedDatanode =
DiskBalancerSubCommandUtil.getDatanodeHostAndIp(
- proto.getNode());
- contentList.add(formattedDatanode);
+ contentList.add(datanodeDisplayNames.get(i));
contentList.add(proto.getRunningStatus().name());
contentList.add(
String.format("%.4f", proto.getDiskBalancerConf().getThreshold()));
@@ -161,11 +168,10 @@ protected String getActionName() {
* @param status the DiskBalancer status proto
* @return JSON result map
*/
- private Map<String, Object> createStatusResult(DatanodeDiskBalancerInfoProto
status) {
+ private Map<String, Object> createStatusResult(
+ String hostName, DatanodeDiskBalancerInfoProto status) {
Map<String, Object> result = new LinkedHashMap<>();
- // Format datanode string with hostname and IP address
- String formattedDatanode = DiskBalancerSubCommandUtil.getDatanodeHostAndIp(
- status.getNode());
+ String formattedDatanode = formatDatanodeDisplayName(hostName,
status.getNode());
result.put("datanode", formattedDatanode);
result.put("action", "status");
result.put("status", "success");
diff --git
a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerSubCommandUtil.java
b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerSubCommandUtil.java
index 09d8cf5e02c..29c47dabff4 100644
---
a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerSubCommandUtil.java
+++
b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/datanode/DiskBalancerSubCommandUtil.java
@@ -22,9 +22,12 @@
import java.io.IOException;
import java.net.InetSocketAddress;
+import java.util.ArrayList;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
+import java.util.UUID;
+import java.util.regex.Pattern;
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
import org.apache.hadoop.hdds.protocol.DatanodeDetails.Port;
@@ -40,9 +43,111 @@
*/
final class DiskBalancerSubCommandUtil {
+ private static final Pattern DATANODE_UUID_PATTERN = Pattern.compile(
+
"^[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}$");
+
+ static final class DatanodeTarget {
+ private final String clientRpcAddress;
+ private final String displayName;
+
+ DatanodeTarget(String clientRpcAddress, String displayName) {
+ this.clientRpcAddress = clientRpcAddress;
+ this.displayName = displayName;
+ }
+
+ String getClientRpcAddress() {
+ return clientRpcAddress;
+ }
+
+ String getDisplayName() {
+ return displayName;
+ }
+ }
+
private DiskBalancerSubCommandUtil() {
}
+ /**
+ * Returns true if the argument is a canonical datanode UUID rather than a
host or address.
+ */
+ static boolean isDatanodeUuid(String nodeArg) {
+ if (!DATANODE_UUID_PATTERN.matcher(nodeArg).matches()) {
+ return false;
+ }
+ try {
+ UUID.fromString(nodeArg);
+ return true;
+ } catch (IllegalArgumentException e) {
+ return false;
+ }
+ }
+
+ /**
+ * Normalizes {@code --node-id} values, including comma-separated lists and
trailing commas when
+ * the shell splits {@code uuid1, uuid2} into separate arguments.
+ */
+ static List<String> normalizeNodeIds(List<String> rawNodeIds) {
+ List<String> normalized = new ArrayList<>();
+ if (rawNodeIds == null) {
+ return normalized;
+ }
+ for (String rawNodeId : rawNodeIds) {
+ if (rawNodeId == null || rawNodeId.isEmpty()) {
+ continue;
+ }
+ for (String nodeId : rawNodeId.split(",\\s*")) {
+ String trimmed = nodeId.trim();
+ if (!trimmed.isEmpty()) {
+ normalized.add(trimmed);
+ }
+ }
+ }
+ return normalized;
+ }
+
+ /**
+ * Resolves a datanode hostname or host:port to a CLIENT_RPC address without
contacting SCM.
+ */
+ static DatanodeTarget resolveDatanodeAddress(String nodeArg) {
+ return new DatanodeTarget(nodeArg, nodeArg);
+ }
+
+ /**
+ * Resolves a datanode UUID to a CLIENT_RPC address via SCM.
+ */
+ static DatanodeTarget resolveDatanodeTargetByUuid(ScmClient scmClient,
String nodeUuid)
+ throws IOException {
+ if (!isDatanodeUuid(nodeUuid)) {
+ throw new IOException("Invalid datanode UUID: " + nodeUuid);
+ }
+
+ HddsProtos.Node node = scmClient.queryNode(UUID.fromString(nodeUuid));
+ HddsProtos.DatanodeDetailsProto nodeId = node.getNodeID();
+ if (!node.hasNodeID() || (!nodeId.hasUuid() && !nodeId.hasUuid128() &&
!nodeId.hasId())) {
+ throw new IOException("Datanode not found: " + nodeUuid);
+ }
+
+ DatanodeDetails details = DatanodeDetails.getFromProtoBuf(nodeId);
+ if (details.getIpAddress() == null || details.getIpAddress().isEmpty()) {
+ throw new IOException("Datanode not found: " + nodeUuid);
+ }
+
+ String address = getClientRpcAddress(details);
+ return new DatanodeTarget(address, nodeUuid);
+ }
+
+ /**
+ * Resolves a datanode identifier to a CLIENT_RPC address.
+ * Accepts datanode UUID, hostname, or host:port.
+ */
+ static DatanodeTarget resolveDatanodeTarget(ScmClient scmClient, String
nodeArg)
+ throws IOException {
+ if (!isDatanodeUuid(nodeArg)) {
+ return resolveDatanodeAddress(nodeArg);
+ }
+ return resolveDatanodeTargetByUuid(scmClient, nodeArg);
+ }
+
/**
* Creates a DiskBalancerProtocol proxy for a single datanode.
*
@@ -95,31 +200,32 @@ public static Map<String, String>
getAllOperableNodesClientRpcAddress(
if (node.getNodeStates(0).equals(HddsProtos.NodeState.DEAD)) {
continue;
}
- Port port = details.getPort(Port.Name.CLIENT_RPC);
- if (port != null) {
- String address = details.getIpAddress() + ":" + port.getValue();
- // Format the display string: "hostname (ip:port)" or "ip:port"
- String hostname = details.getHostName();
- String display = (hostname != null && !hostname.isEmpty()
- && !hostname.equals(details.getIpAddress())) ? hostname + " (" +
address + ")"
- : address;
- addressToDisplay.put(address, display);
- } else {
- System.out.printf("host: %s(%s) %s port not found%n",
- details.getHostName(), details.getIpAddress(),
- Port.Name.CLIENT_RPC.name());
+ try {
+ String address = getClientRpcAddress(details);
+ addressToDisplay.put(address, getDatanodeHostAndIp(node.getNodeID()));
+ } catch (IOException e) {
+ System.err.println(e.getMessage());
}
}
return addressToDisplay;
}
+ static String getClientRpcAddress(DatanodeDetails details) throws
IOException {
+ Port port = details.getPort(Port.Name.CLIENT_RPC);
+ if (port == null) {
+ throw new IOException(String.format("host: %s(%s) %s port not found",
+ details.getHostName(), details.getIpAddress(),
Port.Name.CLIENT_RPC.name()));
+ }
+ return details.getIpAddress() + ":" + port.getValue();
+ }
+
/**
* Returns a formatted string combining hostname and IP address from
DatanodeDetailsProto.
- * If hostname is null or empty, returns just "ip:port".
- *
+ * Format: {@code hostname (ip:port)} or {@code ip:port}.
+ *
* @param nodeProto the DatanodeDetailsProto from the diskbalancer info
- * @return formatted string "hostname (ip:port)" or "ip:port" if hostname is
not available
+ * @return formatted datanode identifier for status/report output
*/
public static String getDatanodeHostAndIp(HddsProtos.DatanodeDetailsProto
nodeProto) {
String hostname = nodeProto.getHostName();
@@ -131,7 +237,6 @@ public static String
getDatanodeHostAndIp(HddsProtos.DatanodeDetailsProto nodePr
.findFirst()
.orElse(HDDS_DATANODE_CLIENT_PORT_DEFAULT); // Default port if not
found
- // Format the output string
String addressPort = ipAddress + ":" + port;
if (hostname != null && !hostname.isEmpty() &&
!hostname.equals(ipAddress)) {
return hostname + " (" + addressPort + ")";
diff --git
a/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/datanode/TestDiskBalancerSubCommandUtil.java
b/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/datanode/TestDiskBalancerSubCommandUtil.java
new file mode 100644
index 00000000000..3e07752fd55
--- /dev/null
+++
b/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/datanode/TestDiskBalancerSubCommandUtil.java
@@ -0,0 +1,135 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hadoop.hdds.scm.cli.datanode;
+
+import static
org.apache.hadoop.hdds.HddsConfigKeys.HDDS_DATANODE_CLIENT_PORT_DEFAULT;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+import java.io.IOException;
+import java.util.Arrays;
+import java.util.List;
+import java.util.UUID;
+import org.apache.hadoop.hdds.protocol.DatanodeDetails;
+import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
+import org.apache.hadoop.hdds.scm.client.ScmClient;
+import org.junit.jupiter.api.Test;
+
+/**
+ * Unit tests for {@link DiskBalancerSubCommandUtil}.
+ */
+public class TestDiskBalancerSubCommandUtil {
+
+ private static final String DN_UUID = "a3b63511-bdf8-4fa1-8ab6-d19c0e806f84";
+
+ @Test
+ public void testIsDatanodeUuid() {
+ assertTrue(DiskBalancerSubCommandUtil.isDatanodeUuid(DN_UUID));
+ assertFalse(DiskBalancerSubCommandUtil.isDatanodeUuid("host-1"));
+ assertFalse(DiskBalancerSubCommandUtil.isDatanodeUuid("host-1:19864"));
+ assertFalse(DiskBalancerSubCommandUtil.isDatanodeUuid("10.140.95.199"));
+ assertFalse(DiskBalancerSubCommandUtil.isDatanodeUuid(
+ "a3b63511bdf84fa18ab6d19c0e806f84"));
+ }
+
+ @Test
+ public void testResolveDatanodeTargetWithHostname() throws IOException {
+ ScmClient scmClient = mock(ScmClient.class);
+
+ DiskBalancerSubCommandUtil.DatanodeTarget target =
+ DiskBalancerSubCommandUtil.resolveDatanodeTarget(scmClient, "host-1");
+
+ assertEquals("host-1", target.getClientRpcAddress());
+ assertEquals("host-1", target.getDisplayName());
+ }
+
+ @Test
+ public void testResolveDatanodeTargetWithUuid() throws IOException {
+ ScmClient scmClient = mock(ScmClient.class);
+ HddsProtos.Node node = buildNode(DN_UUID, "nodename", "10.140.95.199",
HDDS_DATANODE_CLIENT_PORT_DEFAULT);
+ when(scmClient.queryNode(UUID.fromString(DN_UUID))).thenReturn(node);
+
+ DiskBalancerSubCommandUtil.DatanodeTarget target =
+ DiskBalancerSubCommandUtil.resolveDatanodeTarget(scmClient, DN_UUID);
+
+ assertEquals("10.140.95.199:" + HDDS_DATANODE_CLIENT_PORT_DEFAULT,
+ target.getClientRpcAddress());
+ assertEquals(DN_UUID, target.getDisplayName());
+ }
+
+ @Test
+ public void testResolveDatanodeTargetWithUnknownUuid() throws IOException {
+ ScmClient scmClient = mock(ScmClient.class);
+ when(scmClient.queryNode(UUID.fromString(DN_UUID)))
+ .thenReturn(HddsProtos.Node.getDefaultInstance());
+
+ IOException ex = assertThrows(IOException.class,
+ () -> DiskBalancerSubCommandUtil.resolveDatanodeTarget(scmClient,
DN_UUID));
+ assertTrue(ex.getMessage().contains("Datanode not found"));
+ }
+
+ @Test
+ public void testGetClientRpcAddress() throws IOException {
+ DatanodeDetails details = DatanodeDetails.getFromProtoBuf(
+ buildNode(DN_UUID, "nodename", "10.140.95.199",
HDDS_DATANODE_CLIENT_PORT_DEFAULT)
+ .getNodeID());
+
+ assertEquals("10.140.95.199:" + HDDS_DATANODE_CLIENT_PORT_DEFAULT,
+ DiskBalancerSubCommandUtil.getClientRpcAddress(details));
+ }
+
+ @Test
+ public void testGetDatanodeHostAndIp() {
+ HddsProtos.DatanodeDetailsProto nodeProto = buildNode(
+ "6d8157c2-280d-4eb2-a264-d388b05a0a87",
+ "ozone-datanode-2.ozone_default",
+ "172.18.0.6",
+ HDDS_DATANODE_CLIENT_PORT_DEFAULT).getNodeID();
+
+ assertEquals(
+ "ozone-datanode-2.ozone_default (172.18.0.6:" +
HDDS_DATANODE_CLIENT_PORT_DEFAULT + ")",
+ DiskBalancerSubCommandUtil.getDatanodeHostAndIp(nodeProto));
+ }
+
+ @Test
+ public void testNormalizeNodeIds() {
+ List<String> normalized = DiskBalancerSubCommandUtil.normalizeNodeIds(
+ Arrays.asList("uuid1,", " uuid2"));
+ assertEquals(2, normalized.size());
+ assertEquals("uuid1", normalized.get(0));
+ assertEquals("uuid2", normalized.get(1));
+ }
+
+ private static HddsProtos.Node buildNode(
+ String uuid, String hostname, String ipAddress, int clientRpcPort) {
+ HddsProtos.DatanodeDetailsProto dnd =
HddsProtos.DatanodeDetailsProto.newBuilder()
+ .setUuid(uuid)
+ .setHostName(hostname)
+ .setIpAddress(ipAddress)
+ .addPorts(HddsProtos.Port.newBuilder()
+ .setName(DatanodeDetails.Port.Name.CLIENT_RPC.name())
+ .setValue(clientRpcPort)
+ .build())
+ .build();
+ return HddsProtos.Node.newBuilder().setNodeID(dnd).build();
+ }
+}
diff --git
a/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/datanode/TestDiskBalancerSubCommands.java
b/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/datanode/TestDiskBalancerSubCommands.java
index 0e5667e9ea0..1b5ef945ed1 100644
---
a/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/datanode/TestDiskBalancerSubCommands.java
+++
b/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/hdds/scm/cli/datanode/TestDiskBalancerSubCommands.java
@@ -19,6 +19,7 @@
import static
org.apache.hadoop.hdds.HddsConfigKeys.HDDS_DATANODE_CLIENT_PORT_DEFAULT;
import static org.assertj.core.api.Assertions.assertThat;
+import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyString;
@@ -28,6 +29,7 @@
import static org.mockito.Mockito.mockConstruction;
import static org.mockito.Mockito.mockStatic;
import static org.mockito.Mockito.when;
+import static org.mockito.Mockito.withSettings;
import java.io.ByteArrayInputStream;
import java.io.ByteArrayOutputStream;
@@ -40,6 +42,7 @@
import java.util.List;
import java.util.Map;
import java.util.Random;
+import java.util.UUID;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
import java.util.stream.Stream;
@@ -60,6 +63,7 @@
import org.junit.jupiter.params.provider.MethodSource;
import org.mockito.MockedConstruction;
import org.mockito.MockedStatic;
+import org.mockito.Mockito;
import picocli.CommandLine;
/**
@@ -102,6 +106,10 @@ private static class DiskBalancerMocks implements
AutoCloseable {
this.mockedClient = mockedClient;
this.mockedUtil = mockedUtil;
}
+
+ MockedStatic<DiskBalancerSubCommandUtil> getMockedUtil() {
+ return mockedUtil;
+ }
@Override
public void close() {
@@ -123,7 +131,8 @@ private DiskBalancerMocks setupAllMocks() {
mockConstruction(ContainerOperationClient.class);
MockedStatic<DiskBalancerSubCommandUtil> mockedUtil =
- mockStatic(DiskBalancerSubCommandUtil.class);
+ mockStatic(DiskBalancerSubCommandUtil.class,
withSettings().defaultAnswer(
+ Mockito.CALLS_REAL_METHODS));
Map<String, String> addressToDisplay = new LinkedHashMap<>();
for (String addr : inServiceDatanodes) {
addressToDisplay.put(addr, addr);
@@ -134,38 +143,6 @@ private DiskBalancerMocks setupAllMocks() {
mockedUtil.when(() -> DiskBalancerSubCommandUtil
.getSingleNodeDiskBalancerProxy(anyString()))
.thenReturn(mockProtocol);
- // Mock getDatanodeHostAndIp(HddsProtos.DatanodeDetailsProto) to format
the output
- mockedUtil.when(() -> DiskBalancerSubCommandUtil
- .getDatanodeHostAndIp(any(HddsProtos.DatanodeDetailsProto.class)))
- .thenAnswer(invocation -> {
- HddsProtos.DatanodeDetailsProto proto = invocation.getArgument(0);
- return proto.getHostName() + " (" + proto.getIpAddress() + ":" +
- HDDS_DATANODE_CLIENT_PORT_DEFAULT + ")";
- });
- // Mock getDatanodeHostAndIp(String, String, int) to format the output
- // Return value is used by Mockito internally for mock setup
- mockedUtil.when(() -> {
- @SuppressWarnings("RV_RETURN_VALUE_IGNORED_NO_SIDE_EFFECT")
- String ignored = DiskBalancerSubCommandUtil
- .getDatanodeHostAndIp(any(DatanodeDetailsProto.class));
- // Use the value to avoid "ignored return value" static analysis
warnings.
- System.out.println(ignored);
- }).thenAnswer(invocation -> {
- DatanodeDetailsProto proto = invocation.getArgument(0);
- String hostname = proto.getHostName();
- String ipAddress = proto.getIpAddress();
- int port = proto.getPortsList().stream()
- .filter(p -> p.getName().equals(
- DatanodeDetails.Port.Name.CLIENT_RPC.name()))
- .mapToInt(HddsProtos.Port::getValue)
- .findFirst()
- .orElse(HDDS_DATANODE_CLIENT_PORT_DEFAULT);
- String addressPort = ipAddress + ":" + port;
- if (hostname != null && !hostname.isEmpty() &&
!hostname.equals(ipAddress)) {
- return hostname + " (" + addressPort + ")";
- }
- return addressPort;
- });
return new DiskBalancerMocks(mockedClient, mockedUtil);
}
@@ -533,6 +510,293 @@ public void testStatusDiskBalancerWithMultipleNodes()
throws Exception {
}
}
+ @Test
+ public void testStatusDiskBalancerWithDatanodeUuid() throws Exception {
+ final String dnUuid = "a3b63511-bdf8-4fa1-8ab6-d19c0e806f84";
+ final String resolvedAddress = "10.140.95.199:" +
HDDS_DATANODE_CLIENT_PORT_DEFAULT;
+
+ HddsProtos.DatanodeDetailsProto dnd =
HddsProtos.DatanodeDetailsProto.newBuilder()
+ .setUuid(dnUuid)
+ .setHostName("nodename")
+ .setIpAddress("10.140.95.199")
+ .addPorts(HddsProtos.Port.newBuilder()
+ .setName(DatanodeDetails.Port.Name.CLIENT_RPC.name())
+ .setValue(HDDS_DATANODE_CLIENT_PORT_DEFAULT)
+ .build())
+ .build();
+ HddsProtos.Node node = HddsProtos.Node.newBuilder().setNodeID(dnd).build();
+
+ DiskBalancerStatusSubcommand cmd = new DiskBalancerStatusSubcommand();
+ DatanodeDiskBalancerInfoProto statusProto =
generateRandomStatusProto("nodename").toBuilder()
+ .setNode(dnd)
+ .build();
+ when(mockProtocol.getDiskBalancerInfo()).thenReturn(statusProto);
+
+ try (MockedConstruction<ContainerOperationClient> mockedClient =
+ mockConstruction(ContainerOperationClient.class, (mock, context) ->
+ when(mock.queryNode(UUID.fromString(dnUuid))).thenReturn(node));
+ MockedStatic<DiskBalancerSubCommandUtil> mockedUtil =
+ mockStatic(DiskBalancerSubCommandUtil.class,
withSettings().defaultAnswer(
+ Mockito.CALLS_REAL_METHODS))) {
+
+ mockedUtil.when(() -> DiskBalancerSubCommandUtil
+ .getSingleNodeDiskBalancerProxy(resolvedAddress))
+ .thenReturn(mockProtocol);
+
+ CommandLine c = new CommandLine(cmd);
+ c.parseArgs("--node-id", dnUuid);
+ cmd.call();
+
+ String output = outContent.toString(DEFAULT_ENCODING);
+ assertTrue(output.contains("Status result"));
+ assertTrue(output.contains(dnUuid));
+ mockedUtil.verify(() -> DiskBalancerSubCommandUtil
+ .getSingleNodeDiskBalancerProxy(resolvedAddress));
+ }
+ }
+
+ @Test
+ public void testStatusDiskBalancerWithMixedValidAndInvalidUuids() throws
Exception {
+ final String validUuid = "a3b63511-bdf8-4fa1-8ab6-d19c0e806f84";
+ final String invalidUuid = "00000000-0000-0000-0000-000000000000";
+ final String resolvedAddress = "10.140.95.199:" +
HDDS_DATANODE_CLIENT_PORT_DEFAULT;
+
+ HddsProtos.DatanodeDetailsProto dnd =
HddsProtos.DatanodeDetailsProto.newBuilder()
+ .setUuid(validUuid)
+ .setHostName("nodename")
+ .setIpAddress("10.140.95.199")
+ .addPorts(HddsProtos.Port.newBuilder()
+ .setName(DatanodeDetails.Port.Name.CLIENT_RPC.name())
+ .setValue(HDDS_DATANODE_CLIENT_PORT_DEFAULT)
+ .build())
+ .build();
+ HddsProtos.Node node = HddsProtos.Node.newBuilder().setNodeID(dnd).build();
+
+ DiskBalancerStatusSubcommand cmd = new DiskBalancerStatusSubcommand();
+ DatanodeDiskBalancerInfoProto statusProto =
generateRandomStatusProto("nodename").toBuilder()
+ .setNode(dnd)
+ .build();
+ when(mockProtocol.getDiskBalancerInfo())
+ .thenReturn(generateRandomStatusProto("host-1"), statusProto);
+
+ try (MockedConstruction<ContainerOperationClient> mockedClient =
+ mockConstruction(ContainerOperationClient.class, (mock, context) -> {
+ when(mock.queryNode(UUID.fromString(validUuid))).thenReturn(node);
+ when(mock.queryNode(UUID.fromString(invalidUuid)))
+ .thenReturn(HddsProtos.Node.getDefaultInstance());
+ });
+ MockedStatic<DiskBalancerSubCommandUtil> mockedUtil =
+ mockStatic(DiskBalancerSubCommandUtil.class,
withSettings().defaultAnswer(
+ Mockito.CALLS_REAL_METHODS))) {
+
+ mockedUtil.when(() -> DiskBalancerSubCommandUtil
+ .getSingleNodeDiskBalancerProxy(resolvedAddress))
+ .thenReturn(mockProtocol);
+ mockedUtil.when(() -> DiskBalancerSubCommandUtil
+ .getSingleNodeDiskBalancerProxy("host-1"))
+ .thenReturn(mockProtocol);
+
+ CommandLine c = new CommandLine(cmd);
+ c.parseArgs("--node-id", validUuid + "," + invalidUuid, "host-1");
+ cmd.call();
+
+ String output = outContent.toString(DEFAULT_ENCODING);
+ String err = errContent.toString(DEFAULT_ENCODING);
+ assertTrue(output.contains("Status result"));
+ assertTrue(output.contains(validUuid));
+ assertTrue(output.contains("host-1"));
+ assertTrue(err.contains(invalidUuid));
+ assertTrue(err.contains("Datanode not found"));
+ }
+ }
+
+ @Test
+ public void testStartDiskBalancerWithDatanodeUuidJson() throws Exception {
+ final String dnUuid = "a3b63511-bdf8-4fa1-8ab6-d19c0e806f84";
+ final String resolvedAddress = "10.140.95.199:" +
HDDS_DATANODE_CLIENT_PORT_DEFAULT;
+
+ HddsProtos.DatanodeDetailsProto dnd =
HddsProtos.DatanodeDetailsProto.newBuilder()
+ .setUuid(dnUuid)
+ .setHostName("nodename")
+ .setIpAddress("10.140.95.199")
+ .addPorts(HddsProtos.Port.newBuilder()
+ .setName(DatanodeDetails.Port.Name.CLIENT_RPC.name())
+ .setValue(HDDS_DATANODE_CLIENT_PORT_DEFAULT)
+ .build())
+ .build();
+ HddsProtos.Node node = HddsProtos.Node.newBuilder().setNodeID(dnd).build();
+
+ DiskBalancerStartSubcommand cmd = new DiskBalancerStartSubcommand();
+
doNothing().when(mockProtocol).startDiskBalancer(any(DiskBalancerConfigurationProto.class));
+
+ try (MockedConstruction<ContainerOperationClient> mockedClient =
+ mockConstruction(ContainerOperationClient.class, (mock, context) ->
+ when(mock.queryNode(UUID.fromString(dnUuid))).thenReturn(node));
+ MockedStatic<DiskBalancerSubCommandUtil> mockedUtil =
+ mockStatic(DiskBalancerSubCommandUtil.class,
withSettings().defaultAnswer(
+ Mockito.CALLS_REAL_METHODS))) {
+
+ mockedUtil.when(() -> DiskBalancerSubCommandUtil
+ .getSingleNodeDiskBalancerProxy(resolvedAddress))
+ .thenReturn(mockProtocol);
+
+ CommandLine c = new CommandLine(cmd);
+ c.parseArgs("--json", "-t", "0.005", "-b", "100", "--node-id", dnUuid);
+ cmd.call();
+
+ String output = outContent.toString(DEFAULT_ENCODING);
+ assertTrue(output.contains("\"datanode\" : \"" + dnUuid + "\""));
+ }
+ }
+
+ @Test
+ public void testStatusDiskBalancerWithSpaceAfterCommaNodeIds() throws
Exception {
+ final String uuid1 = "59c14bfa-1ccd-45e4-83e6-8c2c3a5de873";
+ final String uuid2 = "0d4a065f-db6c-4649-9906-1a4df09ffbdf";
+ final String resolvedAddress1 = "10.140.95.199:" +
HDDS_DATANODE_CLIENT_PORT_DEFAULT;
+ final String resolvedAddress2 = "10.140.95.200:" +
HDDS_DATANODE_CLIENT_PORT_DEFAULT;
+
+ HddsProtos.Node node1 = buildScmNode(uuid1, "nodename-1", "10.140.95.199");
+ HddsProtos.Node node2 = buildScmNode(uuid2, "nodename-2", "10.140.95.200");
+
+ DiskBalancerStatusSubcommand cmd = new DiskBalancerStatusSubcommand();
+ when(mockProtocol.getDiskBalancerInfo())
+ .thenReturn(generateRandomStatusProto("nodename-1"),
generateRandomStatusProto("nodename-2"));
+
+ try (MockedConstruction<ContainerOperationClient> mockedClient =
+ mockConstruction(ContainerOperationClient.class, (mock, context) -> {
+ when(mock.queryNode(UUID.fromString(uuid1))).thenReturn(node1);
+ when(mock.queryNode(UUID.fromString(uuid2))).thenReturn(node2);
+ });
+ MockedStatic<DiskBalancerSubCommandUtil> mockedUtil =
+ mockStatic(DiskBalancerSubCommandUtil.class,
withSettings().defaultAnswer(
+ Mockito.CALLS_REAL_METHODS))) {
+
+ mockedUtil.when(() -> DiskBalancerSubCommandUtil
+ .getSingleNodeDiskBalancerProxy(resolvedAddress1))
+ .thenReturn(mockProtocol);
+ mockedUtil.when(() -> DiskBalancerSubCommandUtil
+ .getSingleNodeDiskBalancerProxy(resolvedAddress2))
+ .thenReturn(mockProtocol);
+
+ CommandLine c = new CommandLine(cmd);
+ c.parseArgs("--node-id", uuid1 + ",", uuid2);
+ cmd.call();
+
+ String output = outContent.toString(DEFAULT_ENCODING);
+ assertTrue(output.contains("Status result"));
+ assertTrue(output.contains(uuid1));
+ assertTrue(output.contains(uuid2));
+ }
+ }
+
+ @Test
+ public void testPositionalUuidRejected() throws Exception {
+ final String dnUuid = "a3b63511-bdf8-4fa1-8ab6-d19c0e806f84";
+ DiskBalancerStatusSubcommand cmd = new DiskBalancerStatusSubcommand();
+
+ CommandLine c = new CommandLine(cmd);
+ c.parseArgs(dnUuid);
+ cmd.call();
+
+ String err = errContent.toString(DEFAULT_ENCODING);
+ assertTrue(err.contains("Datanode UUID must be specified with --node-id"));
+ }
+
+ @Test
+ public void testResolutionFailuresDoNotLeakAcrossInvocations() throws
Exception {
+ final String invalidUuid = "00000000-0000-0000-0000-000000000000";
+ final String resolvedAddress = "127.0.0.1:" +
HDDS_DATANODE_CLIENT_PORT_DEFAULT;
+ final String expectedDisplay = "host-1 (" + resolvedAddress + ")";
+
+ DiskBalancerStatusSubcommand cmd = new DiskBalancerStatusSubcommand();
+ DatanodeDiskBalancerInfoProto statusProto =
generateRandomStatusProto("host-1");
+
+ try (MockedConstruction<ContainerOperationClient> mockedClient =
+ mockConstruction(ContainerOperationClient.class, (mock, context) ->
+ when(mock.queryNode(UUID.fromString(invalidUuid)))
+ .thenReturn(HddsProtos.Node.getDefaultInstance()));
+ MockedStatic<DiskBalancerSubCommandUtil> mockedUtil =
+ mockStatic(DiskBalancerSubCommandUtil.class,
withSettings().defaultAnswer(
+ Mockito.CALLS_REAL_METHODS))) {
+
+ CommandLine c = new CommandLine(cmd);
+ c.parseArgs("--node-id", invalidUuid);
+ cmd.call();
+ assertTrue(errContent.toString(DEFAULT_ENCODING).contains(invalidUuid));
+
+ outContent.reset();
+ errContent.reset();
+ when(mockProtocol.getDiskBalancerInfo()).thenReturn(statusProto);
+
+ Map<String, String> addressToDisplay = new LinkedHashMap<>();
+ addressToDisplay.put(resolvedAddress, expectedDisplay);
+ mockedUtil.when(() -> DiskBalancerSubCommandUtil
+ .getAllOperableNodesClientRpcAddress(any()))
+ .thenReturn(addressToDisplay);
+ mockedUtil.when(() -> DiskBalancerSubCommandUtil
+ .getSingleNodeDiskBalancerProxy(resolvedAddress))
+ .thenReturn(mockProtocol);
+
+ c.parseArgs("--in-service-datanodes");
+ cmd.call();
+
+ String err = errContent.toString(DEFAULT_ENCODING);
+ assertFalse(err.contains(invalidUuid));
+ assertTrue(outContent.toString(DEFAULT_ENCODING).contains("Status
result"));
+ }
+ }
+
+ @Test
+ public void testStatusStateDoesNotLeakAcrossInvocations() throws Exception {
+ DiskBalancerStatusSubcommand cmd = new DiskBalancerStatusSubcommand();
+ DatanodeDiskBalancerInfoProto statusProto1 =
generateRandomStatusProto("host-1");
+ DatanodeDiskBalancerInfoProto statusProto2 =
generateRandomStatusProto("host-2");
+
+ when(mockProtocol.getDiskBalancerInfo())
+ .thenReturn(statusProto1, statusProto2, statusProto2);
+
+ try (DiskBalancerMocks mocks = setupAllMocks()) {
+ CommandLine c = new CommandLine(cmd);
+ c.parseArgs("host-1", "host-2");
+ cmd.call();
+
+ outContent.reset();
+ errContent.reset();
+ c.parseArgs("host-2");
+ cmd.call();
+
+ String output = outContent.toString(DEFAULT_ENCODING);
+ assertTrue(output.contains("host-2"));
+ assertFalse(output.contains("host-1"));
+ }
+ }
+
+ @Test
+ public void testReportStateDoesNotLeakAcrossInvocations() throws Exception {
+ DiskBalancerReportSubcommand cmd = new DiskBalancerReportSubcommand();
+ DatanodeDiskBalancerInfoProto reportProto1 =
generateRandomReportProto("host-1");
+ DatanodeDiskBalancerInfoProto reportProto2 =
generateRandomReportProto("host-2");
+
+ when(mockProtocol.getDiskBalancerInfo())
+ .thenReturn(reportProto1, reportProto2, reportProto2);
+
+ try (DiskBalancerMocks mocks = setupAllMocks()) {
+ CommandLine c = new CommandLine(cmd);
+ c.parseArgs("host-1", "host-2");
+ cmd.call();
+
+ outContent.reset();
+ errContent.reset();
+ c.parseArgs("host-2");
+ cmd.call();
+
+ String output = outContent.toString(DEFAULT_ENCODING);
+ assertTrue(output.contains("host-2"));
+ assertFalse(output.contains("host-1"));
+ }
+ }
+
@Test
public void testStatusDiskBalancerFailure() throws Exception {
DiskBalancerStatusSubcommand cmd = new DiskBalancerStatusSubcommand();
@@ -778,6 +1042,7 @@ private DatanodeDiskBalancerInfoProto
createStatusProto(String hostname,
DatanodeDetailsProto nodeProto = DatanodeDetailsProto.newBuilder()
.setHostName(hostname)
.setIpAddress("127.0.0.1")
+
.setUuid(UUID.nameUUIDFromBytes(hostname.getBytes(StandardCharsets.UTF_8)).toString())
.addPorts(HddsProtos.Port.newBuilder()
.setName("CLIENT_RPC")
.setValue(HDDS_DATANODE_CLIENT_PORT_DEFAULT)
@@ -831,6 +1096,7 @@ private DatanodeDiskBalancerInfoProto
generateRandomReportProto(String hostname)
DatanodeDetailsProto nodeProto = DatanodeDetailsProto.newBuilder()
.setHostName(hostname)
.setIpAddress("127.0.0.1")
+
.setUuid(UUID.nameUUIDFromBytes(hostname.getBytes(StandardCharsets.UTF_8)).toString())
.addPorts(HddsProtos.Port.newBuilder()
.setName("CLIENT_RPC")
.setValue(HDDS_DATANODE_CLIENT_PORT_DEFAULT)
@@ -912,4 +1178,17 @@ private DiskBalancerConfigurationProto
createConfigProto(double threshold, long
.setStopAfterDiskEven(stopAfterDiskEven)
.build();
}
+
+ private static HddsProtos.Node buildScmNode(String uuid, String hostname,
String ipAddress) {
+ HddsProtos.DatanodeDetailsProto dnd =
HddsProtos.DatanodeDetailsProto.newBuilder()
+ .setUuid(uuid)
+ .setHostName(hostname)
+ .setIpAddress(ipAddress)
+ .addPorts(HddsProtos.Port.newBuilder()
+ .setName(DatanodeDetails.Port.Name.CLIENT_RPC.name())
+ .setValue(HDDS_DATANODE_CLIENT_PORT_DEFAULT)
+ .build())
+ .build();
+ return HddsProtos.Node.newBuilder().setNodeID(dnd).build();
+ }
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]