diff options
| author | ivanmorozov333 <[email protected]> | 2025-04-12 21:20:28 +0300 |
|---|---|---|
| committer | GitHub <[email protected]> | 2025-04-12 21:20:28 +0300 |
| commit | ff4f214225d313fd54806ef83fc9c4016df7ebc8 (patch) | |
| tree | 1411743135bb744b3857267356f6031866b2b6bf | |
| parent | 9eee3dd53ab019dde1cc33f1e0697c06d59137f1 (diff) | |
fix errors processing for json (#17139)
21 files changed, 205 insertions, 71 deletions
diff --git a/ydb/core/formats/arrow/program/graph_execute.cpp b/ydb/core/formats/arrow/program/graph_execute.cpp index a8e2f3f5c6b..6d1a3ebf6be 100644 --- a/ydb/core/formats/arrow/program/graph_execute.cpp +++ b/ydb/core/formats/arrow/program/graph_execute.cpp @@ -3,6 +3,8 @@ #include "graph_optimization.h" #include "visitor.h" +#include <yql/essentials/minikql/mkql_terminator.h> + namespace NKikimr::NArrow::NSSA::NGraph::NExecution { class TResourceUsageInfo { @@ -152,12 +154,13 @@ TCompiledGraph::TCompiledGraph(const NOptimization::TGraph& original, const ICol } } AFL_TRACE(NKikimrServices::SSA_GRAPH_EXECUTION)("graph_constructed", DebugDOT()); -// Cerr << DebugDOT() << Endl; + // Cerr << DebugDOT() << Endl; } TConclusionStatus TCompiledGraph::Apply( const std::shared_ptr<IDataSource>& source, const std::shared_ptr<TAccessorsCollection>& resources) const { TProcessorContext context(source, resources, std::nullopt, false); + NMiniKQL::TThrowingBindTerminator bind; std::shared_ptr<TExecutionVisitor> visitor = std::make_shared<TExecutionVisitor>(context); for (auto it = BuildIterator(visitor); it->IsValid();) { { diff --git a/ydb/core/kqp/ut/olap/json_ut.cpp b/ydb/core/kqp/ut/olap/json_ut.cpp index b4349f04138..cca84a73971 100644 --- a/ydb/core/kqp/ut/olap/json_ut.cpp +++ b/ydb/core/kqp/ut/olap/json_ut.cpp @@ -5,14 +5,19 @@ #include "helpers/writer.h" #include <ydb/core/base/tablet_pipecache.h> +#include <ydb/core/formats/arrow/serializer/native.h> #include <ydb/core/kqp/ut/common/columnshard.h> #include <ydb/core/tx/columnshard/counters/common/object_counter.h> #include <ydb/core/tx/columnshard/engines/reader/common_reader/iterator/source.h> #include <ydb/core/tx/columnshard/hooks/testing/controller.h> +#include <ydb/core/tx/columnshard/test_helper/columnshard_ut_common.h> #include <ydb/core/tx/columnshard/test_helper/controllers.h> #include <ydb/core/tx/limiter/grouped_memory/service/process.h> #include <ydb/core/wrappers/fake_storage.h> +#include <ydb/public/lib/scheme_types/scheme_type_id.h> + +#include <library/cpp/string_utils/base64/base64.h> #include <library/cpp/testing/unittest/registar.h> #include <util/string/strip.h> @@ -243,6 +248,47 @@ Y_UNIT_TEST_SUITE(KqpOlapJson) { } }; + class TBulkUpsertCommand: public ICommand { + private: + TString TableName; + TString ArrowBatch; + Ydb::StatusIds_StatusCode ExpectedCode = Ydb::StatusIds::SUCCESS; + + public: + TBulkUpsertCommand() = default; + + virtual TConclusionStatus DoExecute(TKikimrRunner& kikimr) override { + TLocalHelper lHelper(kikimr); + lHelper.SendDataViaActorSystem(TableName, + NArrow::TStatusValidator::GetValid(NArrow::NSerialization::TNativeSerializer().Deserialize(ArrowBatch)), ExpectedCode); + return TConclusionStatus::Success(); + } + + bool DeserializeFromString(const TString& info) { + auto lines = StringSplitter(info).SplitBySet("\n").SkipEmpty().ToList<TString>(); + if (lines.size() < 2 || lines.size() > 3) { + return false; + } + TableName = Strip(lines[0]); + ArrowBatch = Base64Decode(Strip(lines[1])); + AFL_VERIFY(!!ArrowBatch); + if (lines.size() == 3) { + if (!Ydb::StatusIds_StatusCode_Parse(Strip(lines[2]), &ExpectedCode)) { + return false; + } +// if (lines[2] == "SUCCESS") { +// } else if (lines[2] = "INTERNAL_ERROR") { +// ExpectedCode = Ydb::StatusIds::INTERNAL_ERROR; +// } else if (lines[2] == "BAD_REQUEST") { +// ExpectedCode = Ydb::StatusIds::BAD_REQUEST; +// } else { +// return false; +// } + } + return true; + } + }; + class TScriptExecutor { private: std::vector<std::shared_ptr<ICommand>> Commands; @@ -275,7 +321,12 @@ Y_UNIT_TEST_SUITE(KqpOlapJson) { private: std::vector<TScriptExecutor> Scripts; std::shared_ptr<ICommand> BuildCommand(TString command) { - if (command.StartsWith("SCHEMA:")) { + if (command.StartsWith("BULK_UPSERT:")) { + command = command.substr(12); + auto result = std::make_shared<TBulkUpsertCommand>(); + AFL_VERIFY(result->DeserializeFromString(command)); + return result; + } else if (command.StartsWith("SCHEMA:")) { command = command.substr(7); return std::make_shared<TSchemaCommand>(command); } else if (command.StartsWith("DATA:")) { @@ -559,6 +610,39 @@ Y_UNIT_TEST_SUITE(KqpOlapJson) { TScriptVariator(script).Execute(); } + Y_UNIT_TEST(BrokenJsonWriting) { + NColumnShard::TTableUpdatesBuilder updates(NArrow::MakeArrowSchema( + { { "Col1", NScheme::TTypeInfo(NScheme::NTypeIds::Uint64) }, { "Col2", NScheme::TTypeInfo(NScheme::NTypeIds::Utf8) } })); + updates.AddRow().Add<int64_t>(1).Add("{\"a\" : \"c}"); + auto arrowString = Base64Encode(NArrow::NSerialization::TNativeSerializer().SerializeFull(updates.BuildArrow())); + + TString script = Sprintf(R"( + SCHEMA: + CREATE TABLE `/Root/ColumnTable` ( + Col1 Uint64 NOT NULL, + Col2 JsonDocument, + PRIMARY KEY (Col1) + ) + PARTITION BY HASH(Col1) + WITH (STORE = COLUMN, AUTO_PARTITIONING_MIN_PARTITIONS_COUNT = $$1|2|10$$); + ------ + SCHEMA: + ALTER OBJECT `/Root/ColumnTable` (TYPE TABLE) SET (ACTION=UPSERT_OPTIONS, `SCAN_READER_POLICY_NAME`=`SIMPLE`) + ------ + SCHEMA: + ALTER OBJECT `/Root/ColumnTable` (TYPE TABLE) SET (ACTION=ALTER_COLUMN, NAME=Col2, `DATA_EXTRACTOR_CLASS_NAME`=`JSON_SCANNER`, `SCAN_FIRST_LEVEL_ONLY`=`false`, + `DATA_ACCESSOR_CONSTRUCTOR.CLASS_NAME`=`SUB_COLUMNS`, `FORCE_SIMD_PARSING`=`$$true|false$$`, `COLUMNS_LIMIT`=`$$1024|0|1$$`, + `SPARSED_DETECTOR_KFF`=`$$0|10|1000$$`, `MEM_LIMIT_CHUNK`=`$$0|100|1000000$$`, `OTHERS_ALLOWED_FRACTION`=`$$0|0.5$$`) + ------ + BULK_UPSERT: + /Root/ColumnTable + %s + BAD_REQUEST + )", + arrowString.data()); + TScriptVariator(script).Execute(); + } + Y_UNIT_TEST(RestoreJsonArrayVariants) { TString script = R"( SCHEMA: diff --git a/ydb/core/tx/columnshard/blobs_action/transaction/tx_blobs_written.cpp b/ydb/core/tx/columnshard/blobs_action/transaction/tx_blobs_written.cpp index cd7feae239d..e0c5ab30a6b 100644 --- a/ydb/core/tx/columnshard/blobs_action/transaction/tx_blobs_written.cpp +++ b/ydb/core/tx/columnshard/blobs_action/transaction/tx_blobs_written.cpp @@ -1,7 +1,7 @@ -#include <ydb/core/tx/columnshard/common/path_id.h> #include "tx_blobs_written.h" #include <ydb/core/tx/columnshard/blob_cache.h> +#include <ydb/core/tx/columnshard/common/path_id.h> #include <ydb/core/tx/columnshard/engines/column_engine_logs.h> #include <ydb/core/tx/columnshard/engines/insert_table/user_data.h> #include <ydb/core/tx/columnshard/transactions/locks/write.h> @@ -154,7 +154,9 @@ bool TTxBlobsWritingFailed::DoExecute(TTransactionContext& txc, const TActorCont Self->OperationsManager->AbortTransactionOnExecute(*Self, op->GetLockId(), txc); auto ev = NEvents::TDataEvents::TEvWriteResult::BuildError(Self->TabletID(), op->GetLockId(), - NKikimrDataEvents::TEvWriteResult::STATUS_INTERNAL_ERROR, "cannot write blob: " + ::ToString(PutBlobResult)); + wResult.IsInternalError() ? NKikimrDataEvents::TEvWriteResult::STATUS_INTERNAL_ERROR + : NKikimrDataEvents::TEvWriteResult::STATUS_BAD_REQUEST, + wResult.GetErrorMessage()); Results.emplace_back(std::move(ev), writeMeta.GetSource(), op->GetCookie()); } return true; diff --git a/ydb/core/tx/columnshard/blobs_action/transaction/tx_blobs_written.h b/ydb/core/tx/columnshard/blobs_action/transaction/tx_blobs_written.h index 6bb3549b20d..738ca1e0167 100644 --- a/ydb/core/tx/columnshard/blobs_action/transaction/tx_blobs_written.h +++ b/ydb/core/tx/columnshard/blobs_action/transaction/tx_blobs_written.h @@ -55,7 +55,6 @@ public: class TTxBlobsWritingFailed: public TExtendedTransactionBase { private: using TBase = TExtendedTransactionBase; - const NKikimrProto::EReplyStatus PutBlobResult; TInsertedPortions Pack; class TReplyInfo { @@ -79,9 +78,8 @@ private: std::vector<TReplyInfo> Results; public: - TTxBlobsWritingFailed(TColumnShard* self, const NKikimrProto::EReplyStatus writeStatus, TInsertedPortions&& pack) + TTxBlobsWritingFailed(TColumnShard* self, TInsertedPortions&& pack) : TBase(self) - , PutBlobResult(writeStatus) , Pack(std::move(pack)) { } diff --git a/ydb/core/tx/columnshard/columnshard__write.cpp b/ydb/core/tx/columnshard/columnshard__write.cpp index d3387e831f8..7e407387d86 100644 --- a/ydb/core/tx/columnshard/columnshard__write.cpp +++ b/ydb/core/tx/columnshard/columnshard__write.cpp @@ -121,12 +121,12 @@ void TColumnShard::Handle(NPrivateEvents::NWrite::TEvWritePortionResult::TPtr& e Counters.OnWritePutBlobsFailed(now - i.GetWriteMeta().GetWriteStartInstant(), i.GetRecordsCount()); Counters.GetCSCounters().OnWritePutBlobsFail(now - i.GetWriteMeta().GetWriteStartInstant()); AFL_WARN(NKikimrServices::TX_COLUMNSHARD_WRITE)("writing_size", i.GetDataSize())("event", "data_write_error")( - "writing_id", i.GetWriteMeta().GetId()); + "writing_id", i.GetWriteMeta().GetId())("reason", i.GetErrorMessage()); Counters.GetWritesMonitor()->OnFinishWrite(i.GetDataSize(), 1); i.MutableWriteMeta().OnStage(NEvWrite::EWriteStage::Finished); } - Execute(new TTxBlobsWritingFailed(this, ev->Get()->GetWriteStatus(), std::move(writtenData)), ctx); + Execute(new TTxBlobsWritingFailed(this, std::move(writtenData)), ctx); } } diff --git a/ydb/core/tx/columnshard/columnshard_impl.cpp b/ydb/core/tx/columnshard/columnshard_impl.cpp index ebaf53c86a3..3f089c2b4d0 100644 --- a/ydb/core/tx/columnshard/columnshard_impl.cpp +++ b/ydb/core/tx/columnshard/columnshard_impl.cpp @@ -544,7 +544,7 @@ private: NOlap::TSnapshot LastCompletedTx; protected: - virtual TConclusionStatus DoExecute(const std::shared_ptr<NConveyor::ITask>& /*taskPtr*/) override { + virtual void DoExecute(const std::shared_ptr<NConveyor::ITask>& /*taskPtr*/) override { NActors::TLogContextGuard g(NActors::TLogContextBuilder::Build(NKikimrServices::TX_COLUMNSHARD)("tablet_id", TabletId)("parent_id", ParentActorId)); { NOlap::TConstructionContext context(*TxEvent->IndexInfo, Counters, LastCompletedTx); @@ -554,7 +554,6 @@ protected: } } TActorContext::AsActorContext().Send(ParentActorId, std::move(TxEvent)); - return TConclusionStatus::Success(); } public: @@ -1394,14 +1393,13 @@ private: std::shared_ptr<NOlap::NDataAccessorControl::IAccessorCallback> FetchCallback; std::vector<TPortionConstructorV2> Portions; - virtual TConclusionStatus DoExecute(const std::shared_ptr<ITask>& /*taskPtr*/) override { + virtual void DoExecute(const std::shared_ptr<ITask>& /*taskPtr*/) override { std::vector<NOlap::TPortionDataAccessor> accessors; accessors.reserve(Portions.size()); for (auto&& i : Portions) { accessors.emplace_back(i.BuildAccessor()); } FetchCallback->OnAccessorsFetched(std::move(accessors)); - return TConclusionStatus::Success(); } virtual void DoOnCannotExecute(const TString& reason) override { AFL_VERIFY(false)("cannot parse metadata", reason); diff --git a/ydb/core/tx/columnshard/engines/reader/common/conveyor_task.cpp b/ydb/core/tx/columnshard/engines/reader/common/conveyor_task.cpp index ac778b00a6c..1bcd853bec3 100644 --- a/ydb/core/tx/columnshard/engines/reader/common/conveyor_task.cpp +++ b/ydb/core/tx/columnshard/engines/reader/common/conveyor_task.cpp @@ -4,7 +4,7 @@ namespace NKikimr::NOlap::NReader { -NKikimr::TConclusionStatus IDataTasksProcessor::ITask::DoExecute(const std::shared_ptr<NConveyor::ITask>& taskPtr) { +void IDataTasksProcessor::ITask::DoExecute(const std::shared_ptr<NConveyor::ITask>& taskPtr) { auto result = DoExecuteImpl(); if (result.IsFail()) { NActors::TActivationContext::AsActorContext().Send(OwnerId, new NColumnShard::TEvPrivate::TEvTaskProcessedResult(result)); @@ -12,7 +12,6 @@ NKikimr::TConclusionStatus IDataTasksProcessor::ITask::DoExecute(const std::shar NActors::TActivationContext::AsActorContext().Send( OwnerId, new NColumnShard::TEvPrivate::TEvTaskProcessedResult(static_pointer_cast<IDataTasksProcessor::ITask>(taskPtr))); } - return result; } void IDataTasksProcessor::ITask::DoOnCannotExecute(const TString& reason) { diff --git a/ydb/core/tx/columnshard/engines/reader/common/conveyor_task.h b/ydb/core/tx/columnshard/engines/reader/common/conveyor_task.h index 0342577c255..504bb41b648 100644 --- a/ydb/core/tx/columnshard/engines/reader/common/conveyor_task.h +++ b/ydb/core/tx/columnshard/engines/reader/common/conveyor_task.h @@ -28,7 +28,7 @@ public: virtual TConclusionStatus DoExecuteImpl() = 0; protected: - virtual TConclusionStatus DoExecute(const std::shared_ptr<NConveyor::ITask>& taskPtr) override final; + virtual void DoExecute(const std::shared_ptr<NConveyor::ITask>& taskPtr) override final; virtual void DoOnCannotExecute(const TString& reason) override; public: diff --git a/ydb/core/tx/columnshard/engines/scheme/versions/abstract_scheme.cpp b/ydb/core/tx/columnshard/engines/scheme/versions/abstract_scheme.cpp index b9a16ee0680..833ea12ff37 100644 --- a/ydb/core/tx/columnshard/engines/scheme/versions/abstract_scheme.cpp +++ b/ydb/core/tx/columnshard/engines/scheme/versions/abstract_scheme.cpp @@ -340,6 +340,7 @@ TConclusion<TWritePortionInfoWithBlobsResult> ISnapshotSchema::PrepareForWrite(c TConclusion<std::shared_ptr<NArrow::NAccessor::IChunkedArray>> arrToWrite = loader->GetAccessorConstructor()->Construct(accessor, loader->BuildAccessorContext(accessor->GetRecordsCount())); if (arrToWrite.IsFail()) { + AFL_ERROR(NKikimrServices::TX_COLUMNSHARD)("event", "cannot build accessor")("reason", arrToWrite.GetErrorMessage()); return arrToWrite; } diff --git a/ydb/core/tx/columnshard/normalizer/portion/chunks.cpp b/ydb/core/tx/columnshard/normalizer/portion/chunks.cpp index 7cfdfcd2130..1d7e9656cd4 100644 --- a/ydb/core/tx/columnshard/normalizer/portion/chunks.cpp +++ b/ydb/core/tx/columnshard/normalizer/portion/chunks.cpp @@ -52,7 +52,7 @@ private: TNormalizationContext NormContext; protected: - virtual TConclusionStatus DoExecute(const std::shared_ptr<NConveyor::ITask>& /*taskPtr*/) override { + virtual void DoExecute(const std::shared_ptr<NConveyor::ITask>& /*taskPtr*/) override { for (auto&& chunkInfo : Chunks) { const auto& blobRange = chunkInfo.GetBlobRange(); @@ -73,7 +73,6 @@ protected: auto changes = std::make_shared<TChunksNormalizer::TNormalizerResult>(std::move(Chunks)); TActorContext::AsActorContext().Send( NormContext.GetShardActor(), std::make_unique<NColumnShard::TEvPrivate::TEvNormalizerResult>(changes)); - return TConclusionStatus::Success(); } public: diff --git a/ydb/core/tx/columnshard/operations/batch_builder/builder.cpp b/ydb/core/tx/columnshard/operations/batch_builder/builder.cpp index f5c44815990..7f886c8be62 100644 --- a/ydb/core/tx/columnshard/operations/batch_builder/builder.cpp +++ b/ydb/core/tx/columnshard/operations/batch_builder/builder.cpp @@ -21,19 +21,19 @@ void TBuildBatchesTask::ReplyError(const TString& message, const NColumnShard::T TActorContext::AsActorContext().Send(Context.GetTabletActorId(), result.release()); } -TConclusionStatus TBuildBatchesTask::DoExecute(const std::shared_ptr<ITask>& /*taskPtr*/) { +void TBuildBatchesTask::DoExecute(const std::shared_ptr<ITask>& /*taskPtr*/) { const NActors::TLogContextGuard lGuard = NActors::TLogContextBuilder::Build()("scope", "TBuildBatchesTask::DoExecute"); if (!Context.IsActive()) { AFL_WARN(NKikimrServices::TX_COLUMNSHARD_WRITE)("event", "abort_external"); ReplyError("writing aborted", NColumnShard::TEvPrivate::TEvWriteBlobsResult::EErrorClass::Internal); - return TConclusionStatus::Fail("writing aborted"); + return; } TConclusion<std::shared_ptr<arrow::RecordBatch>> batchConclusion = WriteData.GetData()->ExtractBatch(); if (batchConclusion.IsFail()) { AFL_WARN(NKikimrServices::TX_COLUMNSHARD_WRITE)("event", "abort_on_extract")("reason", batchConclusion.GetErrorMessage()); ReplyError("cannot extract incoming batch: " + batchConclusion.GetErrorMessage(), NColumnShard::TEvPrivate::TEvWriteBlobsResult::EErrorClass::Internal); - return TConclusionStatus::Fail("cannot extract incoming batch: " + batchConclusion.GetErrorMessage()); + return; } Context.GetWritingCounters()->OnIncomingData(NArrow::GetBatchDataSize(*batchConclusion)); @@ -43,7 +43,7 @@ TConclusionStatus TBuildBatchesTask::DoExecute(const std::shared_ptr<ITask>& /*t AFL_WARN(NKikimrServices::TX_COLUMNSHARD_WRITE)("event", "abort_on_prepare")("reason", preparedConclusion.GetErrorMessage()); ReplyError("cannot prepare incoming batch: " + preparedConclusion.GetErrorMessage(), NColumnShard::TEvPrivate::TEvWriteBlobsResult::EErrorClass::Request); - return TConclusionStatus::Fail("cannot prepare incoming batch: " + preparedConclusion.GetErrorMessage()); + return; } auto batch = preparedConclusion.DetachResult(); std::shared_ptr<IMerger> merger; @@ -60,7 +60,7 @@ TConclusionStatus TBuildBatchesTask::DoExecute(const std::shared_ptr<ITask>& /*t new NWritingPortions::TEvAddInsertedDataToBuffer( std::make_shared<NEvWrite::TWriteData>(WriteData), batch, std::make_shared<TWritingContext>(Context))); } - return TConclusionStatus::Success(); + return; } else { auto insertionConclusion = Context.GetActualSchema()->CheckColumnsDefault(defaultFields); auto conclusion = Context.GetActualSchema()->BuildDefaultBatch(Context.GetActualSchema()->GetIndexInfo().ArrowSchema(), 1, true); @@ -92,14 +92,12 @@ TConclusionStatus TBuildBatchesTask::DoExecute(const std::shared_ptr<ITask>& /*t new NWritingPortions::TEvAddInsertedDataToBuffer( std::make_shared<NEvWrite::TWriteData>(WriteData), batch, std::make_shared<TWritingContext>(Context))); } - return TConclusionStatus::Success(); + return; } } std::shared_ptr<NDataReader::IRestoreTask> task = std::make_shared<NOlap::TModificationRestoreTask>(std::move(WriteData), merger, ActualSnapshot, batch, Context); NActors::TActivationContext::AsActorContext().Register(new NDataReader::TActor(task)); - - return TConclusionStatus::Success(); } } // namespace NKikimr::NOlap diff --git a/ydb/core/tx/columnshard/operations/batch_builder/builder.h b/ydb/core/tx/columnshard/operations/batch_builder/builder.h index 1ec9b1ae551..50e15654297 100644 --- a/ydb/core/tx/columnshard/operations/batch_builder/builder.h +++ b/ydb/core/tx/columnshard/operations/batch_builder/builder.h @@ -17,7 +17,7 @@ private: void ReplyError(const TString& message, const NColumnShard::TEvPrivate::TEvWriteBlobsResult::EErrorClass errorClass); protected: - virtual TConclusionStatus DoExecute(const std::shared_ptr<ITask>& taskPtr) override; + virtual void DoExecute(const std::shared_ptr<ITask>& taskPtr) override; public: virtual TString GetTaskClassIdentifier() const override { diff --git a/ydb/core/tx/columnshard/operations/events.h b/ydb/core/tx/columnshard/operations/events.h index 30328c0bfc4..bf0ac210c62 100644 --- a/ydb/core/tx/columnshard/operations/events.h +++ b/ydb/core/tx/columnshard/operations/events.h @@ -29,10 +29,36 @@ private: std::shared_ptr<NEvWrite::TWriteMeta> WriteMeta; YDB_READONLY(ui64, DataSize, 0); YDB_READONLY(bool, NoDataToWrite, false); + TString ErrorMessage; + std::optional<bool> IsInternalErrorFlag; std::shared_ptr<arrow::RecordBatch> PKBatch; ui32 RecordsCount; public: + TWriteResult& SetErrorMessage(const TString& value, const bool isInternal) { + AFL_VERIFY(!ErrorMessage); + IsInternalErrorFlag = isInternal; + ErrorMessage = value; + return *this; + } + + bool IsInternalError() const { + AFL_VERIFY_DEBUG(!!IsInternalErrorFlag); + if (!IsInternalErrorFlag) { + return true; + } + return *IsInternalErrorFlag; + } + + const TString& GetErrorMessage() const { + static TString undefinedMessage = "UNKNOWN_WRITE_RESULT_MESSAGE"; + AFL_VERIFY_DEBUG(!!ErrorMessage); + if (!ErrorMessage) { + return undefinedMessage; + } + return ErrorMessage; + } + const std::shared_ptr<arrow::RecordBatch>& GetPKBatchVerified() const { AFL_VERIFY(PKBatch); return PKBatch; @@ -91,11 +117,16 @@ namespace NKikimr::NColumnShard::NPrivateEvents::NWrite { class TEvWritePortionResult: public TEventLocal<TEvWritePortionResult, TEvPrivate::EvWritePortionResult> { private: YDB_READONLY_DEF(NKikimrProto::EReplyStatus, WriteStatus); - YDB_READONLY_DEF(std::shared_ptr<NOlap::IBlobsWritingAction>, WriteAction); + std::optional<std::shared_ptr<NOlap::IBlobsWritingAction>> WriteAction; bool Detached = false; TInsertedPortions InsertedData; public: + const std::shared_ptr<NOlap::IBlobsWritingAction>& GetWriteAction() const { + AFL_VERIFY(!!WriteAction); + return *WriteAction; + } + const TInsertedPortions& DetachInsertedData() { AFL_VERIFY(!Detached); Detached = true; diff --git a/ydb/core/tx/columnshard/operations/slice_builder/builder.cpp b/ydb/core/tx/columnshard/operations/slice_builder/builder.cpp index 866720bb1a6..fb79e54284b 100644 --- a/ydb/core/tx/columnshard/operations/slice_builder/builder.cpp +++ b/ydb/core/tx/columnshard/operations/slice_builder/builder.cpp @@ -99,19 +99,19 @@ public: } }; -TConclusionStatus TBuildSlicesTask::DoExecute(const std::shared_ptr<ITask>& /*taskPtr*/) { +void TBuildSlicesTask::DoExecute(const std::shared_ptr<ITask>& /*taskPtr*/) { const NActors::TLogContextGuard g = NActors::TLogContextBuilder::Build(NKikimrServices::TX_COLUMNSHARD_WRITE)("tablet_id", TabletId)( "parent_id", Context.GetTabletActorId())("write_id", WriteData.GetWriteMeta().GetWriteId())( "table_id", WriteData.GetWriteMeta().GetTableId()); if (!Context.IsActive()) { AFL_WARN(NKikimrServices::TX_COLUMNSHARD_WRITE)("event", "abort_execution"); ReplyError("execution aborted", NColumnShard::TEvPrivate::TEvWriteBlobsResult::EErrorClass::Internal); - return TConclusionStatus::Fail("execution aborted"); + return; } if (!OriginalBatch) { AFL_WARN(NKikimrServices::TX_COLUMNSHARD_WRITE)("event", "ev_write_bad_data"); ReplyError("no data in batch", NColumnShard::TEvPrivate::TEvWriteBlobsResult::EErrorClass::Internal); - return TConclusionStatus::Fail("no data in batch"); + return; } if (WriteData.GetWritePortions()) { if (OriginalBatch->num_rows() == 0) { @@ -136,7 +136,7 @@ TConclusionStatus TBuildSlicesTask::DoExecute(const std::shared_ptr<ITask>& /*ta WriteData.GetWriteMeta().GetModificationType(), Context.GetStoragesManager(), Context.GetSplitterCounters()); if (portionConclusion.IsFail()) { ReplyError(portionConclusion.GetErrorMessage(), NColumnShard::TEvPrivate::TEvWriteBlobsResult::EErrorClass::Request); - return portionConclusion; + return; } portions.emplace_back(portionConclusion.DetachResult()); } @@ -157,7 +157,7 @@ TConclusionStatus TBuildSlicesTask::DoExecute(const std::shared_ptr<ITask>& /*ta "problem", subsetConclusion.GetErrorMessage()); ReplyError("unadaptable schema: " + subsetConclusion.GetErrorMessage(), NColumnShard::TEvPrivate::TEvWriteBlobsResult::EErrorClass::Internal); - return TConclusionStatus::Fail("cannot reorder schema: " + subsetConclusion.GetErrorMessage()); + return; } NArrow::TSchemaSubset subset = subsetConclusion.DetachResult(); @@ -186,9 +186,8 @@ TConclusionStatus TBuildSlicesTask::DoExecute(const std::shared_ptr<ITask>& /*ta TActorContext::AsActorContext().Send(Context.GetBufferizationInsertionActorId(), result.release()); } else { ReplyError("Cannot slice input to batches", NColumnShard::TEvPrivate::TEvWriteBlobsResult::EErrorClass::Internal); - return TConclusionStatus::Fail("Cannot slice input to batches"); + return; } } - return TConclusionStatus::Success(); } } // namespace NKikimr::NOlap diff --git a/ydb/core/tx/columnshard/operations/slice_builder/builder.h b/ydb/core/tx/columnshard/operations/slice_builder/builder.h index 0efd53378e3..5d6bb415977 100644 --- a/ydb/core/tx/columnshard/operations/slice_builder/builder.h +++ b/ydb/core/tx/columnshard/operations/slice_builder/builder.h @@ -18,7 +18,7 @@ private: void ReplyError(const TString& message, const NColumnShard::TEvPrivate::TEvWriteBlobsResult::EErrorClass errorClass); protected: - virtual TConclusionStatus DoExecute(const std::shared_ptr<ITask>& taskPtr) override; + virtual void DoExecute(const std::shared_ptr<ITask>& taskPtr) override; public: virtual TString GetTaskClassIdentifier() const override { diff --git a/ydb/core/tx/columnshard/operations/slice_builder/pack_builder.cpp b/ydb/core/tx/columnshard/operations/slice_builder/pack_builder.cpp index d3bc168d9b7..9b9b0e351b5 100644 --- a/ydb/core/tx/columnshard/operations/slice_builder/pack_builder.cpp +++ b/ydb/core/tx/columnshard/operations/slice_builder/pack_builder.cpp @@ -44,6 +44,11 @@ private: for (auto&& i : Portions) { portions.emplace_back(i.ExtractPortion()); } + if (putResult->GetPutStatus() != NKikimrProto::OK) { + for (auto&& i : WriteResults) { + i.SetErrorMessage("cannot put blobs: " + ::ToString(putResult->GetPutStatus()), true); + } + } NColumnShard::TInsertedPortions pack(std::move(WriteResults), std::move(portions)); auto result = std::make_unique<NColumnShard::NPrivateEvents::NWrite::TEvWritePortionResult>(putResult->GetPutStatus(), Action, std::move(pack)); @@ -97,6 +102,10 @@ public: if (Batches.size() == 1) { auto portionConclusion = context.GetActualSchema()->PrepareForWrite(context.GetActualSchema(), PathId, Batches.front().GetContainer(), ModificationType, context.GetStoragesManager(), context.GetSplitterCounters()); + if (portionConclusion.IsFail()) { + AFL_ERROR(NKikimrServices::TX_COLUMNSHARD)("event", "cannot prepare for write")("reason", portionConclusion.GetErrorMessage()); + return portionConclusion; + } result.emplace_back(portionConclusion.DetachResult()); } else { ui32 idx = 0; @@ -120,6 +129,8 @@ public: if (itBatchIndexes == i.GetColumnIndexes().end() || *itAllIndexes < *itBatchIndexes) { auto defaultColumn = indexInfo.BuildDefaultColumn(*itAllIndexes, i->num_rows(), false); if (defaultColumn.IsFail()) { + AFL_ERROR(NKikimrServices::TX_COLUMNSHARD)("event", "cannot build default column")( + "reason", defaultColumn.GetErrorMessage()); return defaultColumn; } gContainer->AddField(context.GetActualSchema()->GetFieldByIndexVerified(*itAllIndexes), defaultColumn.DetachResult()) @@ -155,19 +166,19 @@ public: stream.DrainAll(rbBuilder); auto portionConclusion = context.GetActualSchema()->PrepareForWrite(context.GetActualSchema(), PathId, rbBuilder.Finalize(), ModificationType, context.GetStoragesManager(), context.GetSplitterCounters()); + if (portionConclusion.IsFail()) { + AFL_ERROR(NKikimrServices::TX_COLUMNSHARD)("event", "cannot prepare for write")("reason", portionConclusion.GetErrorMessage()); + return portionConclusion; + } result.emplace_back(portionConclusion.DetachResult()); } return TConclusionStatus::Success(); } }; -TConclusionStatus TBuildPackSlicesTask::DoExecute(const std::shared_ptr<ITask>& /*taskPtr*/) { +void TBuildPackSlicesTask::DoExecute(const std::shared_ptr<ITask>& /*taskPtr*/) { const NActors::TLogContextGuard g = NActors::TLogContextBuilder::Build(NKikimrServices::TX_COLUMNSHARD_WRITE)("tablet_id", TabletId)( "parent_id", Context.GetTabletActorId())("path_id", PathId); - if (!Context.IsActive()) { - AFL_WARN(NKikimrServices::TX_COLUMNSHARD_WRITE)("event", "abort_execution"); - return TConclusionStatus::Fail("execution aborted"); - } NArrow::NMerger::TIntervalPositions splitPositions; for (auto&& unit : WriteUnits) { splitPositions.Merge(unit.GetData()->GetData()->GetSeparationPoints()); @@ -195,21 +206,39 @@ TConclusionStatus TBuildPackSlicesTask::DoExecute(const std::shared_ptr<ITask>& } } std::vector<TPortionWriteController::TInsertPortion> portionsToWrite; - for (auto&& i : slicesToMerge) { - auto conclusion = i.Finalize(Context, portionsToWrite); - if (conclusion.IsFail()) { - return conclusion; + TString cancelWritingReason; + if (!Context.IsActive()) { + AFL_WARN(NKikimrServices::TX_COLUMNSHARD_WRITE)("event", "abort_execution"); + cancelWritingReason = "execution aborted"; + } else { + for (auto&& i : slicesToMerge) { + auto conclusion = i.Finalize(Context, portionsToWrite); + if (conclusion.IsFail()) { + AFL_ERROR(NKikimrServices::TX_COLUMNSHARD)("event", "cannot build slice")("reason", conclusion.GetErrorMessage()); + cancelWritingReason = conclusion.GetErrorMessage(); + break; + } } } - auto actions = WriteUnits.front().GetData()->GetBlobsAction(); - auto writeController = - std::make_shared<TPortionWriteController>(Context.GetTabletActorId(), actions, std::move(writeResults), std::move(portionsToWrite)); - if (actions->NeedDraftTransaction()) { - TActorContext::AsActorContext().Send( - Context.GetTabletActorId(), std::make_unique<NColumnShard::TEvPrivate::TEvWriteDraft>(writeController)); + if (!cancelWritingReason) { + auto actions = WriteUnits.front().GetData()->GetBlobsAction(); + auto writeController = + std::make_shared<TPortionWriteController>(Context.GetTabletActorId(), actions, std::move(writeResults), std::move(portionsToWrite)); + if (actions->NeedDraftTransaction()) { + TActorContext::AsActorContext().Send( + Context.GetTabletActorId(), std::make_unique<NColumnShard::TEvPrivate::TEvWriteDraft>(writeController)); + } else { + TActorContext::AsActorContext().Register(NColumnShard::CreateWriteActor(TabletId, writeController, TInstant::Max())); + } } else { - TActorContext::AsActorContext().Register(NColumnShard::CreateWriteActor(TabletId, writeController, TInstant::Max())); + for (auto&& i : writeResults) { + i.SetErrorMessage(cancelWritingReason, false); + } + NColumnShard::TInsertedPortions pack(std::move(writeResults), std::vector<NColumnShard::TInsertedPortion>()); + auto result = + std::make_unique<NColumnShard::NPrivateEvents::NWrite::TEvWritePortionResult>(NKikimrProto::EReplyStatus::ERROR, nullptr, std::move(pack)); + TActorContext::AsActorContext().Send(Context.GetTabletActorId(), result.release()); + } - return TConclusionStatus::Success(); } } // namespace NKikimr::NOlap::NWritingPortions diff --git a/ydb/core/tx/columnshard/operations/slice_builder/pack_builder.h b/ydb/core/tx/columnshard/operations/slice_builder/pack_builder.h index f8f46e35ba6..40ebd564677 100644 --- a/ydb/core/tx/columnshard/operations/slice_builder/pack_builder.h +++ b/ydb/core/tx/columnshard/operations/slice_builder/pack_builder.h @@ -35,7 +35,7 @@ private: std::optional<std::vector<NArrow::TSerializedBatch>> BuildSlices(); protected: - virtual TConclusionStatus DoExecute(const std::shared_ptr<ITask>& taskPtr) override; + virtual void DoExecute(const std::shared_ptr<ITask>& taskPtr) override; public: virtual TString GetTaskClassIdentifier() const override { diff --git a/ydb/core/tx/conveyor/service/worker.cpp b/ydb/core/tx/conveyor/service/worker.cpp index 6450725f15f..a9e68d7b935 100644 --- a/ydb/core/tx/conveyor/service/worker.cpp +++ b/ydb/core/tx/conveyor/service/worker.cpp @@ -8,7 +8,7 @@ void TWorker::ExecuteTask(std::vector<TWorkerTask>&& workerTasks) { std::vector<ui64> processes; instants.emplace_back(TMonotonic::Now()); for (auto&& t : workerTasks) { - Y_UNUSED(t.GetTask()->Execute(t.GetTaskSignals(), t.GetTask())); + t.GetTask()->Execute(t.GetTaskSignals(), t.GetTask()); instants.emplace_back(TMonotonic::Now()); processes.emplace_back(t.GetProcessId()); } diff --git a/ydb/core/tx/conveyor/usage/abstract.cpp b/ydb/core/tx/conveyor/usage/abstract.cpp index 55c19e7bba8..2c670d5a022 100644 --- a/ydb/core/tx/conveyor/usage/abstract.cpp +++ b/ydb/core/tx/conveyor/usage/abstract.cpp @@ -8,30 +8,22 @@ #include <util/string/builder.h> namespace NKikimr::NConveyor { -TConclusionStatus ITask::Execute(std::shared_ptr<TTaskSignals> signals, const std::shared_ptr<ITask>& taskPtr) { +void ITask::Execute(std::shared_ptr<TTaskSignals> signals, const std::shared_ptr<ITask>& taskPtr) { AFL_VERIFY(!ExecutedFlag); ExecutedFlag = true; const TMonotonic start = TMonotonic::Now(); try { - TConclusionStatus result = DoExecute(taskPtr); - if (result.IsFail()) { - if (signals) { - signals->Fails->Add(1); - signals->FailsDuration->Add((TMonotonic::Now() - start).MicroSeconds()); - } - } else { - if (signals) { - signals->Success->Add(1); - signals->SuccessDuration->Add((TMonotonic::Now() - start).MicroSeconds()); - } + DoExecute(taskPtr); + if (signals) { + signals->Success->Add(1); + signals->SuccessDuration->Add((TMonotonic::Now() - start).MicroSeconds()); } - return result; } catch (...) { if (signals) { signals->Fails->Add(1); signals->FailsDuration->Add((TMonotonic::Now() - start).MicroSeconds()); } - return TConclusionStatus::Fail("exception: " + CurrentExceptionMessage()); + AFL_ERROR(NKikimrServices::TX_CONVEYOR)("event", "exception_on_execute")("message", CurrentExceptionMessage()); } } diff --git a/ydb/core/tx/conveyor/usage/abstract.h b/ydb/core/tx/conveyor/usage/abstract.h index d66e203274b..593d0c1feba 100644 --- a/ydb/core/tx/conveyor/usage/abstract.h +++ b/ydb/core/tx/conveyor/usage/abstract.h @@ -65,7 +65,7 @@ private: YDB_ACCESSOR(EPriority, Priority, EPriority::Normal); bool ExecutedFlag = false; protected: - virtual TConclusionStatus DoExecute(const std::shared_ptr<ITask>& taskPtr) = 0; + virtual void DoExecute(const std::shared_ptr<ITask>& taskPtr) = 0; virtual void DoOnCannotExecute(const TString& reason); public: using TPtr = std::shared_ptr<ITask>; @@ -76,7 +76,7 @@ public: void OnCannotExecute(const TString& reason) { return DoOnCannotExecute(reason); } - TConclusionStatus Execute(std::shared_ptr<TTaskSignals> signals, const std::shared_ptr<ITask>& taskPtr); + void Execute(std::shared_ptr<TTaskSignals> signals, const std::shared_ptr<ITask>& taskPtr); }; } diff --git a/ydb/library/conclusion/generic/string_status.h b/ydb/library/conclusion/generic/string_status.h index 81541395d05..c485bcc7e75 100644 --- a/ydb/library/conclusion/generic/string_status.h +++ b/ydb/library/conclusion/generic/string_status.h @@ -44,6 +44,7 @@ public: } [[nodiscard]] TString GetErrorMessage() const { + Y_ABORT_UNLESS(TBase::IsFail()); return TBase::GetErrorDescription(); } }; |
