summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorivanmorozov333 <[email protected]>2025-04-12 21:20:28 +0300
committerGitHub <[email protected]>2025-04-12 21:20:28 +0300
commitff4f214225d313fd54806ef83fc9c4016df7ebc8 (patch)
tree1411743135bb744b3857267356f6031866b2b6bf
parent9eee3dd53ab019dde1cc33f1e0697c06d59137f1 (diff)
fix errors processing for json (#17139)
-rw-r--r--ydb/core/formats/arrow/program/graph_execute.cpp5
-rw-r--r--ydb/core/kqp/ut/olap/json_ut.cpp86
-rw-r--r--ydb/core/tx/columnshard/blobs_action/transaction/tx_blobs_written.cpp6
-rw-r--r--ydb/core/tx/columnshard/blobs_action/transaction/tx_blobs_written.h4
-rw-r--r--ydb/core/tx/columnshard/columnshard__write.cpp4
-rw-r--r--ydb/core/tx/columnshard/columnshard_impl.cpp6
-rw-r--r--ydb/core/tx/columnshard/engines/reader/common/conveyor_task.cpp3
-rw-r--r--ydb/core/tx/columnshard/engines/reader/common/conveyor_task.h2
-rw-r--r--ydb/core/tx/columnshard/engines/scheme/versions/abstract_scheme.cpp1
-rw-r--r--ydb/core/tx/columnshard/normalizer/portion/chunks.cpp3
-rw-r--r--ydb/core/tx/columnshard/operations/batch_builder/builder.cpp14
-rw-r--r--ydb/core/tx/columnshard/operations/batch_builder/builder.h2
-rw-r--r--ydb/core/tx/columnshard/operations/events.h33
-rw-r--r--ydb/core/tx/columnshard/operations/slice_builder/builder.cpp13
-rw-r--r--ydb/core/tx/columnshard/operations/slice_builder/builder.h2
-rw-r--r--ydb/core/tx/columnshard/operations/slice_builder/pack_builder.cpp63
-rw-r--r--ydb/core/tx/columnshard/operations/slice_builder/pack_builder.h2
-rw-r--r--ydb/core/tx/conveyor/service/worker.cpp2
-rw-r--r--ydb/core/tx/conveyor/usage/abstract.cpp20
-rw-r--r--ydb/core/tx/conveyor/usage/abstract.h4
-rw-r--r--ydb/library/conclusion/generic/string_status.h1
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();
}
};