github-actions[bot] commented on code in PR #68630:
URL: https://github.com/apache/doris/pull/68630#discussion_r4132182080
##########
be/src/service/arrow_flight/arrow_flight_batch_reader.cpp:
##########
@@ -55,13 +91,79 @@ arrow::Status
ArrowFlightBatchReaderBase::_return_invalid_status(const std::stri
return arrow::Status::Invalid(status_msg);
}
+bool ArrowFlightBatchReaderBase::is_cancelled() const {
+ return _closed.load() || (_is_cancelled && _is_cancelled());
+}
+
+void ArrowFlightBatchReaderBase::close(const Status& reason) {
+ if (_closed.exchange(true) || _eof.load()) {
+ return;
+ }
+ auto* env = ExecEnv::GetInstance();
+ if (_statement->result_addr.hostname == BackendOptions::get_localhost() &&
+ _statement->result_addr.port == config::brpc_port) {
+ env->result_mgr()->cancel_arrow_flight_query(_statement->query_id,
reason);
+ }
+ // A result endpoint can outlive its local fragment. Find the owning FE by
buffer ID,
+ // then use its query-wide cancellation route, including for older result
BEs.
+ for (const auto& [address, info] : env->get_running_frontends()) {
Review Comment:
[P2] Bound FE discovery before cancelling the owning query. This loop visits
the address-sorted running-FE map synchronously; a stale FE before the owner
can trigger three connection attempts plus backoff on each of two passes, even
though the constructor is given 2000 ms. A reused Thrift client can also retain
the default 60000 ms receive timeout because that argument is applied only when
a client is created. The Flight stream has already closed while the other BEs
keep producing or block on buffers. Use an owner-directed route or one short
total deadline, including cached-client RPCs.
##########
fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/results/FlightSqlChannel.java:
##########
@@ -162,8 +166,16 @@ public long getAllocatedMemory() {
return allocator.getAllocatedMemory();
}
- public void reset() {
+ public synchronized void registerRemoteQuery(TUniqueId queryId,
List<TUniqueId> resultIds,
+ List<TNetworkAddress> backends, int timeoutSeconds) {
+ FlightSqlQueryCancellation.INSTANCE.register(queryId, resultIds,
backends, timeoutSeconds);
+ remoteResultIds.addAll(resultIds);
+ }
+
+ public synchronized void reset() {
resultCache.invalidateAll();
+ FlightSqlQueryCancellation.INSTANCE.unregister(remoteResultIds);
Review Comment:
[P2] Keep cancellation routes for earlier BE tickets across session reset.
`executeQueryStatement` calls `reset()` before every new statement, but a
previous FlightInfo's tickets remain readable directly from BE. If a client
obtains query A's parallel endpoints, starts query B, then closes one A stream
early, this invalidates A's result IDs here; `cancelFlightQuery` returns
NOT_FOUND and A's other BEs can remain blocked until timeout. Keep each route
until its own cancellation or execution-time expiry, and cancel outstanding
routes on session teardown. This is distinct from the prior finished-fragment
thread: the replacement FE route itself is removed here.
##########
be/test/service/arrow_flight/arrow_flight_cancel_test.cpp:
##########
@@ -0,0 +1,567 @@
+// 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.
+
+#include <arrow/flight/client.h>
+#include <arrow/flight/sql/server.h>
+#include <brpc/server.h>
+#include <gtest/gtest.h>
+#include <thrift/protocol/TBinaryProtocol.h>
+#include <thrift/server/TThreadedServer.h>
+#include <thrift/transport/TBufferTransports.h>
+#include <thrift/transport/TServerSocket.h>
+
+#include <future>
+#include <thread>
+
+#include "exec/pipeline/dependency.h"
+#include "load/channel/load_stream_mgr.h"
+#include "runtime/result_buffer_mgr.h"
+#include "service/arrow_flight/arrow_flight_batch_reader.h"
+#include "service/arrow_flight/flight_sql_service.h"
+#include "service/backend_options.h"
+#include "service/internal_service.h"
+#include "testutil/column_helper.h"
+#include "testutil/mock/mock_runtime_state.h"
+#include "util/brpc_client_cache.h"
+#include "util/client_cache.h"
+#include "util/dns_cache.h"
+
+namespace doris::flight {
+
+class ArrowFlightCancelTest : public testing::Test {
+protected:
+ void SetUp() override {
+ _previous_mgr = std::exchange(ExecEnv::GetInstance()->_result_mgr,
&_mgr);
+ _state._batch_size = 1;
+ _id.hi = 1;
+ _id.lo = 2;
+ auto schema = arrow::schema({arrow::field("value", arrow::int64())});
+ std::shared_ptr<ResultBlockBufferBase> buffer;
+ ASSERT_TRUE(_mgr.create_sender(_id, 16, &buffer, &_state, true,
schema).ok());
+ _buffer =
std::dynamic_pointer_cast<ArrowFlightResultBlockBuffer>(buffer);
+ _dep = Dependency::create_shared(0, 0, "Result", true);
+ _buffer->set_dependency(_state.fragment_instance_id(), _dep);
+ TNetworkAddress address;
+ address.hostname = BackendOptions::get_localhost();
+ address.port = config::brpc_port;
+ _statement = std::make_shared<QueryStatement>(_id, address, "select
value");
+ }
+ void TearDown() override { ExecEnv::GetInstance()->_result_mgr =
_previous_mgr; }
Review Comment:
[P2] Scope the global metrics hook to this fixture's manager. Each case
constructs a `ResultBufferMgr`, whose constructor registers
`result_buffer_block_count` with a callback capturing `this`; its default
destructor does not unregister it. `TearDown()` only restores
`ExecEnv::_result_mgr`, so a filtered run of this new suite leaves a hook
pointing into the first destroyed fixture. A later metrics scrape can access
freed storage, and registration by later cases uses `emplace` rather than
replacing that hook. Remove or restore only the hook this fixture owns before
`_mgr` is destroyed; preserve any earlier manager's registration.
##########
fe/fe-core/src/main/java/org/apache/doris/service/arrowflight/FlightSqlQueryCancellation.java:
##########
@@ -0,0 +1,146 @@
+// 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.service.arrowflight;
+
+import org.apache.doris.common.Status;
+import org.apache.doris.proto.InternalService.PCancelPlanFragmentResult;
+import org.apache.doris.rpc.BackendServiceProxy;
+import org.apache.doris.thrift.TNetworkAddress;
+import org.apache.doris.thrift.TStatus;
+import org.apache.doris.thrift.TStatusCode;
+import org.apache.doris.thrift.TUniqueId;
+
+import com.github.benmanes.caffeine.cache.Cache;
+import com.github.benmanes.caffeine.cache.Caffeine;
+import com.github.benmanes.caffeine.cache.Expiry;
+import com.github.benmanes.caffeine.cache.Scheduler;
+import com.github.benmanes.caffeine.cache.Ticker;
+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;
+import java.util.stream.Collectors;
+
+public final class FlightSqlQueryCancellation {
+ private static final Logger LOG =
LogManager.getLogger(FlightSqlQueryCancellation.class);
+ public static final FlightSqlQueryCancellation INSTANCE = new
FlightSqlQueryCancellation(Ticker.systemTicker());
+ private final Cache<TUniqueId, Route> results;
+
+ FlightSqlQueryCancellation(Ticker ticker) {
+ results =
Caffeine.newBuilder().ticker(ticker).scheduler(Scheduler.systemScheduler())
+ .expireAfter(new Expiry<TUniqueId, Route>() {
+ @Override
+ public long expireAfterCreate(TUniqueId key, Route route,
long now) {
+ return route.ttlNanos;
+ }
+
+ @Override
+ public long expireAfterUpdate(TUniqueId key, Route route,
long now, long duration) {
+ return route.ttlNanos;
+ }
+
+ @Override
+ public long expireAfterRead(TUniqueId key, Route route,
long now, long duration) {
+ return duration;
+ }
+ }).build();
+ }
+
+ public void register(TUniqueId queryId, List<TUniqueId> resultIds,
+ List<TNetworkAddress> backends, int timeoutSeconds) {
+ if (backends.isEmpty()) {
+ throw new IllegalArgumentException("Flight query has no
cancellation backends");
+ }
+ // Keep only cancellation addresses, not the coordinator, scan state,
or query queue slot.
+ // Ordinary Flight coordinators are unregistered when GetFlightInfo
returns.
+ Route route = new Route(queryId.deepCopy(),
+
resultIds.stream().map(TUniqueId::deepCopy).distinct().collect(Collectors.toList()),
+
backends.stream().map(TNetworkAddress::deepCopy).distinct().collect(Collectors.toList()),
+ TimeUnit.SECONDS.toNanos(Math.max(0L, timeoutSeconds) + 5));
+ for (TUniqueId resultId : route.resultIds) {
+ results.put(resultId, route);
+ }
+ }
+
+ public void unregister(List<TUniqueId> resultIds) {
+ results.invalidateAll(resultIds);
+ }
+
+ public TStatus cancel(TUniqueId resultId) {
+ Route route = results.getIfPresent(resultId);
+ if (route == null) {
+ return new TStatus(TStatusCode.NOT_FOUND);
+ }
+ Status reason = new Status(TStatusCode.CANCELLED, "Arrow Flight stream
closed before EOF");
+ TStatus status = new TStatus(TStatusCode.OK);
+ List<Future<PCancelPlanFragmentResult>> futures = new ArrayList<>();
+ long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(1);
+ for (TNetworkAddress backend : route.backends) {
+ try {
+ // The existing RPC is understood by older result BEs and
bypasses the Flight read pool.
+ futures.add(BackendServiceProxy.getInstance()
+ .cancelPipelineXPlanFragmentAsync(backend,
route.queryId, reason));
Review Comment:
[P2] Ensure abort reaches a BE when its light pool rejects cancellation.
This new fanout calls `cancel_plan_fragment`, which enters the bounded
`_light_work_pool`; on `try_offer` failure the BE returns a protobuf CANCELLED
before cancelling the query or its Flight buffers. The FE retains this route,
but the closing reader makes only two immediate attempts and sets `_closed`, so
under sustained light-pool pressure no later caller retries and a producer can
remain blocked until timeout. Give this cancel reliable admission or retain a
retry owner until every BE acknowledges it. This is a different pool and RPC
from the earlier Arrow fetch-pool thread.
--
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]