diff options
| author | Pisarenko Grigoriy <[email protected]> | 2026-07-16 15:31:53 +0500 |
|---|---|---|
| committer | GitHub <[email protected]> | 2026-07-16 13:31:53 +0300 |
| commit | 959f00a059c4deb0f6fcbbb5c77cc2dd283b57f8 (patch) | |
| tree | 22455289eadc7f6faa8e3b6b4aa6a8ff6fbcee83 | |
| parent | 69dac9df94547dc0c4716c525fcabfaa19eebaf0 (diff) | |
YQ-5483 fix script executions leak form streaming queries (#46033)
11 files changed, 174 insertions, 78 deletions
diff --git a/ydb/core/fq/libs/row_dispatcher/local_leader_election.cpp b/ydb/core/fq/libs/row_dispatcher/local_leader_election.cpp index 7a691f86f4b..432b2b5508c 100644 --- a/ydb/core/fq/libs/row_dispatcher/local_leader_election.cpp +++ b/ydb/core/fq/libs/row_dispatcher/local_leader_election.cpp @@ -648,38 +648,7 @@ void TLocalLeaderElection::CloseSession(const grpc::Status& status, const TStrin issues = AddRootIssue(message, issues); } - switch (status.error_code()) { - case grpc::OK: - return CloseSession(NYdb::EStatus::SUCCESS, issues); - case grpc::CANCELLED: - return CloseSession(NYdb::EStatus::CANCELLED, issues); - case grpc::UNKNOWN: - return CloseSession(NYdb::EStatus::UNDETERMINED, issues); - case grpc::DEADLINE_EXCEEDED: - return CloseSession(NYdb::EStatus::TIMEOUT, issues); - case grpc::NOT_FOUND: - return CloseSession(NYdb::EStatus::NOT_FOUND, issues); - case grpc::ALREADY_EXISTS: - return CloseSession(NYdb::EStatus::ALREADY_EXISTS, issues); - case grpc::RESOURCE_EXHAUSTED: - return CloseSession(NYdb::EStatus::OVERLOADED, issues); - case grpc::ABORTED: - return CloseSession(NYdb::EStatus::ABORTED, issues); - case grpc::OUT_OF_RANGE: - return CloseSession(NYdb::EStatus::CLIENT_OUT_OF_RANGE, issues); - case grpc::UNIMPLEMENTED: - return CloseSession(NYdb::EStatus::UNSUPPORTED, issues); - case grpc::UNAVAILABLE: - return CloseSession(NYdb::EStatus::UNAVAILABLE, issues); - case grpc::INVALID_ARGUMENT: - case grpc::FAILED_PRECONDITION: - return CloseSession(NYdb::EStatus::PRECONDITION_FAILED, issues); - case grpc::UNAUTHENTICATED: - case grpc::PERMISSION_DENIED: - return CloseSession(NYdb::EStatus::UNAUTHORIZED, issues); - default: - return CloseSession(NYdb::EStatus::INTERNAL_ERROR, issues); - } + CloseSession(static_cast<NYdb::EStatus>(NKikimr::NRpcService::GrpcStatusToYdbStatus(status.error_code())), issues); } void TLocalLeaderElection::ProcessPing(const TRpcOut& message) { diff --git a/ydb/core/grpc_services/local_rpc/local_rpc.cpp b/ydb/core/grpc_services/local_rpc/local_rpc.cpp new file mode 100644 index 00000000000..db3bb69ebbe --- /dev/null +++ b/ydb/core/grpc_services/local_rpc/local_rpc.cpp @@ -0,0 +1,39 @@ +#include "local_rpc.h" + +namespace NKikimr::NRpcService { + +Ydb::StatusIds::StatusCode GrpcStatusToYdbStatus(grpc::StatusCode status) { + switch (status) { + case grpc::OK: + return Ydb::StatusIds::SUCCESS; + case grpc::CANCELLED: + return Ydb::StatusIds::CANCELLED; + case grpc::UNKNOWN: + return Ydb::StatusIds::UNDETERMINED; + case grpc::DEADLINE_EXCEEDED: + return Ydb::StatusIds::TIMEOUT; + case grpc::NOT_FOUND: + return Ydb::StatusIds::NOT_FOUND; + case grpc::ALREADY_EXISTS: + return Ydb::StatusIds::ALREADY_EXISTS; + case grpc::RESOURCE_EXHAUSTED: + return Ydb::StatusIds::OVERLOADED; + case grpc::ABORTED: + return Ydb::StatusIds::ABORTED; + case grpc::UNIMPLEMENTED: + return Ydb::StatusIds::UNSUPPORTED; + case grpc::UNAVAILABLE: + return Ydb::StatusIds::UNAVAILABLE; + case grpc::OUT_OF_RANGE: + case grpc::INVALID_ARGUMENT: + case grpc::FAILED_PRECONDITION: + return Ydb::StatusIds::PRECONDITION_FAILED; + case grpc::UNAUTHENTICATED: + case grpc::PERMISSION_DENIED: + return Ydb::StatusIds::UNAUTHORIZED; + default: + return Ydb::StatusIds::INTERNAL_ERROR; + } +} + +} // namespace NKikimr::NRpcService diff --git a/ydb/core/grpc_services/local_rpc/local_rpc.h b/ydb/core/grpc_services/local_rpc/local_rpc.h index d0533eeb9ab..e4a92bf7fd1 100644 --- a/ydb/core/grpc_services/local_rpc/local_rpc.h +++ b/ydb/core/grpc_services/local_rpc/local_rpc.h @@ -11,6 +11,8 @@ namespace NKikimr::NRpcService { +Ydb::StatusIds::StatusCode GrpcStatusToYdbStatus(grpc::StatusCode status); + template<typename TResponse> class TPromiseWrapper { public: @@ -261,11 +263,19 @@ public: void SetRuHeader(ui64) override { } - // Unimplemented methods - void ReplyWithRpcStatus(grpc::StatusCode, const TString&, const TString&) override { - ReplyWithYdbStatus(Ydb::StatusIds::GENERIC_ERROR); + void ReplyWithRpcStatus(grpc::StatusCode status, const TString& reason, const TString& details) override { + if (reason) { + TBase::IssueManager.RaiseIssue(NYql::TIssue(reason)); + } + + if (details) { + TBase::IssueManager.RaiseIssue(NYql::TIssue(TStringBuilder() << "gRPC Details: " << details)); + } + + ReplyWithYdbStatus(GrpcStatusToYdbStatus(status)); } + // Unimplemented methods void SetStreamingNotify(NYdbGrpc::IRequestContextBase::TOnNextReply&&) override { Y_ABORT("Unimplemented for local rpc"); } diff --git a/ydb/core/grpc_services/local_rpc/ya.make b/ydb/core/grpc_services/local_rpc/ya.make index 8480c08f263..b530d9dc93d 100644 --- a/ydb/core/grpc_services/local_rpc/ya.make +++ b/ydb/core/grpc_services/local_rpc/ya.make @@ -1,7 +1,7 @@ LIBRARY() SRCS( - local_rpc.h + local_rpc.cpp ) PEERDIR( diff --git a/ydb/core/kqp/gateway/behaviour/streaming_query/queries.cpp b/ydb/core/kqp/gateway/behaviour/streaming_query/queries.cpp index b36b0064183..bb3275c20fe 100644 --- a/ydb/core/kqp/gateway/behaviour/streaming_query/queries.cpp +++ b/ydb/core/kqp/gateway/behaviour/streaming_query/queries.cpp @@ -1986,11 +1986,12 @@ public: ) void HandleRemove(TEvPrivate::TEvCleanupStreamingQueryResult::TPtr& ev) { + State = ev->Get()->Info; + if (HandleResult(ev, "Cleanup streaming query")) { return; } - State = ev->Get()->Info; RemoveQuery(); } @@ -2043,11 +2044,12 @@ public: } void Handle(TEvPrivate::TEvStartStreamingQueryResult::TPtr& ev) { + State = ev->Get()->Info; + if (HandleResult(ev, "Start streaming query")) { return; } - State = ev->Get()->Info; Finish(Ydb::StatusIds::SUCCESS); } @@ -2444,11 +2446,12 @@ public: } void HandleSync(TEvPrivate::TEvSyncStreamingQueryResult::TPtr& ev) { + TBase::QueryState = ev->Get()->State; + if (TBase::HandleResult(ev, "Streaming query initialization (recover previous query state, try to repeat request)")) { return; } - TBase::QueryState = ev->Get()->State; if (!ev->Get()->ExistsInSS) { TBase::SchemeInfo = std::nullopt; } @@ -2519,11 +2522,12 @@ public: } void Handle(TEvPrivate::TEvSyncStreamingQueryResult::TPtr& ev) { + QueryState = ev->Get()->State; + if (HandleResult(ev, "Streaming query initialization")) { return; } - QueryState = ev->Get()->State; if (!ev->Get()->ExistsInSS) { SchemeInfo = std::nullopt; } @@ -2662,11 +2666,12 @@ public: } void Handle(TEvPrivate::TEvSyncStreamingQueryResult::TPtr& ev) { + QueryState = ev->Get()->State; + if (HandleResult(ev, "Streaming query alter")) { return; } - QueryState = ev->Get()->State; if (!ev->Get()->ExistsInSS) { SchemeInfo = std::nullopt; } @@ -2773,6 +2778,8 @@ public: } void Handle(TEvPrivate::TEvCleanupStreamingQueryResult::TPtr& ev) { + QueryState = ev->Get()->Info; + if (HandleResult(ev, "Cleanup streaming query")) { return; } diff --git a/ydb/core/kqp/run_script_actor/kqp_script_result_handler.cpp b/ydb/core/kqp/run_script_actor/kqp_script_result_handler.cpp index cdd42f12d28..877774b5448 100644 --- a/ydb/core/kqp/run_script_actor/kqp_script_result_handler.cpp +++ b/ydb/core/kqp/run_script_actor/kqp_script_result_handler.cpp @@ -55,6 +55,7 @@ class TScriptResultHandlerActor final : public TActorBootstrapped<TScriptResultH bool WaitSave = false; bool QueryStatsChanged = true; // We should update status once execution was started bool AstSaved = false; + TInstant SuspendUntil; void UpdateChanged(const std::optional<TString>& from, const TString& to) { if (QueryStatsChanged) { @@ -234,6 +235,7 @@ private: hFunc(TEvKqp::TEvQueryResponse, Handle); hFunc(TEvKqp::TEvCancelQueryResponse, Handle); sFunc(TEvents::TEvPoison, Finish); + sFunc(TEvents::TEvWakeup, ContinueExecute); ) void Handle(TEvSaveScriptExternalEffectRequest::TPtr& ev) { @@ -353,8 +355,10 @@ private: const auto astSaved = ev->Get()->AstSaved; if (const auto status = ev->Get()->Status; status != Ydb::StatusIds::SUCCESS) { - LOG_N("Script progress updated " << ev->Sender << ", fail: " << status << ", issues: " << ev->Get()->Issues.ToOneLineString()); SaveProgressState.QueryStatsChanged = true; + SaveProgressState.SuspendUntil = TInstant::Now() + TDuration::Seconds(1); + Schedule(SaveProgressState.SuspendUntil, new TEvents::TEvWakeup()); + LOG_N("Script progress updated " << ev->Sender << ", fail: " << status << ", suspend until: " << SaveProgressState.SuspendUntil << ", issues: " << ev->Get()->Issues.ToOneLineString()); } else { LOG_T("Script progress updated " << ev->Sender << ", ast saved: " << astSaved); SaveProgressState.AstSaved = SaveProgressState.AstSaved || astSaved; @@ -594,7 +598,7 @@ private: return SaveResultsMeta(); } - if (SaveProgressState.QueryStatsChanged) { + if (SaveProgressState.QueryStatsChanged && SaveProgressState.SuspendUntil <= TInstant::Now()) { return UpdateScriptProgress(); } diff --git a/ydb/core/kqp/ut/federated_query/datastreams/common.cpp b/ydb/core/kqp/ut/federated_query/datastreams/common.cpp index 75d8ce31d6d..f429c895ad6 100644 --- a/ydb/core/kqp/ut/federated_query/datastreams/common.cpp +++ b/ydb/core/kqp/ut/federated_query/datastreams/common.cpp @@ -902,7 +902,11 @@ std::vector<TStreamingSysViewTestFixture::TSysViewResult> TStreamingSysViewTestF const bool expectExecutions = row.Run && IsIn({"RUNNING", "COMPLETED", "CANCELLED", "FAILED"}, row.Status); if (expectExecutions || row.CheckPlan) { - UNIT_ASSERT_STRING_CONTAINS(*resultSet.ColumnParser("Plan").GetOptionalUtf8(), TStringBuilder() << "Write " << PQ_SOURCE); + if (row.CheckPlan || IsIn({"RUNNING", "COMPLETED", "CANCELLED"}, row.Status)) { + UNIT_ASSERT_STRING_CONTAINS(*resultSet.ColumnParser("Plan").GetOptionalUtf8(), TStringBuilder() << "Write " << PQ_SOURCE); + } else { + UNIT_ASSERT(resultSet.ColumnParser("Plan").GetOptionalUtf8()); + } UNIT_ASSERT_STRING_CONTAINS(*resultSet.ColumnParser("Ast").GetOptionalUtf8(), row.Ast ? *row.Ast : JoinPath({"/Root", PQ_SOURCE})); } @@ -926,7 +930,7 @@ std::vector<TStreamingSysViewTestFixture::TSysViewResult> TStreamingSysViewTestF } result.ExecutionId = *resultSet.ColumnParser("LastExecutionId").GetOptionalUtf8(); - UNIT_ASSERT_VALUES_EQUAL(!result.ExecutionId.empty(), expectExecutions); + UNIT_ASSERT_VALUES_EQUAL(!result.ExecutionId.empty(), expectExecutions && IsIn({"RUNNING", "COMPLETED", "CANCELLED"}, row.Status)); const auto previousExecutionIds = *resultSet.ColumnParser("PreviousExecutionIds").GetOptionalUtf8(); NJson::TJsonValue value; diff --git a/ydb/core/kqp/ut/federated_query/datastreams/streaming_ddl_ut.cpp b/ydb/core/kqp/ut/federated_query/datastreams/streaming_ddl_ut.cpp index 49a96f84fed..668e83da734 100644 --- a/ydb/core/kqp/ut/federated_query/datastreams/streaming_ddl_ut.cpp +++ b/ydb/core/kqp/ut/federated_query/datastreams/streaming_ddl_ut.cpp @@ -3567,6 +3567,76 @@ Y_UNIT_TEST_SUITE(KqpStreamingQueriesDdl) { EStatus::GENERIC_ERROR, "names starting with '__ydb_' are reserved for system columns"); } + + Y_UNIT_TEST_F(StreamingQueryInvalidationAfterCreation, TStreamingTestFixture) { + ExecQuery("GRANT ALL ON `/Root` TO `" BUILTIN_ACL_ROOT "`"); + + constexpr char inputTopicName[] = "streamingQueryInvalidationAfterCreationInputTopic1"; + constexpr char outputTopicName[] = "streamingQueryInvalidationAfterCreationOutputTopic"; + CreateTopic(inputTopicName); + CreateTopic(outputTopicName); + + constexpr char pqSourceName[] = "sourceName"; + CreatePqSource(pqSourceName); + + constexpr char queryName[] = "streamingQuery"; + ExecQuery(fmt::format( + R"sql( + CREATE STREAMING QUERY `{query_name}` AS + DO BEGIN + INSERT INTO `{pq_source}`.`{output_topic}` + SELECT * FROM `{pq_source}`.`{input_topic}`; + END DO; + )sql", + "query_name"_a = queryName, + "pq_source"_a = pqSourceName, + "input_topic"_a = inputTopicName, + "output_topic"_a = outputTopicName + )); + + CheckScriptExecutionsCount(1, 1); + Sleep(TDuration::Seconds(1)); + + WriteTopicMessage(inputTopicName, "test_message"); + ReadTopicMessage(outputTopicName, "test_message"); + + constexpr ui64 changesCount = 10; + for (ui64 i = 0; i < changesCount; ++i) { + ExecQuery(fmt::format( + R"sql( + ALTER STREAMING QUERY `{query_name}` SET (FORCE = TRUE) AS + DO BEGIN + PRAGMA ydb.OverridePlanner = "invalid"; + INSERT INTO `{pq_source}`.`{output_topic}` + SELECT * FROM `{pq_source}`.`{input_topic}`; + END DO; + )sql", + "query_name"_a = queryName, + "pq_source"_a = pqSourceName, + "input_topic"_a = inputTopicName, + "output_topic"_a = outputTopicName + ), EStatus::GENERIC_ERROR, "Invalid override planner settings"); + } + + { + const auto& result = ExecQuery("SELECT Status, Issues FROM `.sys/streaming_queries`"); + UNIT_ASSERT_VALUES_EQUAL(result.size(), 1); + + CheckScriptResult(result[0], 2, 1, [&](TResultSetParser& resultSet) { + UNIT_ASSERT_STRING_CONTAINS(resultSet.ColumnParser("Issues").GetOptionalUtf8().value_or(""), "Invalid override planner settings"); + UNIT_ASSERT_VALUES_EQUAL(resultSet.ColumnParser("Status").GetOptionalUtf8().value_or(""), "FAILED"); + }); + } + + { + const auto& result = ExecQuery("SELECT COUNT(*) AS count FROM `.metadata/script_executions`"); + UNIT_ASSERT_VALUES_EQUAL(result.size(), 1); + + CheckScriptResult(result[0], 1, 1, [&](TResultSetParser& resultSet) { + UNIT_ASSERT_VALUES_EQUAL(resultSet.ColumnParser("count").GetUint64(), std::min(changesCount, static_cast<ui64>(4))); + }); + } + } } } // namespace NKikimr::NKqp diff --git a/ydb/core/kqp/ut/federated_query/datastreams/streaming_sys_view_ut.cpp b/ydb/core/kqp/ut/federated_query/datastreams/streaming_sys_view_ut.cpp index ea696e0077b..f7db2cac487 100644 --- a/ydb/core/kqp/ut/federated_query/datastreams/streaming_sys_view_ut.cpp +++ b/ydb/core/kqp/ut/federated_query/datastreams/streaming_sys_view_ut.cpp @@ -249,7 +249,8 @@ Y_UNIT_TEST_SUITE(KqpStreamingQueriesSysView) { .CheckPlan = true, }, { .Name = "B", - .Status = "CREATED", + .Status = "FAILED", + .Issues = "Invalid override planner settings", .Text = texts[1], }}); } diff --git a/ydb/core/local_proxy/local_pq_client/local_topic_io_session_common.h b/ydb/core/local_proxy/local_pq_client/local_topic_io_session_common.h index 64136482cbf..d512e51696d 100644 --- a/ydb/core/local_proxy/local_pq_client/local_topic_io_session_common.h +++ b/ydb/core/local_proxy/local_pq_client/local_topic_io_session_common.h @@ -204,38 +204,7 @@ protected: issues = AddRootIssue(message, issues); } - switch (status.error_code()) { - case grpc::OK: - return CloseSession(NYdb::EStatus::SUCCESS, issues); - case grpc::CANCELLED: - return CloseSession(NYdb::EStatus::CANCELLED, issues); - case grpc::UNKNOWN: - return CloseSession(NYdb::EStatus::UNDETERMINED, issues); - case grpc::DEADLINE_EXCEEDED: - return CloseSession(NYdb::EStatus::TIMEOUT, issues); - case grpc::NOT_FOUND: - return CloseSession(NYdb::EStatus::NOT_FOUND, issues); - case grpc::ALREADY_EXISTS: - return CloseSession(NYdb::EStatus::ALREADY_EXISTS, issues); - case grpc::RESOURCE_EXHAUSTED: - return CloseSession(NYdb::EStatus::OVERLOADED, issues); - case grpc::ABORTED: - return CloseSession(NYdb::EStatus::ABORTED, issues); - case grpc::OUT_OF_RANGE: - return CloseSession(NYdb::EStatus::CLIENT_OUT_OF_RANGE, issues); - case grpc::UNIMPLEMENTED: - return CloseSession(NYdb::EStatus::UNSUPPORTED, issues); - case grpc::UNAVAILABLE: - return CloseSession(NYdb::EStatus::UNAVAILABLE, issues); - case grpc::INVALID_ARGUMENT: - case grpc::FAILED_PRECONDITION: - return CloseSession(NYdb::EStatus::PRECONDITION_FAILED, issues); - case grpc::UNAUTHENTICATED: - case grpc::PERMISSION_DENIED: - return CloseSession(NYdb::EStatus::UNAUTHORIZED, issues); - default: - return CloseSession(NYdb::EStatus::INTERNAL_ERROR, issues); - } + CloseSession(static_cast<NYdb::EStatus>(NRpcService::GrpcStatusToYdbStatus(status.error_code())), issues); } void HandleWakeup() { diff --git a/ydb/library/query_actor/query_actor_ut.cpp b/ydb/library/query_actor/query_actor_ut.cpp index b9ce900d1a1..e99a25b77b3 100644 --- a/ydb/library/query_actor/query_actor_ut.cpp +++ b/ydb/library/query_actor/query_actor_ut.cpp @@ -274,6 +274,29 @@ Y_UNIT_TEST_SUITE(QueryActorTest) { UNIT_ASSERT_VALUES_EQUAL(result.ResultSets[0].RowsCount(), 1000); } } + + Y_UNIT_TEST(StartQueryDuringShutdown) { + TTestServer server; + + auto& runtime = *server.Server->GetRuntime(); + runtime.Send(NKqp::MakeKqpProxyID(runtime.GetFirstNodeId()), {}, new NKqp::TEvKqp::TEvInitiateShutdownRequest( + MakeIntrusive<NKqp::TKqpShutdownState>() + )); + + struct TQuery : public TTestQueryActorBase { + void OnRunQuery() override { + RunDataQuery("SELECT 42"); + } + + void OnQueryResult() override { + Finish(); + } + }; + + auto result = server.RunQueryActor<TQuery>(); + UNIT_ASSERT_VALUES_EQUAL_C(result.StatusCode, Ydb::StatusIds::OVERLOADED, result.Issues.ToOneLineString()); + UNIT_ASSERT_STRING_CONTAINS(result.Issues.ToOneLineString(), "system shutdown requested"); + } } } // namespace NKikimr |
