summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
-rw-r--r--ydb/core/fq/libs/row_dispatcher/local_leader_election.cpp33
-rw-r--r--ydb/core/grpc_services/local_rpc/local_rpc.cpp39
-rw-r--r--ydb/core/grpc_services/local_rpc/local_rpc.h16
-rw-r--r--ydb/core/grpc_services/local_rpc/ya.make2
-rw-r--r--ydb/core/kqp/gateway/behaviour/streaming_query/queries.cpp17
-rw-r--r--ydb/core/kqp/run_script_actor/kqp_script_result_handler.cpp8
-rw-r--r--ydb/core/kqp/ut/federated_query/datastreams/common.cpp8
-rw-r--r--ydb/core/kqp/ut/federated_query/datastreams/streaming_ddl_ut.cpp70
-rw-r--r--ydb/core/kqp/ut/federated_query/datastreams/streaming_sys_view_ut.cpp3
-rw-r--r--ydb/core/local_proxy/local_pq_client/local_topic_io_session_common.h33
-rw-r--r--ydb/library/query_actor/query_actor_ut.cpp23
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