mrhhsg commented on code in PR #68032:
URL: https://github.com/apache/doris/pull/68032#discussion_r4228873732
##########
fe/fe-core/src/main/java/org/apache/doris/common/proc/CurrentQueryStatisticsProcDir.java:
##########
@@ -45,7 +45,9 @@ public class CurrentQueryStatisticsProcDir implements
ProcDirInterface {
.add("ScanBytesFromLocalStorage").add("ScanBytesFromRemoteStorage")
.add("SpillWriteBytesToLocalStorage").add("SpillReadBytesFromLocalStorage")
.add("BytesWriteIntoCache")
- .add("TotalTasks").add("FinishedTasks").add("Progress").build();
+ .add("TotalTasks").add("FinishedTasks").add("Progress")
+ // Appended last: multi-FE aggregation concatenates rows by
position.
Review Comment:
已在 d1b5f50489c 按各 FE 返回的 NodeInfo.columnNames 对齐 current_queries 行;旧 FE
缺失的两个 remote-spill 字段补格式化零值,行宽不合法则明确拒绝,避免错列。新增旧版、新版、重排列名及畸形行单测;FE 定向测试 57/57 通过。
##########
fe/fe-core/src/main/java/org/apache/doris/qe/AuditLogHelper.java:
##########
@@ -274,6 +274,10 @@ private static void logAuditLogImpl(ConnectContext ctx,
String origStmt, Stateme
statistics.getSpillWriteBytesToLocalStorage())
.setSpillReadBytesFromLocalStorage(statistics == null ? 0 :
statistics.getSpillReadBytesFromLocalStorage())
+ // Remote spill bytes are only reported through
TQueryStatistics; for queries
Review Comment:
已在 d1b5f50489c 让普通查询及转发查询传递实际 dispatch 的 BE ID,并复用最终 TQueryStatistics
报告屏障;审计不会仅因普通 audit timeout 到期就提前写入零值。新增延迟超过该 timeout 的两 BE
最终报告测试,验证远端读写字节合并;保留 BE 永不汇报时原有的有界回退。FE 定向测试 57/57 通过。
##########
fe/fe-core/src/main/java/org/apache/doris/cloud/catalog/RemoteSpillStatsPoller.java:
##########
@@ -0,0 +1,161 @@
+// 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.doris.cloud.catalog;
+
+import org.apache.doris.catalog.Env;
+import org.apache.doris.common.AnalysisException;
+import org.apache.doris.common.Config;
+import org.apache.doris.common.Pair;
+import org.apache.doris.common.Status;
+import org.apache.doris.common.util.MasterDaemon;
+import org.apache.doris.proto.InternalService;
+import org.apache.doris.rpc.BackendServiceProxy;
+import org.apache.doris.system.Backend;
+import org.apache.doris.thrift.TStatusCode;
+
+import com.google.common.annotations.VisibleForTesting;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+
+/**
+ * Polls the bytes of query spill held in object storage
(spill_storage_type=s3) from the alive
+ * backends of all clusters, so that SHOW DATA reads them from memory. Runs on
every FE on its own
+ * schedule (cloud_spill_stats_poll_interval_second): the freshness of the
value must not depend on
+ * how long a round of the tablet stats takes.
+ */
+public class RemoteSpillStatsPoller extends MasterDaemon {
+ private static final Logger LOG =
LogManager.getLogger(RemoteSpillStatsPoller.class);
+
+ private static final int RPC_TIMEOUT_SECOND = 5;
+
+ /**
+ * One successful poll: the value and when it was fetched, on the
monotonic clock so that a wall
+ * clock moved backward cannot keep a stale value fresh.
+ */
+ private static final class RemoteSpillStats {
+ private final long bytes;
+ private final long fetchTimeNanos;
+
+ private RemoteSpillStats(long bytes, long fetchTimeNanos) {
+ this.bytes = bytes;
+ this.fetchTimeNanos = fetchTimeNanos;
+ }
+ }
+
+ // Summed over the alive BEs of all clusters. Null until the first
successful poll. A BE that is
+ // gone no longer contributes.
+ private volatile RemoteSpillStats remoteSpillStats = null;
+
+ public RemoteSpillStatsPoller() {
+ super("remote spill stats poller", pollIntervalMs());
+ }
+
+ private static long pollIntervalMs() {
+ return Math.max(1, Config.cloud_spill_stats_poll_interval_second) *
1000L;
+ }
+
+ /**
+ * A value older than this is not served: the configured max age, but at
least three poll
+ * intervals so that a longer interval cannot make every value stale.
+ */
+ @VisibleForTesting
+ static long maxAgeSecond() {
+ return Math.max(Config.cloud_spill_stats_max_age_second,
+ 3L * Math.max(1,
Config.cloud_spill_stats_poll_interval_second));
+ }
+
+ @Override
+ protected void runAfterCatalogReady() {
+ refresh();
+ // The interval is mutable.
+ setInterval(pollIntervalMs());
+ }
+
+ private void refresh() {
+ List<Backend> backends;
+ try {
+ backends =
Env.getCurrentSystemInfo().getAllBackendsByAllCluster().values().asList();
+ } catch (AnalysisException e) {
+ LOG.warn("failed to list the backends for the remote spill stats",
e);
+ return;
+ }
+ InternalService.PGetBeResourceRequest request =
InternalService.PGetBeResourceRequest.newBuilder().build();
+ List<Pair<Backend, Future<InternalService.PGetBeResourceResponse>>>
futures = new ArrayList<>();
+ for (Backend be : backends) {
+ if (!be.isAlive()) {
Review Comment:
已在 d1b5f50489c 移除心跳 isAlive 对统计 BRPC 的预先过滤:所有已注册 BE 均尝试 RPC,心跳失败但 BRPC
可达的字节会计入;RPC 发送或返回失败则保留旧完整值且不刷新时间戳。新增心跳失败而 BRPC 成功、以及失败轮次保留旧值的测试;FE 定向测试 57/57
通过。
--
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]