yandrey321 commented on code in PR #11294:
URL: https://github.com/apache/ozone/pull/11294#discussion_r4147164305
##########
hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/XceiverClientShortCircuit.java:
##########
@@ -510,33 +481,37 @@ public void run() {
entry.getFuture().complete(responseProto);
}
long currentTime = System.nanoTime();
- long endToEndCost = currentTime - entry.getCreateTimeNs();
- long sentCost = entry.getSentTimeNs() - entry.getCreateTimeNs();
- long receiveCost = processStartTime - receiveStartTime;
- long processCost = currentTime - processStartTime;
+ endToEndCost = currentTime - entry.getCreateTimeNs();
if (LOG.isDebugEnabled()) {
- LOG.debug("Executed command {} {}:{} on datanode {}, end-to-end {}
ns, sent {} ns, receive {} ns, " +
- "process {} ns", type,
entry.getRequest().getClientId().toStringUtf8(),
- entry.getRequest().getCallId(), dn, endToEndCost, sentCost,
receiveCost, processCost);
+ LOG.debug("Executed command {} {}:{} on datanode {}, end-to-end
{}ns, receive {}ns, process {}ns",
+ type, entry.getRequest().getClientId().toStringUtf8(),
entry.getRequest().getCallId(), dn,
+ endToEndCost,
+ processStartTime - receiveStartTime,
+ currentTime - processStartTime);
}
- responseReceived++;
- metrics.decrPendingContainerOpsMetrics(type);
- metrics.addContainerOpsLatency(type, endToEndCost);
+ responseReceived.incrementAndGet();
} catch (SocketTimeoutException | EOFException |
ClosedChannelException e) {
isDomainSocketOpen.set(false);
- LOG.info("{} receiveResponseTask is closed after send {} requests
and received {} responses, due to {}",
- domainSocket.toString(), requestSent, responseReceived,
e.getClass().getName(), e);
+ LOG.info("ReceiveResponseTask closed: (sent {}, received {}): {}",
+ requestSent, responseReceived, XceiverClientShortCircuit.this,
e);
// fail all requests pending responses
sentRequests.values().forEach(i -> i.fail(e));
} catch (Throwable e) {
isDomainSocketOpen.set(false);
- LOG.error("{} failed after send {} requests and received {}
responses",
- domainSocket.toString(), requestSent, responseReceived, e);
+ LOG.error("ReceiveResponseTask failed: (sent {}, received {}): {}",
+ requestSent, responseReceived, XceiverClientShortCircuit.this,
e);
if (entry != null) {
entry.getFuture().completeExceptionally(e);
}
sentRequests.values().forEach(i -> i.fail(e));
break;
+ } finally {
Review Comment:
The new `finally` fixes the real leak on the GetBlock FD-exchange `continue`
path, but
keying the decrement off `type != null` also fires it on exceptional exits
where no
pending op is outstanding — so `numPending<Type>` is now decremented twice:
1. Late response after a timeout (exactly the case the comment at L426-428
describes):
`requestTimeout` already did
`metrics.decrPendingContainerOpsMetrics(type)` at L333 and
removed the entry, so `sentRequests.remove(...)` at L424 returns null,
L429 throws,
`catch (Throwable)` runs, and this `finally` decrements a second time.
2. A response truncated after the type short was read (EOF/timeout inside
`parseFrom`):
`type` is non-null, so we decrement here, while the entry is still in
`sentRequests`
and gets `fail()`ed at L498 with no decrement of its own.
`decrPendingContainerOpsMetrics` is a bare `MutableCounterLong.incr(-1)`
(`XceiverClientMetrics:96-98`) on the process-wide singleton from
`XceiverClientManager.getXceiverClientMetrics()`, so the skew is not
per-client and is not
self-correcting — `numPendingGetBlock` can go negative and stays wrong for
the JVM lifetime.
Could the decrement be tied to entry ownership instead of to `type`? i.e.
decrement exactly
where an entry leaves the pending set: on a non-null
`sentRequests.remove(...)` at L424,
in `requestTimeout`, and in `fail`. Then `finally` only needs
`addContainerOpsLatency`.
##########
hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/XceiverClientShortCircuit.java:
##########
@@ -653,4 +624,35 @@ public boolean equals(Object obj) {
&& this.blockLocalId == that.blockLocalId;
}
}
+
+ static final class TimeoutScheduler {
+ private Timer timer;
+
+ synchronized void init(String prefix) {
+ if (timer != null) {
+ timer.cancel();
+ }
+ timer = new Timer(prefix + "-Timer");
+ }
+
+ synchronized void schedule(TimerTask task, int timeoutMs) {
Review Comment:
`Objects.requireNonNull(timer, "Timer is null")` makes the send/close race
throw NPE, and
`sendRequest` only catches `IOException` (L381), so the NPE escapes
`sendRequest` ->
`sendCommandInternal` -> `sendCommandWithoutTraceID` (whose outer try
catches only
`ExecutionException`/`InterruptedException`) and reaches the caller as an
unchecked
exception, with `entry`'s future left uncompleted.
The window is real: `checkOpen()` (L298) and `close()` (L158) are both
synchronized, but
`scheduler.schedule` at L349 runs outside that monitor, so a concurrent
`close()` can null
the timer in between. Before this PR the same race produced
`IllegalStateException("Timer already cancelled")` — also unchecked — so
this is a change of
symptom rather than a new bug class, but since the PR owns this lifecycle
now it would be
good to close it: either have `schedule` throw `IOException` (or return
false) on a
closed scheduler so `sendRequest`'s existing handler completes the future
exceptionally,
or keep the cancelled `Timer` instance instead of nulling it in `close()`.
##########
hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/XceiverClientShortCircuit.java:
##########
@@ -137,11 +142,13 @@ public void connect() throws IOException {
return;
}
domainSocket = domainSocketFactory.createSocket(readTimeoutMs,
writeTimeoutMs, dnAddr);
+ updateName("Connected-" + domainSocket);
isDomainSocketOpen.set(true);
- prefix = XceiverClientShortCircuit.class.getSimpleName() + "-" +
domainSocket.toString();
- timer = new Timer(prefix + "-Timer");
+ final String prefix = XceiverClientShortCircuit.class.getSimpleName() +
"-" + domainSocket;
+ scheduler.init(prefix);
Review Comment:
`scheduler.init` is written for reconnect — it cancels and replaces the
`Timer` — but two
things make reconnect look unreachable:
- `readDaemon` is a single `Daemon` instance, so a second `connect()` hits
`readDaemon.start()` at L150 with `IllegalThreadStateException`.
- `init()` cancelling the old timer silently drops the timeout tasks of any
requests still
in `sentRequests`; those futures then never complete and `response.get()`
at L262 blocks
forever.
So either reconnect is unreachable, in which case the `if (timer != null)
timer.cancel()`
branch is dead and an assertion would document the invariant better; or it
is reachable and
the daemon restart is the actual bug. Which is intended here?
##########
hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/XceiverClientShortCircuit.java:
##########
@@ -193,23 +200,12 @@ public long getCallId() {
@Override
public ContainerCommandResponseProto
sendCommand(ContainerCommandRequestProto request) throws IOException {
- try {
- return sendCommandWithTraceID(request, null).getResponse().get();
- } catch (ExecutionException e) {
- throw getIOExceptionForSendCommand(request, e);
- } catch (InterruptedException e) {
- LOG.error("Command execution was interrupted.");
- Thread.currentThread().interrupt();
- throw (IOException) new InterruptedIOException(
- "Command " + processForDebug(request) + " was interrupted.")
- .initCause(e);
- }
+ return sendCommandWithTraceID(request);
}
@Override
- public Map<DatanodeDetails, ContainerCommandResponseProto>
- sendCommandOnAllNodes(
- ContainerCommandRequestProto request) throws IOException {
+ public Map<DatanodeDetails, ContainerCommandResponseProto>
sendCommandOnAllNodes(
Review Comment:
Dropping `throws IOException` from this override is unrelated to the stated
goal of the PR.
It is source-compatible (no subclasses, callers go through
`XceiverClientSpi`), but it is
churn in a change that is otherwise about logging and the timer.
##########
hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/XceiverClientShortCircuit.java:
##########
@@ -380,7 +346,7 @@ public void run() {
}
};
entry.setTimerTask(task);
- timer.schedule(task, readTimeoutMs);
+ scheduler.schedule(task, readTimeoutMs);
Review Comment:
The timer task is scheduled (L349) before the entry is published (L350). If
the task fires
in that window, `requestTimeout` removes nothing and the entry is then
inserted carrying a
task that has already run — the request never times out and `response.get()`
at L262 blocks
indefinitely. The window needs `readTimeoutMs` to elapse between two
adjacent statements, so
it is remote in practice, but swapping `put` and `schedule` removes it at no
cost, and this
PR is already rewriting these lines.
##########
hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/XceiverClientShortCircuit.java:
##########
@@ -150,15 +157,15 @@ public void connect() throws IOException {
@Override
public synchronized void close() {
closed = true;
- timer.cancel();
+ scheduler.close();
if (domainSocket != null) {
try {
isDomainSocketOpen.set(false);
domainSocket.close();
- LOG.info("{} is closed for {} with {} requests sent and {} responses
received",
- domainSocket.toString(), dn, requestSent, responseReceived);
+ LOG.info("Closed successfully (sent {}, received {}): {}",
requestSent, responseReceived, this);
Review Comment:
`updateName("Closed-" + domainSocket)` sits inside the `try` after
`domainSocket.close()`,
so when close throws, the name stays `Connected-…` while `closed == true`,
and a subsequent
`checkOpen()` reports `Closed: XceiverClientShortCircuit[ Connected-…, dn]`.
Moving the
`updateName` above the `close()` call (or into a `finally`) keeps the name
honest — the
whole point of the field.
##########
hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/XceiverClientShortCircuit.java:
##########
@@ -404,25 +370,29 @@ public void run() {
dataOut.flush();
} finally {
Review Comment:
Pre-existing, but now more visible: `requestSent.incrementAndGet()` is in
the `finally` of
the write block, so a request whose write threw still counts as sent — and
that counter now
appears in three log messages as "sent N". Moving it after `dataOut.flush()`
makes the
number mean what the messages claim.
##########
hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/XceiverClientShortCircuit.java:
##########
@@ -605,24 +579,21 @@ public long getCreateTimeNs() {
return createTimeNs;
}
- public long getSentTimeNs() {
- return sentTimeNs;
- }
-
- public void setSentTimeNs() {
- sentTimeNs = System.nanoTime();
- }
-
public void setTimerTask(TimerTask task) {
- timerTask = task;
+ final boolean set = timerTask.compareAndSet(null, task);
+ Preconditions.assertTrue(set);
}
- public TimerTask getTimerTask() {
- return timerTask;
+ boolean cancelTimerTask() {
+ final TimerTask t = timerTask.getAndSet(null);
+ if (t == null) {
+ return false;
+ }
+ return t.cancel();
}
public void fail(Throwable e) {
Review Comment:
`fail` completes the future but neither removes the entry from
`sentRequests` nor
decrements the pending-ops metric, and both callers are
`sentRequests.values().forEach(i -> i.fail(e))` — so after a socket error
every in-flight
request is retained in the map and its `numPending<Type>` is never released.
This is not bounded by the client's lifetime: `XceiverClientManager` keeps
clients in a
Guava `clientCache` (`XceiverClientManager.java:69,93`), and a short-circuit
client whose
socket died stays cached until eviction with `isDomainSocketOpen == false`,
so the map and
the counters persist. This confirms @echonesis' retention comment and adds
the metrics side
of it.
Draining rather than iterating would cover both, e.g.
for (Iterator<RequestEntry> i = sentRequests.values().iterator();
i.hasNext(); ) {
RequestEntry pending = i.next();
i.remove();
pending.fail(e);
}
with the pending-op decrement moved inside `fail`.
##########
hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/XceiverClientShortCircuit.java:
##########
@@ -350,19 +325,10 @@ public HddsProtos.ReplicationType getPipelineType() {
return HddsProtos.ReplicationType.STAND_ALONE;
}
- public ConfigurationSource getConfig() {
- return config;
- }
-
- @VisibleForTesting
- public static Logger getLogger() {
- return LOG;
- }
-
void requestTimeout(RequestKey requestKey) {
final RequestEntry entry = sentRequests.remove(requestKey);
if (entry != null) {
- LOG.warn("Timeout to receive response for command {}",
entry.getRequest());
+ LOG.warn("Timeout: {}", processForDebug(entry.getRequest()));
Review Comment:
This message loses the client identity that the rest of the PR is adding
everywhere else —
with several clients per JVM there is nothing to tell them apart.
`LOG.warn("Timeout: {}
{}", processForDebug(entry.getRequest()), this)` would match the other sites.
--
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]