From 330ddcc6d02da905b4e96c6779acdd3d7e7f83ec Mon Sep 17 00:00:00 2001 From: Hor911 Date: Fri, 28 Mar 2025 21:20:49 +0300 Subject: Fix timeouts on heavy loads (#16340) --- .../kqp/compute_actor/kqp_scan_compute_manager.cpp | 3 +- .../kqp/compute_actor/kqp_scan_compute_manager.h | 12 ++++++++ .../kqp/compute_actor/kqp_scan_fetcher_actor.cpp | 6 ++++ .../kqp/compute_actor/kqp_scan_fetcher_actor.h | 3 ++ .../tx/columnshard/engines/reader/actor/actor.cpp | 33 ++++++++++++++-------- .../tx/columnshard/engines/reader/actor/actor.h | 4 ++- 6 files changed, 47 insertions(+), 14 deletions(-) diff --git a/ydb/core/kqp/compute_actor/kqp_scan_compute_manager.cpp b/ydb/core/kqp/compute_actor/kqp_scan_compute_manager.cpp index 7eab5ef7819..1b408cdf6f8 100644 --- a/ydb/core/kqp/compute_actor/kqp_scan_compute_manager.cpp +++ b/ydb/core/kqp/compute_actor/kqp_scan_compute_manager.cpp @@ -27,10 +27,9 @@ std::vector> TShardScannerInfo::OnReceiveData( AFL_ENSURE(data.Finished); result.emplace_back(std::make_unique(selfPtr, std::make_unique(TabletId, data.LocksInfo))); } else if (data.SplittedBatches.size() > 1) { - ui32 idx = 0; AFL_ENSURE(data.ArrowBatch); for (auto&& i : data.SplittedBatches) { - result.emplace_back(std::make_unique(selfPtr, std::make_unique(data.ArrowBatch, TabletId, std::move(i), data.LocksInfo), idx++)); + result.emplace_back(std::make_unique(selfPtr, std::make_unique(data.ArrowBatch, TabletId, std::move(i), data.LocksInfo))); } } else if (data.ArrowBatch) { result.emplace_back(std::make_unique(selfPtr, std::make_unique(data.ArrowBatch, TabletId, data.LocksInfo))); diff --git a/ydb/core/kqp/compute_actor/kqp_scan_compute_manager.h b/ydb/core/kqp/compute_actor/kqp_scan_compute_manager.h index 67b7ff64bee..db4b753e68f 100644 --- a/ydb/core/kqp/compute_actor/kqp_scan_compute_manager.h +++ b/ydb/core/kqp/compute_actor/kqp_scan_compute_manager.h @@ -105,6 +105,12 @@ public: return !ActorId.has_value(); } + void Ping() { + if (ActorId) { + NActors::TActivationContext::AsActorContext().Send(*ActorId, new TEvKqpCompute::TEvScanDataAck(0, 0, 0), IEventHandle::FlagTrackDelivery, TabletId); + } + } + void Start(const TActorId& actorId) { AFL_DEBUG(NKikimrServices::KQP_COMPUTE)("event", "start_scanner")("actor_id", actorId); AFL_ENSURE(!ActorId); @@ -284,6 +290,12 @@ public: } } + void PingAllScanners() { + for (auto&& itTablet : ShardScanners) { + itTablet.second->Ping(); + } + } + std::shared_ptr GetShardStateByActorId(const NActors::TActorId& actorId) const { auto it = ShardsByActorId.find(actorId); if (it == ShardsByActorId.end()) { diff --git a/ydb/core/kqp/compute_actor/kqp_scan_fetcher_actor.cpp b/ydb/core/kqp/compute_actor/kqp_scan_fetcher_actor.cpp index 0a383b154a6..d3a2ccc5d36 100644 --- a/ydb/core/kqp/compute_actor/kqp_scan_fetcher_actor.cpp +++ b/ydb/core/kqp/compute_actor/kqp_scan_fetcher_actor.cpp @@ -82,6 +82,7 @@ void TKqpScanFetcherActor::Bootstrap() { AFL_DEBUG(NKikimrServices::KQP_COMPUTE)("event", "bootstrap")("compute", ComputeActorIds.size())("shards", PendingShards.size()); StartTableScan(); Become(&TKqpScanFetcherActor::StateFunc); + Schedule(TDuration::Seconds(30), new NActors::TEvents::TEvWakeup()); } void TKqpScanFetcherActor::HandleExecute(TEvScanExchange::TEvAckData::TPtr& ev) { @@ -676,4 +677,9 @@ void TKqpScanFetcherActor::CheckFinish() { } } +void TKqpScanFetcherActor::HandleExecute(NActors::TEvents::TEvWakeup::TPtr&) { + InFlightShards.PingAllScanners(); + Schedule(TDuration::Seconds(30), new NActors::TEvents::TEvWakeup()); +} + } diff --git a/ydb/core/kqp/compute_actor/kqp_scan_fetcher_actor.h b/ydb/core/kqp/compute_actor/kqp_scan_fetcher_actor.h index 73ead0a5ef0..18d6e050039 100644 --- a/ydb/core/kqp/compute_actor/kqp_scan_fetcher_actor.h +++ b/ydb/core/kqp/compute_actor/kqp_scan_fetcher_actor.h @@ -82,6 +82,7 @@ public: hFunc(TEvInterconnect::TEvNodeDisconnected, HandleExecute); hFunc(TEvScanExchange::TEvTerminateFromCompute, HandleExecute); hFunc(TEvScanExchange::TEvAckData, HandleExecute); + hFunc(NActors::TEvents::TEvWakeup, HandleExecute); IgnoreFunc(TEvInterconnect::TEvNodeConnected); IgnoreFunc(TEvTxProxySchemeCache::TEvInvalidateTableResult); default: @@ -97,6 +98,8 @@ public: void HandleExecute(TEvScanExchange::TEvTerminateFromCompute::TPtr& ev); + void HandleExecute(NActors::TEvents::TEvWakeup::TPtr& ev); + private: void CheckFinish(); diff --git a/ydb/core/tx/columnshard/engines/reader/actor/actor.cpp b/ydb/core/tx/columnshard/engines/reader/actor/actor.cpp index d0c96097a6a..fb589150996 100644 --- a/ydb/core/tx/columnshard/engines/reader/actor/actor.cpp +++ b/ydb/core/tx/columnshard/engines/reader/actor/actor.cpp @@ -7,8 +7,8 @@ #include namespace NKikimr::NOlap::NReader { -constexpr TDuration SCAN_HARD_TIMEOUT = TDuration::Minutes(10); -constexpr TDuration SCAN_HARD_TIMEOUT_GAP = TDuration::Seconds(5); +constexpr TDuration SCAN_HARD_TIMEOUT = TDuration::Minutes(60); +constexpr TDuration COMPUTE_HARD_TIMEOUT = TDuration::Minutes(10); void TColumnShardScan::PassAway() { Send(ResourceSubscribeActorId, new TEvents::TEvPoisonPill); @@ -34,7 +34,7 @@ TColumnShardScan::TColumnShardScan(const TActorId& columnShardActorId, const TAc , DataFormat(dataFormat) , TabletId(tabletId) , ReadMetadataRange(readMetadataRange) - , Timeout(timeout ? timeout + SCAN_HARD_TIMEOUT_GAP : SCAN_HARD_TIMEOUT) + , Timeout(timeout ? timeout : COMPUTE_HARD_TIMEOUT) , ScanCountersPool(scanCountersPool, TValidator::CheckNotNull(ReadMetadataRange)->GetProgram().GetGraphOptional()) , Stats(NTracing::TTraceClient::GetLocalClient("SHARD", ::ToString(TabletId) /*, "SCAN_TXID:" + ::ToString(TxId)*/)) , ComputeShardingPolicy(computeShardingPolicy) { @@ -93,6 +93,14 @@ void TColumnShardScan::HandleScan(NColumnShard::TEvPrivate::TEvTaskProcessedResu void TColumnShardScan::HandleScan(NKqp::TEvKqpCompute::TEvScanDataAck::TPtr& ev) { auto g = Stats->MakeGuard("ack"); + + if (ev->Get()->FreeSpace == 0 && ev->Get()->MaxChunksCount == 0) { + if (!AckReceivedInstant) { + LastResultInstant = TMonotonic::Now(); + } + return; + } + AFL_VERIFY(!AckReceivedInstant); AckReceivedInstant = TMonotonic::Now(); @@ -165,10 +173,10 @@ void TColumnShardScan::HandleScan(TEvents::TEvWakeup::TPtr& /*ev*/) { << " txId: " << TxId << " scanId: " << ScanId << " gen: " << ScanGen << " tablet: " << TabletId); CheckHanging(true); - if (!!AckReceivedInstant && TMonotonic::Now() >= GetDeadline() + Timeout * 0.5) { + if (!!AckReceivedInstant && TMonotonic::Now() >= GetScanDeadline()) { SendScanError("ColumnShard scanner timeout: HAS_ACK=1"); Finish(NColumnShard::TScanCounters::EStatusFinish::Deadline); - } else if (!AckReceivedInstant && TMonotonic::Now() >= GetDeadline()) { + } else if (!AckReceivedInstant && TMonotonic::Now() >= GetComputeDeadline()) { SendScanError("ColumnShard scanner timeout: HAS_ACK=0"); Finish(NColumnShard::TScanCounters::EStatusFinish::Deadline); } else { @@ -423,11 +431,14 @@ void TColumnShardScan::ScheduleWakeup(const TMonotonic deadline) { } } -TMonotonic TColumnShardScan::GetDeadline() const { - AFL_VERIFY(StartInstant); - if (LastResultInstant) { - return *LastResultInstant + Timeout; - } - return *StartInstant + Timeout; +TMonotonic TColumnShardScan::GetScanDeadline() const { + AFL_VERIFY(!!AckReceivedInstant); + return *AckReceivedInstant + SCAN_HARD_TIMEOUT; } + +TMonotonic TColumnShardScan::GetComputeDeadline() const { + AFL_VERIFY(!AckReceivedInstant); + return (LastResultInstant ? *LastResultInstant : *StartInstant) + Timeout; +} + } // namespace NKikimr::NOlap::NReader diff --git a/ydb/core/tx/columnshard/engines/reader/actor/actor.h b/ydb/core/tx/columnshard/engines/reader/actor/actor.h index a66d1ad7f6a..c0a859ccbac 100644 --- a/ydb/core/tx/columnshard/engines/reader/actor/actor.h +++ b/ydb/core/tx/columnshard/engines/reader/actor/actor.h @@ -109,7 +109,9 @@ private: void ScheduleWakeup(const TMonotonic deadline); - TMonotonic GetDeadline() const; + TMonotonic GetScanDeadline() const; + + TMonotonic GetComputeDeadline() const; private: const TActorId ColumnShardActorId; -- cgit v1.3