From 24eae22a6dca53c86ffbb77a942e7f2dac840df0 Mon Sep 17 00:00:00 2001 From: Pisarenko Grigoriy Date: Fri, 26 Sep 2025 17:37:57 +0500 Subject: YQ-4664 used AS threads in topic sdk IO operations (#25668) --- ydb/core/fq/libs/row_dispatcher/topic_session.cpp | 28 ++++++--- ydb/core/fq/libs/row_dispatcher/ya.make | 1 + .../yql/providers/pq/async_io/dq_pq_read_actor.cpp | 53 +++++++++------- .../providers/pq/async_io/dq_pq_write_actor.cpp | 38 ++++++----- .../yql/providers/pq/common/pq_events_processor.h | 73 ++++++++++++++++++++++ ydb/library/yql/providers/pq/common/ya.make | 2 + 6 files changed, 146 insertions(+), 49 deletions(-) create mode 100644 ydb/library/yql/providers/pq/common/pq_events_processor.h diff --git a/ydb/core/fq/libs/row_dispatcher/topic_session.cpp b/ydb/core/fq/libs/row_dispatcher/topic_session.cpp index dd5b739c495..635c349deaa 100644 --- a/ydb/core/fq/libs/row_dispatcher/topic_session.cpp +++ b/ydb/core/fq/libs/row_dispatcher/topic_session.cpp @@ -5,13 +5,11 @@ #include #include #include - #include #include #include - +#include #include - #include #include @@ -62,6 +60,7 @@ struct TEvPrivate { EvCreateSession, EvSendStatistic, EvReconnectSession, + EvExecuteTopicEvent, EvEnd }; static_assert(EvEnd < EventSpaceEnd(TEvents::ES_PRIVATE), "expect EvEnd < EventSpaceEnd(TEvents::ES_PRIVATE)"); @@ -71,13 +70,16 @@ struct TEvPrivate { struct TEvCreateSession : public TEventLocal {}; struct TEvSendStatistic : public TEventLocal {}; struct TEvReconnectSession : public TEventLocal {}; + struct TEvExecuteTopicEvent : public NYql::TTopicEventBase { + using TTopicEventBase::TTopicEventBase; + }; }; constexpr ui64 SendStatisticPeriodSec = 2; constexpr ui64 MaxHandledEventsCount = 1000; constexpr ui64 MaxHandledEventsSize = 1000000; -class TTopicSession : public TActorBootstrapped { +class TTopicSession : public TActorBootstrapped, NYql::TTopicEventProcessor { private: using TBase = TActorBootstrapped; @@ -315,7 +317,7 @@ public: [[maybe_unused]] static constexpr char ActorName[] = "FQ_ROW_DISPATCHER_SESSION"; private: - NYdb::NTopic::TTopicClientSettings GetTopicClientSettings(bool useSsl) const; + NYdb::NTopic::TTopicClientSettings GetTopicClientSettings(bool useSsl); NYql::ITopicClient& GetTopicClient(bool useSsl); NYdb::NTopic::TReadSessionSettings GetReadSessionSettings(const TString& consumerName) const; void CreateTopicSession(); @@ -351,6 +353,7 @@ private: private: STRICT_STFUNC_EXC(StateFunc, + hFunc(TEvPrivate::TEvExecuteTopicEvent, HandleTopicEvent); hFunc(NFq::TEvPrivate::TEvPqEventsReady, Handle); hFunc(NFq::TEvPrivate::TEvCreateSession, Handle); hFunc(NFq::TEvPrivate::TEvSendStatistic, Handle); @@ -364,6 +367,7 @@ private: STRICT_STFUNC_EXC(ErrorState, cFunc(TEvents::TEvPoisonPill::EventType, PassAway); + hFunc(TEvPrivate::TEvExecuteTopicEvent, HandleTopicEvent); IgnoreFunc(NFq::TEvPrivate::TEvPqEventsReady); IgnoreFunc(NFq::TEvPrivate::TEvCreateSession); IgnoreFunc(TEvRowDispatcher::TEvGetNextBatch); @@ -446,12 +450,16 @@ void TTopicSession::SubscribeOnNextEvent() { }); } -NYdb::NTopic::TTopicClientSettings TTopicSession::GetTopicClientSettings(bool useSsl) const { - return PqGateway->GetTopicClientSettings() - .Database(Database) +NYdb::NTopic::TTopicClientSettings TTopicSession::GetTopicClientSettings(bool useSsl) { + auto opts = PqGateway->GetTopicClientSettings(); + SetupTopicClientSettings(ActorContext().ActorSystem(), SelfId(), opts); + + opts.Database(Database) .DiscoveryEndpoint(Endpoint) .SslCredentials(NYdb::TSslCredentials(useSsl)) .CredentialsProviderFactory(CredentialsProviderFactory); + + return opts; } NYql::ITopicClient& TTopicSession::GetTopicClient(bool useSsl) { @@ -983,7 +991,7 @@ void TTopicSession::RefreshParsers() { } } -} // anonymous namespace +} // anonymous namespace //////////////////////////////////////////////////////////////////////////////// @@ -1006,4 +1014,4 @@ std::unique_ptr NewTopicSession( return std::unique_ptr(new TTopicSession(readGroup, topicPath, endpoint, database, config, functionRegistry, rowDispatcherActorId, compileServiceActorId, partitionId, std::move(driver), credentialsProviderFactory, counters, countersRoot, pqGateway, maxBufferSize)); } -} // namespace NFq +} // namespace NFq diff --git a/ydb/core/fq/libs/row_dispatcher/ya.make b/ydb/core/fq/libs/row_dispatcher/ya.make index 73a6aad5493..f5ce290ff23 100644 --- a/ydb/core/fq/libs/row_dispatcher/ya.make +++ b/ydb/core/fq/libs/row_dispatcher/ya.make @@ -29,6 +29,7 @@ PEERDIR( ydb/library/yql/dq/actors/common ydb/library/yql/dq/actors/compute ydb/library/yql/dq/proto + ydb/library/yql/providers/pq/common ydb/library/yql/providers/pq/provider ydb/public/sdk/cpp/adapters/issue diff --git a/ydb/library/yql/providers/pq/async_io/dq_pq_read_actor.cpp b/ydb/library/yql/providers/pq/async_io/dq_pq_read_actor.cpp index cf91eef3065..c8d5c4f489a 100644 --- a/ydb/library/yql/providers/pq/async_io/dq_pq_read_actor.cpp +++ b/ydb/library/yql/providers/pq/async_io/dq_pq_read_actor.cpp @@ -1,47 +1,44 @@ #include "dq_pq_read_actor.h" +#include "dq_pq_meta_extractor.h" +#include "dq_pq_rd_read_actor.h" +#include "dq_pq_read_actor_base.h" #include "probes.h" +#include +#include +#include +#include +#include +#include +#include #include #include #include #include #include #include - -#include -#include -#include -#include -#include -#include +#include #include #include #include -#include -#include - +#include #include #include #include -#include -#include -#include -#include -#include -#include -#include - -#include +#include +#include +#include +#include +#include -#include +#include #include #include #include #include - #include #include @@ -81,6 +78,7 @@ struct TEvPrivate { EvReconnectSession, EvReceivedClusters, EvDescribeTopicResult, + EvExecuteTopicEvent, EvEnd }; @@ -116,10 +114,14 @@ struct TEvPrivate { ui32 PartitionsCount; TMaybe Status; }; + struct TEvExecuteTopicEvent : public TTopicEventBase { + using TTopicEventBase::TTopicEventBase; + }; }; -} // namespace -class TDqPqReadActor : public NActors::TActor, public NYql::NDq::NInternal::TDqPqReadActorBase { +} // anonymous namespace + +class TDqPqReadActor : public NActors::TActor, public NYql::NDq::NInternal::TDqPqReadActorBase, TTopicEventProcessor { static constexpr bool StaticDiscovery = true; struct TMetrics { TMetrics( @@ -232,8 +234,10 @@ public: return opts; } - NYdb::NTopic::TTopicClientSettings GetTopicClientSettings(TClusterState& state) const { + NYdb::NTopic::TTopicClientSettings GetTopicClientSettings(TClusterState& state) { NYdb::NTopic::TTopicClientSettings opts = PqGateway->GetTopicClientSettings(); + SetupTopicClientSettings(ActorContext().ActorSystem(), SelfId(), opts); + opts.Database(SourceParams.GetDatabase()) .DiscoveryEndpoint(SourceParams.GetEndpoint()) .SslCredentials(NYdb::TSslCredentials(SourceParams.GetUseSsl())) @@ -330,6 +334,7 @@ private: hFunc(TEvPrivate::TEvReconnectSession, Handle); hFunc(TEvPrivate::TEvReceivedClusters, Handle); hFunc(TEvPrivate::TEvDescribeTopicResult, Handle); + hFunc(TEvPrivate::TEvExecuteTopicEvent, HandleTopicEvent); ) void Handle(TEvPrivate::TEvSourceDataReady::TPtr& ev) { diff --git a/ydb/library/yql/providers/pq/async_io/dq_pq_write_actor.cpp b/ydb/library/yql/providers/pq/async_io/dq_pq_write_actor.cpp index 56c422a6f96..3590f3eeb22 100644 --- a/ydb/library/yql/providers/pq/async_io/dq_pq_write_actor.cpp +++ b/ydb/library/yql/providers/pq/async_io/dq_pq_write_actor.cpp @@ -1,27 +1,27 @@ #include "dq_pq_write_actor.h" #include "probes.h" +#include +#include +#include +#include +#include +#include #include #include #include -#include +#include +#include +#include +#include +#include -#include #include #include #include -#include +#include #include -#include -#include -#include - -#include -#include -#include -#include -#include #include #include @@ -78,6 +78,7 @@ struct TEvPrivate { EvBegin = EventSpaceBegin(NActors::TEvents::ES_PRIVATE), EvPqEventsReady = EvBegin, + EvExecuteTopicEvent, EvEnd }; @@ -87,15 +88,19 @@ struct TEvPrivate { // Events struct TEvPqEventsReady : public TEventLocal {}; + + struct TEvExecuteTopicEvent : public TTopicEventBase { + using TTopicEventBase::TTopicEventBase; + }; }; TString MakeStringForLog(const NDqProto::TCheckpoint& checkpoint) { return TStringBuilder() << "[Checkpoint " << checkpoint.GetGeneration() << "." << checkpoint.GetId() << "] "; } -} // namespace +} // anonymous namespace -class TDqPqWriteActor : public NActors::TActor, public IDqComputeActorAsyncOutput { +class TDqPqWriteActor : public NActors::TActor, public IDqComputeActorAsyncOutput, TTopicEventProcessor { struct TMetrics { TMetrics(const TTxId& txId, ui64 taskId, const ::NMonitoring::TDynamicCounterPtr& counters) : TxId(std::visit([](auto arg) { return ToString(arg); }, txId)) @@ -282,6 +287,7 @@ public: private: STRICT_STFUNC(StateFunc, hFunc(TEvPrivate::TEvPqEventsReady, Handle); + hFunc(TEvPrivate::TEvExecuteTopicEvent, HandleTopicEvent); ) void Handle(TEvPrivate::TEvPqEventsReady::TPtr&) { @@ -329,8 +335,10 @@ private: return *FederatedTopicClient; } - NYdb::NFederatedTopic::TFederatedTopicClientSettings GetFederatedTopicClientSettings() const { + NYdb::NFederatedTopic::TFederatedTopicClientSettings GetFederatedTopicClientSettings() { NYdb::NFederatedTopic::TFederatedTopicClientSettings opts = PqGateway->GetFederatedTopicClientSettings(); + SetupTopicClientSettings(ActorContext().ActorSystem(), SelfId(), opts); + opts.Database(SinkParams.GetDatabase()) .DiscoveryEndpoint(SinkParams.GetEndpoint()) .SslCredentials(NYdb::TSslCredentials(SinkParams.GetUseSsl())) diff --git a/ydb/library/yql/providers/pq/common/pq_events_processor.h b/ydb/library/yql/providers/pq/common/pq_events_processor.h new file mode 100644 index 00000000000..ee013a27a3e --- /dev/null +++ b/ydb/library/yql/providers/pq/common/pq_events_processor.h @@ -0,0 +1,73 @@ +#pragma once + +#include +#include +#include +#include + +namespace NYql { + +template +class TTopicEventBase : public NActors::TEventLocal { +public: + explicit TTopicEventBase(NYdb::NTopic::IExecutor::TFunction&& f) + : Function(std::move(f)) + {} + + void Execute() { + Function(); + } + +private: + NYdb::NTopic::IExecutor::TFunction Function; +}; + +template +class TTopicEventProcessor { + class TEventProxy final : public NYdb::NTopic::IExecutor { + public: + TEventProxy(NActors::TActorSystem* actorSystem, const NActors::TActorId& executerId) + : ActorSystem(actorSystem) + , ExecuterId(executerId) + { + Y_ENSURE(actorSystem); + } + + bool IsAsync() const final { + return true; + } + + void Post(TFunction&& f) final { + ActorSystem->Send(ExecuterId, new TTopicEvent(std::move(f))); + } + + private: + void DoStart() final { + } + + private: + NActors::TActorSystem* ActorSystem = nullptr; + const NActors::TActorId ExecuterId; + }; + +public: + template + void SetupTopicClientSettings(NActors::TActorSystem* actorSystem, const NActors::TActorId& selfId, TSettings& settings) { + if (!ExecuterProxy) { + ExecuterProxy = MakeIntrusive(actorSystem, selfId); + } + + settings.DefaultHandlersExecutor(ExecuterProxy); + settings.DefaultCompressionExecutor(ExecuterProxy); + } + +protected: + void HandleTopicEvent(TTopicEvent::TPtr& event) { + event->Get()->Execute(); + } + +private: + NYdb::NTopic::IExecutor::TPtr ExecuterProxy; +}; + +} // namespace NYql diff --git a/ydb/library/yql/providers/pq/common/ya.make b/ydb/library/yql/providers/pq/common/ya.make index e2dbc03f530..380338c2d62 100644 --- a/ydb/library/yql/providers/pq/common/ya.make +++ b/ydb/library/yql/providers/pq/common/ya.make @@ -7,6 +7,8 @@ SRCS( ) PEERDIR( + ydb/library/actors/core + ydb/public/sdk/cpp/src/client/topic yql/essentials/public/types ) -- cgit v1.3