diff options
| author | ivanmorozov333 <[email protected]> | 2025-02-21 12:55:07 +0300 |
|---|---|---|
| committer | GitHub <[email protected]> | 2025-02-21 12:55:07 +0300 |
| commit | 2d9caae0e483e34ed4549ccee02cd2d392b65074 (patch) | |
| tree | 2f55109ba4691d9e86056420c438bf4182e6cd63 | |
| parent | 1b2ad9e5d9552a68bb9c79164de7980d87b11e93 (diff) | |
json optimization usage fixes (#14748)
111 files changed, 1064 insertions, 772 deletions
diff --git a/ydb/library/formats/arrow/accessor/abstract/accessor.cpp b/ydb/core/formats/arrow/accessor/abstract/accessor.cpp index a27310e8116..926f3aa8192 100644 --- a/ydb/library/formats/arrow/accessor/abstract/accessor.cpp +++ b/ydb/core/formats/arrow/accessor/abstract/accessor.cpp @@ -119,6 +119,16 @@ std::shared_ptr<IChunkedArray> IChunkedArray::DoApplyFilter(const TColumnFilter& } } +std::shared_ptr<IChunkedArray> IChunkedArray::ApplyFilter(const TColumnFilter& filter, const std::shared_ptr<IChunkedArray>& selfPtr) const { + if (filter.IsTotalAllowFilter()) { + return selfPtr; + } + if (filter.IsTotalDenyFilter()) { + return TTrivialArray::BuildEmpty(GetDataType()); + } + return DoApplyFilter(filter); +} + TString IChunkedArray::TReader::DebugString(const ui32 position) const { auto address = GetReadChunk(position); return NArrow::DebugString(address.GetArray(), address.GetPosition()); diff --git a/ydb/library/formats/arrow/accessor/abstract/accessor.h b/ydb/core/formats/arrow/accessor/abstract/accessor.h index 15091640923..5cbbd37e6fc 100644 --- a/ydb/library/formats/arrow/accessor/abstract/accessor.h +++ b/ydb/core/formats/arrow/accessor/abstract/accessor.h @@ -320,9 +320,7 @@ protected: } public: - std::shared_ptr<IChunkedArray> ApplyFilter(const TColumnFilter& filter) const { - return DoApplyFilter(filter); - } + std::shared_ptr<IChunkedArray> ApplyFilter(const TColumnFilter& filter, const std::shared_ptr<IChunkedArray>& selfPtr) const; NJson::TJsonValue DebugJson() const { NJson::TJsonValue result = NJson::JSON_MAP; diff --git a/ydb/core/formats/arrow/accessor/abstract/constructor.h b/ydb/core/formats/arrow/accessor/abstract/constructor.h index c8c4aeb3449..f942e42e581 100644 --- a/ydb/core/formats/arrow/accessor/abstract/constructor.h +++ b/ydb/core/formats/arrow/accessor/abstract/constructor.h @@ -1,7 +1,8 @@ #pragma once -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> -#include <ydb/library/formats/arrow/accessor/common/chunk_data.h> +#include <ydb/core/formats/arrow/accessor/abstract/accessor.h> +#include <ydb/core/formats/arrow/accessor/common/chunk_data.h> + #include <ydb/library/formats/arrow/protos/accessor.pb.h> #include <ydb/services/bg_tasks/abstract/interface.h> diff --git a/ydb/core/formats/arrow/accessor/abstract/ya.make b/ydb/core/formats/arrow/accessor/abstract/ya.make index 4e07801b228..6b239c4ad15 100644 --- a/ydb/core/formats/arrow/accessor/abstract/ya.make +++ b/ydb/core/formats/arrow/accessor/abstract/ya.make @@ -5,14 +5,16 @@ PEERDIR( ydb/library/conclusion ydb/services/metadata/abstract ydb/library/actors/core - ydb/library/formats/arrow/accessor/abstract - ydb/library/formats/arrow/accessor/common + ydb/core/formats/arrow/accessor/common ydb/library/formats/arrow/protos ) SRCS( constructor.cpp request.cpp + accessor.cpp ) +GENERATE_ENUM_SERIALIZATION(accessor.h) + END() diff --git a/ydb/library/formats/arrow/accessor/common/chunk_data.cpp b/ydb/core/formats/arrow/accessor/common/chunk_data.cpp index 2356af1bd43..2356af1bd43 100644 --- a/ydb/library/formats/arrow/accessor/common/chunk_data.cpp +++ b/ydb/core/formats/arrow/accessor/common/chunk_data.cpp diff --git a/ydb/library/formats/arrow/accessor/common/chunk_data.h b/ydb/core/formats/arrow/accessor/common/chunk_data.h index ccbb318faaa..ccbb318faaa 100644 --- a/ydb/library/formats/arrow/accessor/common/chunk_data.h +++ b/ydb/core/formats/arrow/accessor/common/chunk_data.h diff --git a/ydb/library/formats/arrow/accessor/common/const.cpp b/ydb/core/formats/arrow/accessor/common/const.cpp index 926a9ca94de..926a9ca94de 100644 --- a/ydb/library/formats/arrow/accessor/common/const.cpp +++ b/ydb/core/formats/arrow/accessor/common/const.cpp diff --git a/ydb/library/formats/arrow/accessor/common/const.h b/ydb/core/formats/arrow/accessor/common/const.h index 3de41402c01..3de41402c01 100644 --- a/ydb/library/formats/arrow/accessor/common/const.h +++ b/ydb/core/formats/arrow/accessor/common/const.h diff --git a/ydb/library/formats/arrow/accessor/common/ya.make b/ydb/core/formats/arrow/accessor/common/ya.make index 9f2c4e95b85..9f2c4e95b85 100644 --- a/ydb/library/formats/arrow/accessor/common/ya.make +++ b/ydb/core/formats/arrow/accessor/common/ya.make diff --git a/ydb/library/formats/arrow/accessor/composite/accessor.cpp b/ydb/core/formats/arrow/accessor/composite/accessor.cpp index 85e4396b994..85e4396b994 100644 --- a/ydb/library/formats/arrow/accessor/composite/accessor.cpp +++ b/ydb/core/formats/arrow/accessor/composite/accessor.cpp diff --git a/ydb/library/formats/arrow/accessor/composite/accessor.h b/ydb/core/formats/arrow/accessor/composite/accessor.h index 18787eff5ea..bf3e59b94d8 100644 --- a/ydb/library/formats/arrow/accessor/composite/accessor.h +++ b/ydb/core/formats/arrow/accessor/composite/accessor.h @@ -1,6 +1,7 @@ #pragma once +#include <ydb/core/formats/arrow/accessor/abstract/accessor.h> + #include <ydb/library/accessor/accessor.h> -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> namespace NKikimr::NArrow::NAccessor { @@ -59,6 +60,7 @@ public: std::vector<std::shared_ptr<NArrow::NAccessor::IChunkedArray>> Chunks; const std::shared_ptr<arrow::DataType> Type; bool Finished = false; + public: TBuilder(const std::shared_ptr<arrow::DataType>& type) : Type(type) { diff --git a/ydb/library/formats/arrow/accessor/composite/ya.make b/ydb/core/formats/arrow/accessor/composite/ya.make index d4192f0aa0b..d4192f0aa0b 100644 --- a/ydb/library/formats/arrow/accessor/composite/ya.make +++ b/ydb/core/formats/arrow/accessor/composite/ya.make diff --git a/ydb/core/formats/arrow/accessor/composite_serial/accessor.h b/ydb/core/formats/arrow/accessor/composite_serial/accessor.h index ab09f0faaf0..3721754007a 100644 --- a/ydb/core/formats/arrow/accessor/composite_serial/accessor.h +++ b/ydb/core/formats/arrow/accessor/composite_serial/accessor.h @@ -1,9 +1,8 @@ #pragma once +#include <ydb/core/formats/arrow/accessor/abstract/accessor.h> +#include <ydb/core/formats/arrow/accessor/composite/accessor.h> #include <ydb/core/formats/arrow/save_load/loader.h> -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> -#include <ydb/library/formats/arrow/accessor/composite/accessor.h> - namespace NKikimr::NArrow::NAccessor { class TDeserializeChunkedArray: public ICompositeChunkedArray { @@ -16,6 +15,7 @@ public: YDB_READONLY(ui32, RecordsCount, 0); std::shared_ptr<IChunkedArray> PredefinedArray; const TString Data; + const TStringBuf DataBuffer; public: TChunk(const std::shared_ptr<IChunkedArray>& predefinedArray) @@ -29,11 +29,21 @@ public: , Data(data) { } + TChunk(const ui32 recordsCount, const TStringBuf dataBuffer) + : RecordsCount(recordsCount) + , DataBuffer(dataBuffer) { + } + std::shared_ptr<IChunkedArray> GetArrayVerified(const std::shared_ptr<TColumnLoader>& loader) const { if (PredefinedArray) { return PredefinedArray; } - return loader->ApplyVerified(Data, RecordsCount); + if (!!Data) { + return loader->ApplyVerified(Data, RecordsCount); + } else { + AFL_VERIFY(!!DataBuffer); + return loader->ApplyVerified(TString(DataBuffer.data(), DataBuffer.size()), RecordsCount); + } } }; @@ -77,12 +87,12 @@ protected: } public: - TDeserializeChunkedArray(const ui64 recordsCount, const std::shared_ptr<TColumnLoader>& loader, std::vector<TChunk>&& chunks, const bool forLazyInitialization = false) + TDeserializeChunkedArray(const ui64 recordsCount, const std::shared_ptr<TColumnLoader>& loader, std::vector<TChunk>&& chunks, + const bool forLazyInitialization = false) : TBase(recordsCount, NArrow::NAccessor::IChunkedArray::EType::SerializedChunkedArray, loader->GetField()->type()) , Loader(loader) , Chunks(std::move(chunks)) - , ForLazyInitialization(forLazyInitialization) - { + , ForLazyInitialization(forLazyInitialization) { AFL_VERIFY(Loader); } }; diff --git a/ydb/core/formats/arrow/accessor/composite_serial/ya.make b/ydb/core/formats/arrow/accessor/composite_serial/ya.make index e8095e99028..5551f3565a2 100644 --- a/ydb/core/formats/arrow/accessor/composite_serial/ya.make +++ b/ydb/core/formats/arrow/accessor/composite_serial/ya.make @@ -2,7 +2,7 @@ LIBRARY() PEERDIR( contrib/libs/apache/arrow - ydb/library/formats/arrow/accessor/abstract + ydb/core/formats/arrow/accessor/abstract ydb/core/formats/arrow/common ydb/core/formats/arrow/save_load ) diff --git a/ydb/core/formats/arrow/accessor/plain/accessor.cpp b/ydb/core/formats/arrow/accessor/plain/accessor.cpp index e2e3cd672d4..262343577c5 100644 --- a/ydb/core/formats/arrow/accessor/plain/accessor.cpp +++ b/ydb/core/formats/arrow/accessor/plain/accessor.cpp @@ -5,6 +5,8 @@ #include <ydb/core/formats/arrow/size_calcer.h> #include <ydb/core/formats/arrow/splitter/simple.h> +#include <ydb/library/formats/arrow/simple_arrays_cache.h> + namespace NKikimr::NArrow::NAccessor { std::optional<ui64> TTrivialArray::DoGetRawSize() const { @@ -20,6 +22,10 @@ ui32 TTrivialArray::DoGetValueRawBytes() const { return NArrow::GetArrayDataSize(Array); } +std::shared_ptr<TTrivialArray> TTrivialArray::BuildEmpty(const std::shared_ptr<arrow::DataType>& type) { + return std::make_shared<TTrivialArray>(TThreadSimpleArraysCache::GetNull(type, 0)); +} + namespace { class TChunkAccessor { private: diff --git a/ydb/core/formats/arrow/accessor/plain/accessor.h b/ydb/core/formats/arrow/accessor/plain/accessor.h index 9927beed2f0..9096b57ff88 100644 --- a/ydb/core/formats/arrow/accessor/plain/accessor.h +++ b/ydb/core/formats/arrow/accessor/plain/accessor.h @@ -1,6 +1,8 @@ #pragma once -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> +#include <ydb/core/formats/arrow/accessor/abstract/accessor.h> + #include <ydb/library/formats/arrow/arrow_helpers.h> +#include <ydb/library/formats/arrow/switch/switch_type.h> #include <ydb/library/formats/arrow/validation/validation.h> namespace NKikimr::NArrow::NAccessor { @@ -34,6 +36,8 @@ public: return Array; } + static std::shared_ptr<TTrivialArray> BuildEmpty(const std::shared_ptr<arrow::DataType>& type); + TTrivialArray(const std::shared_ptr<arrow::Array>& data) : TBase(data->length(), EType::Array, data->type()) , Array(data) { diff --git a/ydb/core/formats/arrow/accessor/plain/constructor.cpp b/ydb/core/formats/arrow/accessor/plain/constructor.cpp index d73c906911d..b4ff335a24e 100644 --- a/ydb/core/formats/arrow/accessor/plain/constructor.cpp +++ b/ydb/core/formats/arrow/accessor/plain/constructor.cpp @@ -1,9 +1,9 @@ #include "accessor.h" #include "constructor.h" +#include <ydb/core/formats/arrow/accessor/abstract/accessor.h> #include <ydb/core/formats/arrow/serializer/abstract.h> -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> #include <ydb/library/formats/arrow/arrow_helpers.h> #include <ydb/library/formats/arrow/simple_arrays_cache.h> diff --git a/ydb/core/formats/arrow/accessor/plain/constructor.h b/ydb/core/formats/arrow/accessor/plain/constructor.h index e16f8810d6d..884beff48e2 100644 --- a/ydb/core/formats/arrow/accessor/plain/constructor.h +++ b/ydb/core/formats/arrow/accessor/plain/constructor.h @@ -1,7 +1,6 @@ #pragma once #include <ydb/core/formats/arrow/accessor/abstract/constructor.h> - -#include <ydb/library/formats/arrow/accessor/common/const.h> +#include <ydb/core/formats/arrow/accessor/common/const.h> namespace NKikimr::NArrow::NAccessor::NPlain { diff --git a/ydb/core/formats/arrow/accessor/plain/request.h b/ydb/core/formats/arrow/accessor/plain/request.h index 19a8390f2df..02f6cce8560 100644 --- a/ydb/core/formats/arrow/accessor/plain/request.h +++ b/ydb/core/formats/arrow/accessor/plain/request.h @@ -1,6 +1,6 @@ #pragma once #include <ydb/core/formats/arrow/accessor/abstract/request.h> -#include <ydb/library/formats/arrow/accessor/common/const.h> +#include <ydb/core/formats/arrow/accessor/common/const.h> namespace NKikimr::NArrow::NAccessor::NPlain { diff --git a/ydb/core/formats/arrow/accessor/sparsed/accessor.cpp b/ydb/core/formats/arrow/accessor/sparsed/accessor.cpp index 41b0c7fbd27..4411a3dea50 100644 --- a/ydb/core/formats/arrow/accessor/sparsed/accessor.cpp +++ b/ydb/core/formats/arrow/accessor/sparsed/accessor.cpp @@ -1,5 +1,6 @@ #include "accessor.h" +#include <ydb/core/formats/arrow/arrow_filter.h> #include <ydb/core/formats/arrow/save_load/loader.h> #include <ydb/core/formats/arrow/size_calcer.h> #include <ydb/core/formats/arrow/splitter/simple.h> @@ -8,26 +9,24 @@ namespace NKikimr::NArrow::NAccessor { -TSparsedArray::TSparsedArray(const IChunkedArray& defaultArray, const std::shared_ptr<arrow::Scalar>& defaultValue) - : TBase(defaultArray.GetRecordsCount(), EType::SparsedArray, defaultArray.GetDataType()) - , DefaultValue(defaultValue) { - if (DefaultValue) { - AFL_VERIFY(DefaultValue->type->id() == defaultArray.GetDataType()->id()); +std::shared_ptr<TSparsedArray> TSparsedArray::Make(const IChunkedArray& defaultArray, const std::shared_ptr<arrow::Scalar>& defaultValue) { + if (defaultValue) { + AFL_VERIFY(defaultValue->type->id() == defaultArray.GetDataType()->id()); } std::optional<TFullDataAddress> current; std::shared_ptr<arrow::RecordBatch> records; ui32 sparsedRecordsCount = 0; - AFL_VERIFY(SwitchType(GetDataType()->id(), [&](const auto& type) { + AFL_VERIFY(SwitchType(defaultArray.GetDataType()->id(), [&](const auto& type) { using TWrap = std::decay_t<decltype(type)>; using TScalar = typename arrow::TypeTraits<typename TWrap::T>::ScalarType; using TArray = typename arrow::TypeTraits<typename TWrap::T>::ArrayType; using TBuilder = typename arrow::TypeTraits<typename TWrap::T>::BuilderType; - auto builderValue = NArrow::MakeBuilder(GetDataType()); + auto builderValue = NArrow::MakeBuilder(defaultArray.GetDataType()); TBuilder* builderValueImpl = (TBuilder*)builderValue.get(); auto builderIndex = NArrow::MakeBuilder(arrow::uint32()); arrow::UInt32Builder* builderIndexImpl = (arrow::UInt32Builder*)builderIndex.get(); - auto scalar = static_pointer_cast<TScalar>(DefaultValue); - for (ui32 pos = 0; pos < GetRecordsCount();) { + auto scalar = static_pointer_cast<TScalar>(defaultValue); + for (ui32 pos = 0; pos < defaultArray.GetRecordsCount();) { current = defaultArray.GetChunk(current, pos); auto typedArray = static_pointer_cast<TArray>(current->GetArray()); for (ui32 i = 0; i < typedArray->length(); ++i) { @@ -38,7 +37,7 @@ TSparsedArray::TSparsedArray(const IChunkedArray& defaultArray, const std::share } else if constexpr (arrow::has_c_type<typename TWrap::T>()) { isDefault = scalar->value == typedArray->Value(i); } else { - AFL_VERIFY(false)("type", GetDataType()->ToString()); + AFL_VERIFY(false)("type", defaultArray.GetDataType()->ToString()); } } else { isDefault = typedArray->IsNull(i); @@ -53,35 +52,25 @@ TSparsedArray::TSparsedArray(const IChunkedArray& defaultArray, const std::share NArrow::TStatusValidator::Validate(builderIndexImpl->Append(pos + i)); ++sparsedRecordsCount; } else { - AFL_VERIFY(false)("type", GetDataType()->ToString()); + AFL_VERIFY(false)("type", defaultArray.GetDataType()->ToString()); } } } pos = current->GetAddress().GetGlobalFinishPosition(); - AFL_VERIFY(pos <= GetRecordsCount()); + AFL_VERIFY(pos <= defaultArray.GetRecordsCount()); } std::vector<std::shared_ptr<arrow::Array>> columns = { NArrow::TStatusValidator::GetValid(builderIndex->Finish()), NArrow::TStatusValidator::GetValid(builderValue->Finish()) }; - records = arrow::RecordBatch::Make(BuildSchema(GetDataType()), sparsedRecordsCount, columns); + records = arrow::RecordBatch::Make(BuildSchema(defaultArray.GetDataType()), sparsedRecordsCount, columns); AFL_VERIFY_DEBUG(records->ValidateFull().ok()); return true; })); - AFL_VERIFY(records); - Records.emplace_back(0, GetRecordsCount(), records, DefaultValue); + TSparsedArrayChunk chunk(defaultArray.GetRecordsCount(), records, defaultValue); + return std::shared_ptr<TSparsedArray>(new TSparsedArray(std::move(chunk), defaultValue, defaultArray.GetDataType())); } std::shared_ptr<arrow::Scalar> TSparsedArray::DoGetMaxScalar() const { - std::shared_ptr<arrow::Scalar> result; - for (auto&& i : Records) { - auto scalarCurrent = i.GetMaxScalar(); - if (!scalarCurrent) { - continue; - } - if (!result || ScalarCompare(result, scalarCurrent) < 0) { - result = scalarCurrent; - } - } - return result; + return Record.GetMaxScalar(); } ui32 TSparsedArray::GetLastIndex(const std::shared_ptr<arrow::RecordBatch>& batch) { @@ -105,11 +94,11 @@ TSparsedArrayChunk TSparsedArray::MakeDefaultChunk( it = SimpleBatchesCache.emplace(type->ToString(), NArrow::MakeEmptyBatch(BuildSchema(type))).first; AFL_VERIFY(it->second->ValidateFull().ok()); } - return TSparsedArrayChunk(0, recordsCount, it->second, defaultValue); + return TSparsedArrayChunk(recordsCount, it->second, defaultValue); } IChunkedArray::TLocalDataAddress TSparsedArrayChunk::GetChunk( - const std::optional<IChunkedArray::TCommonChunkAddress>& /*chunkCurrent*/, const ui64 position, const ui32 chunkIdx) const { + const std::optional<IChunkedArray::TCommonChunkAddress>& /*chunkCurrent*/, const ui64 position) const { const auto predCompare = [](const ui32 position, const TInternalChunkInfo& item) { return position < item.GetStartExt(); }; @@ -118,18 +107,18 @@ IChunkedArray::TLocalDataAddress TSparsedArrayChunk::GetChunk( --it; if (it->GetIsDefault()) { return IChunkedArray::TLocalDataAddress( - NArrow::TThreadSimpleArraysCache::Get(ColValue->type(), DefaultValue, it->GetSize()), StartPosition + it->GetStartExt(), chunkIdx); + NArrow::TThreadSimpleArraysCache::Get(ColValue->type(), DefaultValue, it->GetSize()), it->GetStartExt(), 0); } else { - return IChunkedArray::TLocalDataAddress(ColValue->Slice(it->GetStartInt(), it->GetSize()), StartPosition + it->GetStartExt(), chunkIdx); + return IChunkedArray::TLocalDataAddress(ColValue->Slice(it->GetStartInt(), it->GetSize()), it->GetStartExt(), 0); } } -TSparsedArrayChunk::TSparsedArrayChunk(const ui32 posStart, const ui32 recordsCount, const std::shared_ptr<arrow::RecordBatch>& records, - const std::shared_ptr<arrow::Scalar>& defaultValue) +TSparsedArrayChunk::TSparsedArrayChunk( + const ui32 recordsCount, const std::shared_ptr<arrow::RecordBatch>& records, const std::shared_ptr<arrow::Scalar>& defaultValue) : RecordsCount(recordsCount) - , StartPosition(posStart) , Records(records) , DefaultValue(defaultValue) { + AFL_VERIFY(Records); AFL_VERIFY(records->num_columns() == 2); ColIndex = Records->GetColumnByName("index"); AFL_VERIFY(ColIndex); @@ -192,9 +181,9 @@ std::shared_ptr<arrow::Scalar> TSparsedArrayChunk::GetScalar(const ui32 index) c ui32 TSparsedArrayChunk::GetFirstIndexNotDefault() const { if (UI32ColIndex->length()) { - return StartPosition + GetUI32ColIndex()->Value(0); + return GetUI32ColIndex()->Value(0); } else { - return StartPosition + GetRecordsCount(); + return GetRecordsCount(); } } @@ -210,6 +199,69 @@ std::shared_ptr<arrow::Scalar> TSparsedArrayChunk::GetMaxScalar() const { return DefaultValue; } +TSparsedArrayChunk TSparsedArrayChunk::ApplyFilter(const TColumnFilter& filter) const { + AFL_VERIFY(!filter.IsTotalAllowFilter()); + AFL_VERIFY(!filter.IsTotalDenyFilter()); + if (UI32ColIndex->length() == 0) { + return TSparsedArrayChunk(filter.GetFilteredCountVerified(), Records, DefaultValue); + } + AFL_VERIFY(filter.GetRecordsCountVerified() == RecordsCount)("filter", filter.GetRecordsCountVerified())("chunk", RecordsCount); + ui32 recordIndex = 0; + ui32 filterIntervalStart = 0; + ui32 skippedCount = 0; + ui32 filteredCount = 0; + bool currentAcceptance = filter.GetStartValue(); + TColumnFilter filterNew = TColumnFilter::BuildAllowFilter(); + auto indexesBuilder = NArrow::MakeBuilder(arrow::uint32()); + auto valuesBuilder = NArrow::MakeBuilder(ColValue->type()); + for (auto it = filter.GetFilter().begin(); it != filter.GetFilter().end(); ++it) { + for (; recordIndex < UI32ColIndex->length(); ++recordIndex) { + if (UI32ColIndex->Value(recordIndex) < filterIntervalStart) { + Y_ABORT_UNLESS(false); + } else if (UI32ColIndex->Value(recordIndex) < filterIntervalStart + *it) { + if (currentAcceptance) { + AFL_VERIFY(NArrow::Append<arrow::UInt32Type>(*indexesBuilder, UI32ColIndex->Value(recordIndex) - skippedCount)); + AFL_VERIFY(NArrow::Append(*valuesBuilder, *ColValue, recordIndex)); + ++filteredCount; + } + } else { + break; + } + } + if (!currentAcceptance) { + skippedCount += *it; + } + currentAcceptance = !currentAcceptance; + filterIntervalStart += *it; + } + AFL_VERIFY(filteredCount <= filter.GetFilteredCountVerified())("count", filteredCount)("filtered", filter.GetFilteredCountVerified()); + AFL_VERIFY(recordIndex == UI32ColIndex->length()); + auto indexesArr = NArrow::FinishBuilder(std::move(indexesBuilder)); + auto valuesArr = NArrow::FinishBuilder(std::move(valuesBuilder)); + std::shared_ptr<arrow::RecordBatch> result = + arrow::RecordBatch::Make(TSparsedArray::BuildSchema(ColValue->type()), filteredCount, { indexesArr, valuesArr }); + return TSparsedArrayChunk(filter.GetFilteredCountVerified(), result, DefaultValue); +} + +TSparsedArrayChunk TSparsedArrayChunk::Slice(const ui32 offset, const ui32 count) const { + AFL_VERIFY(offset + count <= RecordsCount)("offset", offset)("count", count)("records", RecordsCount); + std::optional<ui32> startPosition = NArrow::FindUpperOrEqualPosition(*UI32ColIndex, offset); + std::optional<ui32> finishPosition = NArrow::FindUpperOrEqualPosition(*UI32ColIndex, offset + count); + if (!startPosition || startPosition == finishPosition) { + return TSparsedArrayChunk(count, NArrow::MakeEmptyBatch(Records->schema(), 0), DefaultValue); + } else { + AFL_VERIFY(startPosition); + auto builder = NArrow::MakeBuilder(arrow::uint32()); + for (ui32 i = *startPosition; i < finishPosition.value_or(Records->num_rows()); ++i) { + NArrow::Append<arrow::UInt32Type>(*builder, UI32ColIndex->Value(i) - offset); + } + auto arrIndexes = NArrow::FinishBuilder(std::move(builder)); + auto arrValue = ColValue->Slice(*startPosition, finishPosition.value_or(Records->num_rows()) - *startPosition); + auto sliceRecords = arrow::RecordBatch::Make(Records->schema(), arrValue->length(), { arrIndexes, arrValue }); + return TSparsedArrayChunk(count, sliceRecords, DefaultValue); + } +} + void TSparsedArray::TBuilder::AddChunk(const ui32 recordsCount, const std::shared_ptr<arrow::RecordBatch>& data) { AFL_VERIFY(data); AFL_VERIFY(recordsCount); @@ -222,8 +274,9 @@ void TSparsedArray::TBuilder::AddChunk(const ui32 recordsCount, const std::share auto* arr = static_cast<const arrow::UInt32Array*>(data->column(0).get()); AFL_VERIFY(arr->Value(arr->length() - 1) < recordsCount)("val", arr->Value(arr->length() - 1))("count", recordsCount); } - Chunks.emplace_back(RecordsCount, recordsCount, data, DefaultValue); + Chunks.emplace_back(recordsCount, data, DefaultValue); RecordsCount += recordsCount; + AFL_VERIFY(Chunks.size() == 1); } void TSparsedArray::TBuilder::AddChunk( @@ -240,8 +293,9 @@ void TSparsedArray::TBuilder::AddChunk( AFL_VERIFY(arr->Value(arr->length() - 1) < recordsCount)("val", arr->Value(arr->length() - 1))("count", recordsCount); } Chunks.emplace_back( - RecordsCount, recordsCount, arrow::RecordBatch::Make(BuildSchema(Type), indexes->length(), { indexes, values }), DefaultValue); + recordsCount, arrow::RecordBatch::Make(BuildSchema(Type), indexes->length(), { indexes, values }), DefaultValue); RecordsCount += recordsCount; + AFL_VERIFY(Chunks.size() == 1); } } // namespace NKikimr::NArrow::NAccessor diff --git a/ydb/core/formats/arrow/accessor/sparsed/accessor.h b/ydb/core/formats/arrow/accessor/sparsed/accessor.h index 99ddc3b5ec0..66b38b846b4 100644 --- a/ydb/core/formats/arrow/accessor/sparsed/accessor.h +++ b/ydb/core/formats/arrow/accessor/sparsed/accessor.h @@ -1,22 +1,27 @@ #pragma once +#include <ydb/core/formats/arrow/accessor/abstract/accessor.h> #include <ydb/core/formats/arrow/arrow_helpers.h> #include <ydb/library/accessor/accessor.h> -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> #include <ydb/library/formats/arrow/size_calcer.h> +#include <ydb/library/formats/arrow/switch/switch_type.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/array/array_base.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/record_batch.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/type_fwd.h> +namespace NKikimr::NArrow { +class TColumnFilter; +} + namespace NKikimr::NArrow::NAccessor { class TSparsedArrayChunk: public TMoveOnly { private: YDB_READONLY(ui32, RecordsCount, 0); - YDB_READONLY(ui32, StartPosition, 0); YDB_READONLY_DEF(std::shared_ptr<arrow::RecordBatch>, Records); std::shared_ptr<arrow::Scalar> DefaultValue; + std::shared_ptr<arrow::DataType> DataType; std::shared_ptr<arrow::Array> ColIndex; const ui32* RawValues = nullptr; @@ -49,7 +54,7 @@ private: public: ui32 GetFinishPosition() const { - return StartPosition + RecordsCount; + return RecordsCount; } ui32 GetNotDefaultRecordsCount() const { @@ -66,11 +71,10 @@ public: std::shared_ptr<arrow::Scalar> GetScalar(const ui32 index) const; - IChunkedArray::TLocalDataAddress GetChunk( - const std::optional<IChunkedArray::TCommonChunkAddress>& chunkCurrent, const ui64 position, const ui32 chunkIdx) const; + IChunkedArray::TLocalDataAddress GetChunk(const std::optional<IChunkedArray::TCommonChunkAddress>& chunkCurrent, const ui64 position) const; - TSparsedArrayChunk(const ui32 posStart, const ui32 recordsCount, const std::shared_ptr<arrow::RecordBatch>& records, - const std::shared_ptr<arrow::Scalar>& defaultValue); + TSparsedArrayChunk( + const ui32 recordsCount, const std::shared_ptr<arrow::RecordBatch>& records, const std::shared_ptr<arrow::Scalar>& defaultValue); ui64 GetRawSize() const; @@ -86,94 +90,50 @@ public: return NArrow::GetArrayDataSize(ColValue); } - TSparsedArrayChunk Slice(const ui32 newStart, const ui32 offset, const ui32 count) const { - AFL_VERIFY(offset + count <= RecordsCount)("offset", offset)("count", count)("records", RecordsCount); - std::optional<ui32> startPosition = NArrow::FindUpperOrEqualPosition(*UI32ColIndex, offset); - std::optional<ui32> finishPosition = NArrow::FindUpperOrEqualPosition(*UI32ColIndex, offset + count); - if (!startPosition || startPosition == finishPosition) { - return TSparsedArrayChunk(newStart, count, NArrow::MakeEmptyBatch(Records->schema(), 0), DefaultValue); - } else { - AFL_VERIFY(startPosition); - auto builder = NArrow::MakeBuilder(arrow::uint32()); - for (ui32 i = *startPosition; i < finishPosition.value_or(Records->num_rows()); ++i) { - NArrow::Append<arrow::UInt32Type>(*builder, UI32ColIndex->Value(i) - offset); - } - auto arrIndexes = NArrow::FinishBuilder(std::move(builder)); - auto arrValue = ColValue->Slice(*startPosition, finishPosition.value_or(Records->num_rows()) - *startPosition); - auto sliceRecords = arrow::RecordBatch::Make(Records->schema(), arrValue->length(), { arrIndexes, arrValue }); - return TSparsedArrayChunk(newStart, count, sliceRecords, DefaultValue); - } - } + TSparsedArrayChunk ApplyFilter(const TColumnFilter& filter) const; + TSparsedArrayChunk Slice(const ui32 offset, const ui32 count) const; }; class TSparsedArray: public IChunkedArray { private: using TBase = IChunkedArray; YDB_READONLY_DEF(std::shared_ptr<arrow::Scalar>, DefaultValue); - std::vector<TSparsedArrayChunk> Records; + TSparsedArrayChunk Record; + friend class TSparsedArrayChunk; protected: virtual std::shared_ptr<arrow::Scalar> DoGetMaxScalar() const override; virtual ui32 DoGetNullsCount() const override { - ui32 result = 0; - for (auto&& i : Records) { - result += i.GetNullsCount(); - } - return result; + return Record.GetNullsCount(); } virtual ui32 DoGetValueRawBytes() const override { - ui32 result = 0; - for (auto&& i : Records) { - result += i.GetValueRawBytes(); - } - return result; + return Record.GetValueRawBytes(); } virtual std::shared_ptr<IChunkedArray> DoISlice(const ui32 offset, const ui32 count) const override { TBuilder builder(DefaultValue, GetDataType()); - ui32 newStart = 0; - for (ui32 i = 0; i < Records.size(); ++i) { - if (Records[i].GetStartPosition() + Records[i].GetRecordsCount() <= offset) { - continue; - } - if (offset + count <= Records[i].GetStartPosition()) { - continue; - } - const ui32 chunkStart = (offset < Records[i].GetStartPosition()) ? 0 : (offset - Records[i].GetStartPosition()); - const ui32 chunkCount = (offset + count <= Records[i].GetFinishPosition()) - ? (offset + count - Records[i].GetStartPosition() - chunkStart) - : (Records[i].GetFinishPosition() - chunkStart); - builder.AddChunk(Records[i].Slice(newStart, chunkStart, chunkCount)); - newStart += chunkCount; - } + builder.AddChunk(Record.Slice(offset, count)); + return builder.Finish(); + } + + virtual std::shared_ptr<IChunkedArray> DoApplyFilter(const TColumnFilter& filter) const override { + TBuilder builder(DefaultValue, GetDataType()); + builder.AddChunk(Record.ApplyFilter(filter)); return builder.Finish(); } virtual TLocalDataAddress DoGetLocalData(const std::optional<TCommonChunkAddress>& chunkCurrent, const ui64 position) const override { - ui32 currentIdx = 0; - for (ui32 i = 0; i < Records.size(); ++i) { - if (currentIdx <= position && position < currentIdx + Records[i].GetRecordsCount()) { - return Records[i].GetChunk(chunkCurrent, position - currentIdx, i); - } - currentIdx += Records[i].GetRecordsCount(); - } - AFL_VERIFY(false); - return TLocalDataAddress(nullptr, 0, 0); + return Record.GetChunk(chunkCurrent, position); } virtual std::optional<ui64> DoGetRawSize() const override { - ui64 bytes = 0; - for (auto&& i : Records) { - bytes += i.GetRawSize(); - } - return bytes; + return Record.GetRawSize(); } - TSparsedArray(std::vector<TSparsedArrayChunk>&& data, const std::shared_ptr<arrow::Scalar>& defaultValue, - const std::shared_ptr<arrow::DataType>& type, const ui32 recordsCount) - : TBase(recordsCount, EType::SparsedArray, type) + TSparsedArray(TSparsedArrayChunk&& data, const std::shared_ptr<arrow::Scalar>& defaultValue, const std::shared_ptr<arrow::DataType>& type) + : TBase(data.GetRecordsCount(), EType::SparsedArray, type) , DefaultValue(defaultValue) - , Records(std::move(data)) { + , Record(std::move(data)) { } static ui32 GetLastIndex(const std::shared_ptr<arrow::RecordBatch>& batch); @@ -188,34 +148,25 @@ protected: const std::shared_ptr<arrow::Scalar>& defaultValue, const std::shared_ptr<arrow::DataType>& type, const ui32 recordsCount); public: - TSparsedArray(const IChunkedArray& defaultArray, const std::shared_ptr<arrow::Scalar>& defaultValue); + static std::shared_ptr<TSparsedArray> Make(const IChunkedArray& defaultArray, const std::shared_ptr<arrow::Scalar>& defaultValue); TSparsedArray(const std::shared_ptr<arrow::Scalar>& defaultValue, const std::shared_ptr<arrow::DataType>& type, const ui32 recordsCount) : TBase(recordsCount, EType::SparsedArray, type) - , DefaultValue(defaultValue) { - Records.emplace_back(MakeDefaultChunk(defaultValue, type, recordsCount)); + , DefaultValue(defaultValue) + , Record(MakeDefaultChunk(defaultValue, type, recordsCount)) { } - virtual std::shared_ptr<arrow::Scalar> DoGetScalar(const ui32 index) const override { - auto& chunk = GetSparsedChunk(index); - return chunk.GetScalar(index - chunk.GetStartPosition()); + const TSparsedArrayChunk& GetSparsedChunk(const ui64 position) const { + AFL_VERIFY(position < Record.GetRecordsCount()); + return Record; } - std::shared_ptr<arrow::RecordBatch> GetRecordBatchVerified() const { - AFL_VERIFY(Records.size() == 1)("size", Records.size()); - return Records.front().GetRecords(); + virtual std::shared_ptr<arrow::Scalar> DoGetScalar(const ui32 index) const override { + return Record.GetScalar(index); } - const TSparsedArrayChunk& GetSparsedChunk(const ui64 position) const { - const auto pred = [](const ui64 position, const TSparsedArrayChunk& item) { - return position < item.GetStartPosition(); - }; - auto it = std::upper_bound(Records.begin(), Records.end(), position, pred); - AFL_VERIFY(it != Records.begin()); - --it; - AFL_VERIFY(position < it->GetStartPosition() + it->GetRecordsCount()); - AFL_VERIFY(it->GetStartPosition() <= position); - return *it; + std::shared_ptr<arrow::RecordBatch> GetRecordBatchVerified() const { + return Record.GetRecords(); } template <class TDataType> @@ -272,10 +223,12 @@ public: void AddChunk(TSparsedArrayChunk&& chunk) { RecordsCount += chunk.GetRecordsCount(); Chunks.emplace_back(std::move(chunk)); + AFL_VERIFY(Chunks.size() == 1); } std::shared_ptr<TSparsedArray> Finish() { - return std::shared_ptr<TSparsedArray>(new TSparsedArray(std::move(Chunks), DefaultValue, Type, RecordsCount)); + AFL_VERIFY(Chunks.size() == 1); + return std::shared_ptr<TSparsedArray>(new TSparsedArray(std::move(Chunks.front()), DefaultValue, Type)); } }; }; diff --git a/ydb/core/formats/arrow/accessor/sparsed/constructor.cpp b/ydb/core/formats/arrow/accessor/sparsed/constructor.cpp index 7a76ce60f23..3d4a3574afd 100644 --- a/ydb/core/formats/arrow/accessor/sparsed/constructor.cpp +++ b/ydb/core/formats/arrow/accessor/sparsed/constructor.cpp @@ -43,7 +43,7 @@ TString TConstructor::DoSerializeToString(const std::shared_ptr<IChunkedArray>& TConclusion<std::shared_ptr<IChunkedArray>> TConstructor::DoConstruct( const std::shared_ptr<IChunkedArray>& originalArray, const TChunkConstructionData& externalInfo) const { AFL_VERIFY(originalArray); - return std::make_shared<TSparsedArray>(*originalArray, externalInfo.GetDefaultValue()); + return TSparsedArray::Make(*originalArray, externalInfo.GetDefaultValue()); } } // namespace NKikimr::NArrow::NAccessor::NSparsed diff --git a/ydb/core/formats/arrow/accessor/sparsed/constructor.h b/ydb/core/formats/arrow/accessor/sparsed/constructor.h index 6591b31c769..45493c6d649 100644 --- a/ydb/core/formats/arrow/accessor/sparsed/constructor.h +++ b/ydb/core/formats/arrow/accessor/sparsed/constructor.h @@ -1,13 +1,13 @@ #pragma once #include <ydb/core/formats/arrow/accessor/abstract/constructor.h> - -#include <ydb/library/formats/arrow/accessor/common/const.h> +#include <ydb/core/formats/arrow/accessor/common/const.h> namespace NKikimr::NArrow::NAccessor::NSparsed { class TConstructor: public IConstructor { private: using TBase = IConstructor; + public: static TString GetClassNameStatic() { return TGlobalConst::SparsedDataAccessorName; diff --git a/ydb/core/formats/arrow/accessor/sparsed/request.h b/ydb/core/formats/arrow/accessor/sparsed/request.h index 4be2d897b09..205949bca97 100644 --- a/ydb/core/formats/arrow/accessor/sparsed/request.h +++ b/ydb/core/formats/arrow/accessor/sparsed/request.h @@ -1,6 +1,6 @@ #pragma once #include <ydb/core/formats/arrow/accessor/abstract/request.h> -#include <ydb/library/formats/arrow/accessor/common/const.h> +#include <ydb/core/formats/arrow/accessor/common/const.h> namespace NKikimr::NArrow::NAccessor::NSparsed { diff --git a/ydb/core/formats/arrow/accessor/sparsed/ut/ut_sparsed.cpp b/ydb/core/formats/arrow/accessor/sparsed/ut/ut_sparsed.cpp index 781424597a3..cbcad77a31a 100644 --- a/ydb/core/formats/arrow/accessor/sparsed/ut/ut_sparsed.cpp +++ b/ydb/core/formats/arrow/accessor/sparsed/ut/ut_sparsed.cpp @@ -1,4 +1,5 @@ #include <ydb/core/formats/arrow/accessor/sparsed/accessor.h> +#include <ydb/core/formats/arrow/arrow_filter.h> #include <library/cpp/testing/unittest/registar.h> @@ -81,4 +82,55 @@ Y_UNIT_TEST_SUITE(SparsedArrayAccessor) { AFL_VERIFY(arrString == "[[\"abcd6\"],[null],[\"abcde8\"]]")("string", arrString); } } + Y_UNIT_TEST(FiltersDef) { + TSparsedArray::TSparsedBuilder<arrow::StringType> builder(nullptr, 10, 0); + builder.AddRecord(5, "abc5"); + builder.AddRecord(6, "abcd6"); + builder.AddRecord(8, "abcde8"); + auto arr = builder.Finish(10); + { + TColumnFilter filter = TColumnFilter::BuildAllowFilter(); + filter.Add(true, 1); + filter.Add(false, 1); + filter.Add(true, 1); + filter.Add(false, 1); + filter.Add(true, 1); + filter.Add(false, 5); + auto arrFiltered = filter.Apply(arr)->GetChunkedArray(); + auto arrSlice = PrepareToCompare(arrFiltered->ToString()); + AFL_VERIFY(PrepareToCompare(arrFiltered->ToString()) == R"([[null,null,null]])")("string", PrepareToCompare(arrFiltered->ToString())); + } + { + TColumnFilter filter = TColumnFilter::BuildAllowFilter(); + filter.Add(false, 1); + filter.Add(true, 1); + filter.Add(false, 1); + filter.Add(true, 1); + filter.Add(false, 1); + filter.Add(true, 5); + auto arrFiltered = filter.Apply(arr)->GetChunkedArray(); + auto arrSlice = PrepareToCompare(arrFiltered->ToString()); + AFL_VERIFY(PrepareToCompare(arrFiltered->ToString()) == R"([[null,null],["abc5","abcd6"],[null],["abcde8"],[null]])")( + "string", PrepareToCompare(arrFiltered->ToString())); + } + { + TColumnFilter filter = TColumnFilter::BuildAllowFilter(); + filter.Add(true, 6); + filter.Add(false, 4); + auto arrFiltered = filter.Apply(arr)->GetChunkedArray(); + auto arrSlice = PrepareToCompare(arrFiltered->ToString()); + AFL_VERIFY(PrepareToCompare(arrFiltered->ToString()) == R"([[null,null,null,null,null],["abc5"]])")( + "string", PrepareToCompare(arrFiltered->ToString())); + } + { + TColumnFilter filter = TColumnFilter::BuildAllowFilter(); + filter.Add(false, 5); + filter.Add(true, 1); + filter.Add(false, 4); + auto arrFiltered = filter.Apply(arr)->GetChunkedArray(); + auto arrSlice = PrepareToCompare(arrFiltered->ToString()); + AFL_VERIFY(PrepareToCompare(arrFiltered->ToString()) == R"([["abc5"]])")( + "string", PrepareToCompare(arrFiltered->ToString())); + } + } }; diff --git a/ydb/core/formats/arrow/accessor/sparsed/ut/ya.make b/ydb/core/formats/arrow/accessor/sparsed/ut/ya.make index 276aaddeb88..29e11a1035a 100644 --- a/ydb/core/formats/arrow/accessor/sparsed/ut/ya.make +++ b/ydb/core/formats/arrow/accessor/sparsed/ut/ya.make @@ -7,6 +7,7 @@ PEERDIR( ydb/core/formats/arrow/accessor/plain ydb/core/formats/arrow yql/essentials/public/udf/service/stub + ydb/core/formats/arrow ) YQL_LAST_ABI_VERSION() diff --git a/ydb/core/formats/arrow/accessor/sparsed/ya.make b/ydb/core/formats/arrow/accessor/sparsed/ya.make index 93d6886d6a9..28e34bc073d 100644 --- a/ydb/core/formats/arrow/accessor/sparsed/ya.make +++ b/ydb/core/formats/arrow/accessor/sparsed/ya.make @@ -7,7 +7,7 @@ PEERDIR( ydb/core/formats/arrow/save_load ydb/core/formats/arrow/serializer ydb/core/formats/arrow/splitter - ydb/library/formats/arrow/accessor/common + ydb/core/formats/arrow/accessor/common ) SRCS( diff --git a/ydb/core/formats/arrow/accessor/sub_columns/accessor.h b/ydb/core/formats/arrow/accessor/sub_columns/accessor.h index 1b4cb7fec95..af334b53ae2 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/accessor.h +++ b/ydb/core/formats/arrow/accessor/sub_columns/accessor.h @@ -5,12 +5,13 @@ #include "others_storage.h" #include "settings.h" +#include <ydb/core/formats/arrow/accessor/abstract/accessor.h> +#include <ydb/core/formats/arrow/accessor/common/chunk_data.h> +#include <ydb/core/formats/arrow/arrow_filter.h> #include <ydb/core/formats/arrow/arrow_helpers.h> #include <ydb/core/formats/arrow/common/container.h> #include <ydb/library/accessor/accessor.h> -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> -#include <ydb/library/formats/arrow/accessor/common/chunk_data.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/array/array_base.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/record_batch.h> @@ -24,6 +25,7 @@ private: NSubColumns::TColumnsData ColumnsData; NSubColumns::TOthersData OthersData; const NSubColumns::TSettings Settings; + TString SourceDeserializationString; protected: virtual NJson::TJsonValue DoDebugJson() const override { @@ -50,12 +52,22 @@ protected: virtual std::optional<ui64> DoGetRawSize() const override { return ColumnsData.GetRawSize() + OthersData.GetRawSize(); } + virtual std::shared_ptr<IChunkedArray> DoApplyFilter(const TColumnFilter& filter) const override { + return std::make_shared<TSubColumnsArray>(ColumnsData.ApplyFilter(filter), OthersData.ApplyFilter(filter, Settings), GetDataType(), + filter.GetFilteredCountVerified(), Settings); + } + virtual std::shared_ptr<IChunkedArray> DoISlice(const ui32 offset, const ui32 count) const override { return std::make_shared<TSubColumnsArray>( ColumnsData.Slice(offset, count), OthersData.Slice(offset, count, Settings), GetDataType(), count, Settings); } public: + void StoreSourceString(const TString& sourceDeserializationString) { + AFL_VERIFY(!SourceDeserializationString); + SourceDeserializationString = sourceDeserializationString; + } + std::shared_ptr<NSubColumns::TReadIteratorOrderedKeys> BuildOrderedIterator() const { return std::make_shared<NSubColumns::TReadIteratorOrderedKeys>(ColumnsData, OthersData); } diff --git a/ydb/core/formats/arrow/accessor/sub_columns/columns_storage.cpp b/ydb/core/formats/arrow/accessor/sub_columns/columns_storage.cpp index 4a98bbfb7b1..401b040d582 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/columns_storage.cpp +++ b/ydb/core/formats/arrow/accessor/sub_columns/columns_storage.cpp @@ -1,17 +1,48 @@ #include "columns_storage.h" +#include <ydb/core/formats/arrow/arrow_filter.h> + namespace NKikimr::NArrow::NAccessor::NSubColumns { TColumnsData TColumnsData::Slice(const ui32 offset, const ui32 count) const { - auto sliceRecords = Records->Slice(offset, count); - if (sliceRecords.GetRecordsCount()) { + auto records = Records->Slice(offset, count); + if (records.GetRecordsCount()) { + TDictStats::TBuilder builder; + ui32 idx = 0; + std::vector<ui32> indexesToRemove; + for (auto&& i : records.GetColumns()) { + AFL_VERIFY(Stats.GetColumnName(idx) == records.GetSchema()->field(idx)->name()); + if (i->GetRecordsCount() > i->GetNullsCount()) { + builder.Add(Stats.GetColumnName(idx), i->GetRecordsCount() - i->GetNullsCount(), i->GetValueRawBytes(), i->GetType()); + } else { + indexesToRemove.emplace_back(idx); + } + ++idx; + } + return TColumnsData(builder.Finish(), std::make_shared<TGeneralContainer>(std::move(records))); + + } else { + return TColumnsData(TDictStats::BuildEmpty(), std::make_shared<TGeneralContainer>(0)); + } +} + +TColumnsData TColumnsData::ApplyFilter(const TColumnFilter& filter) const { + auto records = Records; + AFL_VERIFY(filter.Apply(records)); + if (records->GetRecordsCount()) { TDictStats::TBuilder builder; ui32 idx = 0; - for (auto&& i : sliceRecords.GetColumns()) { - AFL_VERIFY(Stats.GetColumnName(idx) == sliceRecords.GetSchema()->field(idx)->name()); - builder.Add(Stats.GetColumnName(idx), i->GetRecordsCount() - i->GetNullsCount(), i->GetValueRawBytes(), i->GetType()); + std::vector<ui32> indexesToRemove; + for (auto&& i : records->GetColumns()) { + AFL_VERIFY(Stats.GetColumnName(idx) == records->GetSchema()->field(idx)->name()); + if (i->GetRecordsCount() > i->GetNullsCount()) { + builder.Add(Stats.GetColumnName(idx), i->GetRecordsCount() - i->GetNullsCount(), i->GetValueRawBytes(), i->GetType()); + } else { + indexesToRemove.emplace_back(idx); + } ++idx; } - return TColumnsData(builder.Finish(), std::make_shared<TGeneralContainer>(std::move(sliceRecords))); + records->DeleteFieldsByIndex(indexesToRemove); + return TColumnsData(builder.Finish(), std::move(records)); } else { return TColumnsData(TDictStats::BuildEmpty(), std::make_shared<TGeneralContainer>(0)); diff --git a/ydb/core/formats/arrow/accessor/sub_columns/columns_storage.h b/ydb/core/formats/arrow/accessor/sub_columns/columns_storage.h index 65587ee49d8..35557e8a766 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/columns_storage.h +++ b/ydb/core/formats/arrow/accessor/sub_columns/columns_storage.h @@ -33,6 +33,8 @@ public: return result; } + TColumnsData ApplyFilter(const TColumnFilter& filter) const; + TColumnsData Slice(const ui32 offset, const ui32 count) const; static TColumnsData BuildEmpty(const ui32 recordsCount) { diff --git a/ydb/core/formats/arrow/accessor/sub_columns/constructor.cpp b/ydb/core/formats/arrow/accessor/sub_columns/constructor.cpp index ffc5b155819..9a624c9f29f 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/constructor.cpp +++ b/ydb/core/formats/arrow/accessor/sub_columns/constructor.cpp @@ -54,7 +54,7 @@ TConclusion<std::shared_ptr<IChunkedArray>> TConstructor::DoDeserializeFromStrin std::shared_ptr<TColumnLoader> columnLoader = std::make_shared<TColumnLoader>( externalInfo.GetDefaultSerializer(), columnStats.GetAccessorConstructor(i), schema->field(i), nullptr, 0); std::vector<TDeserializeChunkedArray::TChunk> chunks = { TDeserializeChunkedArray::TChunk( - externalInfo.GetRecordsCount(), originalData.substr(currentIndex, proto.GetKeyColumns(i).GetSize())) }; + externalInfo.GetRecordsCount(), TStringBuf(originalData.data() + currentIndex, proto.GetKeyColumns(i).GetSize())) }; columns.emplace_back(std::make_shared<TDeserializeChunkedArray>(externalInfo.GetRecordsCount(), columnLoader, std::move(chunks), true)); currentIndex += proto.GetKeyColumns(i).GetSize(); } @@ -71,7 +71,7 @@ TConclusion<std::shared_ptr<IChunkedArray>> TConstructor::DoDeserializeFromStrin std::shared_ptr<TColumnLoader> columnLoader = std::make_shared<TColumnLoader>( externalInfo.GetDefaultSerializer(), std::make_shared<NPlain::TConstructor>(), schema->field(i), nullptr, 0); std::vector<TDeserializeChunkedArray::TChunk> chunks = { TDeserializeChunkedArray::TChunk( - proto.GetOtherRecordsCount(), originalData.substr(currentIndex, proto.GetOtherColumns(i).GetSize())) }; + proto.GetOtherRecordsCount(), TStringBuf(originalData.data() + currentIndex, proto.GetOtherColumns(i).GetSize())) }; columns.emplace_back(std::make_shared<TDeserializeChunkedArray>(proto.GetOtherRecordsCount(), columnLoader, std::move(chunks), true)); currentIndex += proto.GetOtherColumns(i).GetSize(); } @@ -79,8 +79,10 @@ TConclusion<std::shared_ptr<IChunkedArray>> TConstructor::DoDeserializeFromStrin otherData = TOthersData(otherStats, otherKeysContainer); } TColumnsData columnData(columnStats, columnKeysContainer); - return std::make_shared<TSubColumnsArray>( + auto result = std::make_shared<TSubColumnsArray>( std::move(columnData), std::move(otherData), externalInfo.GetColumnType(), externalInfo.GetRecordsCount(), Settings); + result->StoreSourceString(originalData); + return result; } NKikimrArrowAccessorProto::TConstructor TConstructor::DoSerializeToProto() const { diff --git a/ydb/core/formats/arrow/accessor/sub_columns/constructor.h b/ydb/core/formats/arrow/accessor/sub_columns/constructor.h index 1f2c1d20b4a..f5e7b09dd92 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/constructor.h +++ b/ydb/core/formats/arrow/accessor/sub_columns/constructor.h @@ -2,8 +2,7 @@ #include "data_extractor.h" #include <ydb/core/formats/arrow/accessor/abstract/constructor.h> - -#include <ydb/library/formats/arrow/accessor/common/const.h> +#include <ydb/core/formats/arrow/accessor/common/const.h> namespace NKikimr::NArrow::NAccessor::NSubColumns { diff --git a/ydb/core/formats/arrow/accessor/sub_columns/data_extractor.h b/ydb/core/formats/arrow/accessor/sub_columns/data_extractor.h index 272414d0a43..bf22e47b1f8 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/data_extractor.h +++ b/ydb/core/formats/arrow/accessor/sub_columns/data_extractor.h @@ -1,10 +1,9 @@ #pragma once #include "direct_builder.h" +#include <ydb/core/formats/arrow/accessor/abstract/accessor.h> #include <ydb/core/formats/arrow/arrow_helpers.h> -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> - #include <contrib/libs/apache/arrow/cpp/src/arrow/array/builder_base.h> namespace NKikimr::NArrow::NAccessor::NSubColumns { @@ -28,4 +27,4 @@ private: public: }; -} // namespace NKikimr::NArrow::NAccessor +} // namespace NKikimr::NArrow::NAccessor::NSubColumns diff --git a/ydb/core/formats/arrow/accessor/sub_columns/direct_builder.h b/ydb/core/formats/arrow/accessor/sub_columns/direct_builder.h index 33430589cc8..5e7365b271c 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/direct_builder.h +++ b/ydb/core/formats/arrow/accessor/sub_columns/direct_builder.h @@ -3,10 +3,9 @@ #include "settings.h" #include "stats.h" +#include <ydb/core/formats/arrow/accessor/abstract/accessor.h> #include <ydb/core/formats/arrow/arrow_helpers.h> -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> - #include <contrib/libs/apache/arrow/cpp/src/arrow/array/builder_base.h> namespace NKikimr::NArrow::NAccessor { @@ -60,8 +59,7 @@ private: public: TDataBuilder(const std::shared_ptr<arrow::DataType>& type, const TSettings& settings) : Type(type) - , Settings(settings) - { + , Settings(settings) { } void StartNextRecord() { @@ -134,7 +132,7 @@ public: TDictStats BuildStats(const std::vector<TColumnElements*>& keys, const TSettings& settings, const ui32 recordsCount) const { auto builder = TDictStats::MakeBuilder(); for (auto&& i : keys) { - builder.Add(i->GetKeyName(), i->GetRecordIndexes().size(), i->GetDataSize(), + builder.Add(i->GetKeyName(), i->GetRecordIndexes().size(), i->GetDataSize(), settings.IsSparsed(i->GetRecordIndexes().size(), recordsCount) ? IChunkedArray::EType::SparsedArray : IChunkedArray::EType::Array); } diff --git a/ydb/core/formats/arrow/accessor/sub_columns/others_storage.cpp b/ydb/core/formats/arrow/accessor/sub_columns/others_storage.cpp index f06839e334a..3444d72c167 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/others_storage.cpp +++ b/ydb/core/formats/arrow/accessor/sub_columns/others_storage.cpp @@ -76,6 +76,97 @@ TOthersData TOthersData::TBuilderWithStats::Finish(const TFinishContext& finishC return TOthersData(*resultStats, std::make_shared<TGeneralContainer>(arrow::RecordBatch::Make(GetSchema(), RecordsCount, arrays))); } +class TUsedKeysCollection { +private: + std::map<ui32, TDictStats::TRTStats> UsedKeys; + const TDictStats& Stats; + std::vector<ui32> OriginalKeys; + +public: + TUsedKeysCollection(const TDictStats& stats) + : Stats(stats) { + } + + std::vector<ui32> BuildDecoder() const { + std::vector<ui32> keyIndexDecoder; + if (!UsedKeys.size()) { + return keyIndexDecoder; + } + keyIndexDecoder.resize(UsedKeys.rbegin()->first + 1, Max<ui32>()); + ui32 idx = 0; + for (auto&& i : UsedKeys) { + keyIndexDecoder[i.first] = idx++; + } + return keyIndexDecoder; + } + + void AddKeyInfo(const ui32 keyIndex, const std::string_view value) { + auto itUsedKey = UsedKeys.find(keyIndex); + if (itUsedKey == UsedKeys.end()) { + itUsedKey = UsedKeys.emplace(keyIndex, Stats.GetColumnName(keyIndex)).first; + } + itUsedKey->second.AddValue(value); + } + + TDictStats BuildStats(const TSettings& settings, const ui32 recordsCount) const { + TDictStats::TBuilder statBuilder; + for (auto&& i : UsedKeys) { + statBuilder.Add( + i.second.GetKeyName(), i.second.GetRecordsCount(), i.second.GetDataSize(), i.second.GetAccessorType(settings, recordsCount)); + } + return statBuilder.Finish(); + } +}; + +TOthersData TOthersData::ApplyFilter(const TColumnFilter& filter, const TSettings& settings) const { + TOthersData::TIterator itOthersData = BuildIterator(); + bool currentAcceptance = filter.GetStartValue(); + ui32 filterIntervalStart = 0; + ui32 shiftSkippedCount = 0; + TColumnFilter newFilter = TColumnFilter::BuildAllowFilter(); + auto recordIndexBuilder = NArrow::MakeBuilder(arrow::uint32()); + auto valuesBuilder = NArrow::MakeBuilder(arrow::utf8()); + TUsedKeysCollection usedKeys(Stats); + std::vector<ui32> originalKeys; + for (auto it = filter.GetFilter().begin(); itOthersData.IsValid() && it != filter.GetFilter().end(); ++it) { + for (; itOthersData.IsValid(); itOthersData.Next()) { + if (itOthersData.GetRecordIndex() < filterIntervalStart) { + AFL_VERIFY(false); + } else if (itOthersData.GetRecordIndex() < filterIntervalStart + *it) { + if (currentAcceptance) { + AFL_VERIFY(shiftSkippedCount <= itOthersData.GetRecordIndex()); + const ui32 filteredRecordIndex = itOthersData.GetRecordIndex() - shiftSkippedCount; + NArrow::Append<arrow::UInt32Type>(*recordIndexBuilder, filteredRecordIndex); + NArrow::Append<arrow::StringType>(*valuesBuilder, arrow::util::string_view(itOthersData.GetValue().data(), itOthersData.GetValue().size())); + originalKeys.emplace_back(itOthersData.GetKeyIndex()); + usedKeys.AddKeyInfo(itOthersData.GetKeyIndex(), itOthersData.GetValue()); + } + } else { + break; + } + } + if (!currentAcceptance) { + shiftSkippedCount += *it; + } + currentAcceptance = !currentAcceptance; + filterIntervalStart += *it; + } + auto stats = usedKeys.BuildStats(settings, filter.GetFilteredCountVerified()); + const std::vector<ui32> decoder = usedKeys.BuildDecoder(); + auto keyIndexBuilder = NArrow::MakeBuilder(arrow::uint32()); + for (auto&& i : originalKeys) { + NArrow::Append<arrow::UInt32Type>(*keyIndexBuilder, decoder[i]); + } + auto recordIndexes = NArrow::FinishBuilder(std::move(recordIndexBuilder)); + auto keyIndexes = NArrow::FinishBuilder(std::move(keyIndexBuilder)); + auto values = NArrow::FinishBuilder(std::move(valuesBuilder)); + + std::vector<std::shared_ptr<IChunkedArray>> arrays = { std::make_shared<TTrivialArray>(recordIndexes), + std::make_shared<TTrivialArray>(keyIndexes), std::make_shared<TTrivialArray>(values) }; + auto records = std::make_shared<TGeneralContainer>(GetSchema(), std::move(arrays)); + return TOthersData(stats, records); +} + TOthersData TOthersData::Slice(const ui32 offset, const ui32 count, const TSettings& settings) const { AFL_VERIFY(Records->GetColumnsCount() == 3); if (!count) { @@ -87,30 +178,15 @@ TOthersData TOthersData::Slice(const ui32 offset, const ui32 count, const TSetti if (!startPosition || startPosition == finishPosition) { return TOthersData(TDictStats::BuildEmpty(), std::make_shared<TGeneralContainer>(0)); } - std::map<ui32, TDictStats::TRTStats> usedKeys; + TUsedKeysCollection usedKeys(Stats); { itOthersData.MoveToPosition(*startPosition); for (; itOthersData.IsValid() && itOthersData.GetRecordIndex() < offset + count; itOthersData.Next()) { - auto itUsedKey = usedKeys.find(itOthersData.GetKeyIndex()); - if (itUsedKey == usedKeys.end()) { - itUsedKey = usedKeys.emplace(itOthersData.GetKeyIndex(), Stats.GetColumnName(itOthersData.GetKeyIndex())).first; - } - itUsedKey->second.AddValue(itOthersData.GetValue()); + usedKeys.AddKeyInfo(itOthersData.GetKeyIndex(), itOthersData.GetValue()); } } - std::vector<ui32> keyIndexDecoder; - if (usedKeys.size()) { - keyIndexDecoder.resize(usedKeys.rbegin()->first + 1, Max<ui32>()); - ui32 idx = 0; - for (auto&& i : usedKeys) { - keyIndexDecoder[i.first] = idx++; - } - } - TDictStats::TBuilder statBuilder; - for (auto&& i : usedKeys) { - statBuilder.Add(i.second.GetKeyName(), i.second.GetRecordsCount(), i.second.GetDataSize(), i.second.GetAccessorType(settings, count)); - } - TDictStats sliceStats = statBuilder.Finish(); + const std::vector<ui32> keyIndexDecoder = usedKeys.BuildDecoder(); + const TDictStats sliceStats = usedKeys.BuildStats(settings, count); { auto recordIndexBuilder = NArrow::MakeBuilder(arrow::uint32()); diff --git a/ydb/core/formats/arrow/accessor/sub_columns/others_storage.h b/ydb/core/formats/arrow/accessor/sub_columns/others_storage.h index 765337c2cdc..7de486a31d1 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/others_storage.h +++ b/ydb/core/formats/arrow/accessor/sub_columns/others_storage.h @@ -28,6 +28,7 @@ public: } TOthersData Slice(const ui32 offset, const ui32 count, const TSettings& settings) const; + TOthersData ApplyFilter(const TColumnFilter& filter, const TSettings& settings) const; static TOthersData BuildEmpty(); diff --git a/ydb/core/formats/arrow/accessor/sub_columns/request.h b/ydb/core/formats/arrow/accessor/sub_columns/request.h index 264edd95ac2..923befcab81 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/request.h +++ b/ydb/core/formats/arrow/accessor/sub_columns/request.h @@ -2,8 +2,7 @@ #include "settings.h" #include <ydb/core/formats/arrow/accessor/abstract/request.h> - -#include <ydb/library/formats/arrow/accessor/common/const.h> +#include <ydb/core/formats/arrow/accessor/common/const.h> namespace NKikimr::NArrow::NAccessor::NSubColumns { diff --git a/ydb/core/formats/arrow/accessor/sub_columns/settings.h b/ydb/core/formats/arrow/accessor/sub_columns/settings.h index c4ee71dcb38..46808e42b07 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/settings.h +++ b/ydb/core/formats/arrow/accessor/sub_columns/settings.h @@ -1,7 +1,7 @@ #pragma once +#include <ydb/core/formats/arrow/accessor/abstract/accessor.h> #include <ydb/core/formats/arrow/arrow_helpers.h> -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> #include <ydb/library/formats/arrow/protos/accessor.pb.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/array/builder_base.h> diff --git a/ydb/core/formats/arrow/accessor/sub_columns/stats.cpp b/ydb/core/formats/arrow/accessor/sub_columns/stats.cpp index d854ec278fa..fc86ca63881 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/stats.cpp +++ b/ydb/core/formats/arrow/accessor/sub_columns/stats.cpp @@ -101,7 +101,7 @@ TConstructorContainer TDictStats::GetAccessorConstructor(const ui32 columnIndex) case IChunkedArray::EType::CompositeChunkedArray: case IChunkedArray::EType::SubColumnsArray: case IChunkedArray::EType::ChunkedArray: - AFL_VERIFY(false); + AFL_VERIFY(false)("type", GetAccessorType(columnIndex)); return TConstructorContainer(); } } @@ -141,6 +141,7 @@ void TDictStats::TBuilder::Add(const TString& name, const ui32 recordsCount, con AFL_VERIFY(*LastKeyName < name)("last", LastKeyName)("name", name); } AFL_VERIFY(recordsCount); + AFL_VERIFY(accessorType == IChunkedArray::EType::Array || accessorType == IChunkedArray::EType::SparsedArray)("type", accessorType); TStatusValidator::Validate(Names->Append(name.data(), name.size())); TStatusValidator::Validate(Records->Append(recordsCount)); TStatusValidator::Validate(DataSize->Append(dataSize)); diff --git a/ydb/core/formats/arrow/accessor/sub_columns/ut/ut_sub_columns.cpp b/ydb/core/formats/arrow/accessor/sub_columns/ut/ut_sub_columns.cpp index 5bb31befcbb..c8998590e51 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/ut/ut_sub_columns.cpp +++ b/ydb/core/formats/arrow/accessor/sub_columns/ut/ut_sub_columns.cpp @@ -41,7 +41,7 @@ Y_UNIT_TEST_SUITE(SubColumnsArrayAccessor) { } Y_UNIT_TEST(SlicesDef) { - NSubColumns::TSettings settings(4, 1, 0); + NSubColumns::TSettings settings(4, 1, 0, 0); const std::vector<TString> jsons = { R"({"a" : 1, "b" : 1, "c" : "111"})", @@ -113,4 +113,81 @@ Y_UNIT_TEST_SUITE(SubColumnsArrayAccessor) { ("string", arrSlice->DebugJson().GetStringRobust()); } } + + Y_UNIT_TEST(FiltersDef) { + NSubColumns::TSettings settings(4, 1, 0, 0); + + const std::vector<TString> jsons = { + R"({"a" : 1, "b" : 1, "c" : "111"})", + "null", + R"({"a1" : 2, "b" : 2, "c" : "222"})", + R"({"a" : 3, "b" : 3, "c" : "333"})", + "null", + R"({"a" : 5, "b1" : 5})", + }; + + TTrivialArray::TPlainBuilder<arrow::BinaryType> arrBuilder; + ui32 idx = 0; + for (auto&& i : jsons) { + if (i != "null") { + auto v = NBinaryJson::SerializeToBinaryJson(i); + NBinaryJson::TBinaryJson* bJson = std::get_if<NBinaryJson::TBinaryJson>(&v); + arrBuilder.AddRecord(idx, std::string_view(bJson->data(), bJson->size())); + } + ++idx; + } + auto bJsonArr = arrBuilder.Finish(jsons.size()); + auto arrData = TSubColumnsArray::Make(bJsonArr, std::make_shared<NSubColumns::TFirstLevelSchemaData>(), settings).DetachResult(); + Cerr << arrData->DebugJson() << Endl; + AFL_VERIFY(PrintBinaryJsons(arrData->GetChunkedArray()) == R"([[{"a":"1","b":"1","c":"111"},null,{"a1":"2","b":"2","c":"222"},{"a":"3","b":"3","c":"333"},null,{"a":"5","b1":"5"}]])")( + "string", PrintBinaryJsons(arrData->GetChunkedArray())); + { + TColumnFilter filter = TColumnFilter::BuildAllowFilter(); + filter.Add(true, 1); + filter.Add(false, 1); + filter.Add(true, 1); + filter.Add(false, 1); + filter.Add(true, 1); + filter.Add(false, 1); + auto arrSlice = filter.Apply(arrData); + AFL_VERIFY(PrintBinaryJsons(arrSlice->GetChunkedArray()) == R"([[{"a":"1","b":"1","c":"111"},{"a1":"2","b":"2","c":"222"},null]])")( + "string", PrintBinaryJsons(arrSlice->GetChunkedArray())); + } + { + TColumnFilter filter = TColumnFilter::BuildAllowFilter(); + filter.Add(false, 1); + filter.Add(true, 1); + filter.Add(false, 1); + filter.Add(true, 1); + filter.Add(false, 1); + filter.Add(true, 1); + auto arrSlice = filter.Apply(arrData); + AFL_VERIFY(PrintBinaryJsons(arrSlice->GetChunkedArray()) == R"([[null,{"a":"3","b":"3","c":"333"},{"a":"5","b1":"5"}]])")( + "string", PrintBinaryJsons(arrSlice->GetChunkedArray())); + } + { + TColumnFilter filter = TColumnFilter::BuildAllowFilter(); + filter.Add(false, 1); + filter.Add(true, 3); + filter.Add(false, 2); + auto arrSlice = filter.Apply(arrData); + AFL_VERIFY(PrintBinaryJsons(arrSlice->GetChunkedArray()) == R"([[null,{"a1":"2","b":"2","c":"222"},{"a":"3","b":"3","c":"333"}]])")( + "string", PrintBinaryJsons(arrSlice->GetChunkedArray())); + } + { + TColumnFilter filter = TColumnFilter::BuildAllowFilter(); + filter.Add(false, 1); + filter.Add(true, 1); + filter.Add(false, 4); + auto arrSlice = filter.Apply(arrData); + AFL_VERIFY(PrintBinaryJsons(arrSlice->GetChunkedArray()) == R"([[null]])")("string", PrintBinaryJsons(arrSlice->GetChunkedArray())); + } + { + TColumnFilter filter = TColumnFilter::BuildAllowFilter(); + filter.Add(true, 1); + filter.Add(false, 5); + auto arrSlice = filter.Apply(arrData); + AFL_VERIFY(PrintBinaryJsons(arrSlice->GetChunkedArray()) == R"([[{"a":"1","b":"1","c":"111"}]])")("string", PrintBinaryJsons(arrSlice->GetChunkedArray())); + } + } }; diff --git a/ydb/core/formats/arrow/accessor/sub_columns/ut/ya.make b/ydb/core/formats/arrow/accessor/sub_columns/ut/ya.make index 68c75fb208d..09608ce94eb 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/ut/ya.make +++ b/ydb/core/formats/arrow/accessor/sub_columns/ut/ya.make @@ -5,6 +5,7 @@ SIZE(SMALL) PEERDIR( ydb/core/formats/arrow/accessor/sub_columns yql/essentials/public/udf/service/stub + ydb/core/formats/arrow ) SRCS( diff --git a/ydb/core/formats/arrow/accessor/sub_columns/ya.make b/ydb/core/formats/arrow/accessor/sub_columns/ya.make index 18e798da3f7..4226c84e592 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/ya.make +++ b/ydb/core/formats/arrow/accessor/sub_columns/ya.make @@ -28,3 +28,7 @@ SRCS( YQL_LAST_ABI_VERSION() END() + +RECURSE_FOR_TESTS( + ut +) diff --git a/ydb/core/formats/arrow/accessor/ya.make b/ydb/core/formats/arrow/accessor/ya.make index add9a81bf85..522bbb84f0b 100644 --- a/ydb/core/formats/arrow/accessor/ya.make +++ b/ydb/core/formats/arrow/accessor/ya.make @@ -4,6 +4,7 @@ PEERDIR( ydb/core/formats/arrow/accessor/abstract ydb/core/formats/arrow/accessor/plain ydb/core/formats/arrow/accessor/composite_serial + ydb/core/formats/arrow/accessor/composite ydb/core/formats/arrow/accessor/sparsed ydb/core/formats/arrow/accessor/sub_columns ) diff --git a/ydb/core/formats/arrow/arrow_filter.cpp b/ydb/core/formats/arrow/arrow_filter.cpp index 0b88315029c..1e019b82374 100644 --- a/ydb/core/formats/arrow/arrow_filter.cpp +++ b/ydb/core/formats/arrow/arrow_filter.cpp @@ -441,6 +441,16 @@ void TColumnFilter::Apply(const ui32 expectedRecordsCount, std::vector<arrow::Da } } +std::shared_ptr<NAccessor::IChunkedArray> TColumnFilter::Apply( + const std::shared_ptr<NAccessor::IChunkedArray>& source, const TApplyContext& context /*= Default<TApplyContext>()*/) const { + if (context.HasSlice()) { + auto sliceArray = source->ISlice(*context.GetStartPos(), *context.GetCount()); + return sliceArray->ApplyFilter(*this, sliceArray); + } else { + return source->ApplyFilter(*this, source); + } +} + const std::vector<bool>& TColumnFilter::BuildSimpleFilter() const { if (!FilterPlain) { Y_ABORT_UNLESS(RecordsCount); @@ -639,8 +649,7 @@ TColumnFilter::TIterator TColumnFilter::GetIterator(const bool reverse, const ui } else if (IsTotalDenyFilter()) { return TIterator(reverse, expectedSize, false); } else { - AFL_VERIFY(expectedSize == GetRecordsCountVerified())("expected", expectedSize)("count", GetRecordsCountVerified())( - "reverse", reverse); + AFL_VERIFY(expectedSize == GetRecordsCountVerified())("expected", expectedSize)("count", GetRecordsCountVerified())("reverse", reverse); return TIterator(reverse, Filter, GetStartValue(reverse)); } } diff --git a/ydb/core/formats/arrow/arrow_filter.h b/ydb/core/formats/arrow/arrow_filter.h index f2b9641d1c0..59ff5ffb0ad 100644 --- a/ydb/core/formats/arrow/arrow_filter.h +++ b/ydb/core/formats/arrow/arrow_filter.h @@ -7,6 +7,10 @@ #include <deque> +namespace NKikimr::NArrow::NAccessor { +class IChunkedArray; +} + namespace NKikimr::NArrow { class TGeneralContainer; @@ -39,6 +43,7 @@ private: } public: + class TSlicesIterator { private: const TColumnFilter& Owner; @@ -267,10 +272,12 @@ public: TApplyContext& Slice(const ui32 start, const ui32 count); }; - bool Apply(std::shared_ptr<TGeneralContainer>& batch, const TApplyContext& context = Default<TApplyContext>()) const; - bool Apply(std::shared_ptr<arrow::Table>& batch, const TApplyContext& context = Default<TApplyContext>()) const; - bool Apply(std::shared_ptr<arrow::RecordBatch>& batch, const TApplyContext& context = Default<TApplyContext>()) const; + [[nodiscard]] bool Apply(std::shared_ptr<TGeneralContainer>& batch, const TApplyContext& context = Default<TApplyContext>()) const; + [[nodiscard]] bool Apply(std::shared_ptr<arrow::Table>& batch, const TApplyContext& context = Default<TApplyContext>()) const; + [[nodiscard]] bool Apply(std::shared_ptr<arrow::RecordBatch>& batch, const TApplyContext& context = Default<TApplyContext>()) const; void Apply(const ui32 expectedRecordsCount, std::vector<arrow::Datum*>& datums) const; + [[nodiscard]] std::shared_ptr<NAccessor::IChunkedArray> Apply( + const std::shared_ptr<NAccessor::IChunkedArray>& source, const TApplyContext& context = Default<TApplyContext>()) const; // Combines filters by 'and' operator (extFilter count is true positions count in self, thought extFitler patch exactly that positions) TColumnFilter CombineSequentialAnd(const TColumnFilter& extFilter) const Y_WARN_UNUSED_RESULT; diff --git a/ydb/core/formats/arrow/arrow_helpers.cpp b/ydb/core/formats/arrow/arrow_helpers.cpp index c8e3045ff14..e2089e484c4 100644 --- a/ydb/core/formats/arrow/arrow_helpers.cpp +++ b/ydb/core/formats/arrow/arrow_helpers.cpp @@ -1,27 +1,29 @@ #include "arrow_helpers.h" -#include "switch/switch_type.h" #include "permutations.h" + #include "common/adapter.h" -#include "serializer/native.h" #include "serializer/abstract.h" +#include "serializer/native.h" #include "serializer/stream.h" +#include "switch/switch_type.h" -#include <ydb/library/formats/arrow/common/validation.h> -#include <ydb/library/formats/arrow/simple_arrays_cache.h> +#include <ydb/library/actors/core/log.h> #include <ydb/library/formats/arrow/replace_key.h> -#include <ydb/library/yverify_stream/yverify_stream.h> +#include <ydb/library/formats/arrow/simple_arrays_cache.h> +#include <ydb/library/formats/arrow/validation/validation.h> #include <ydb/library/services/services.pb.h> +#include <ydb/library/yverify_stream/yverify_stream.h> -#include <util/system/yassert.h> -#include <util/string/join.h> -#include <contrib/libs/apache/arrow/cpp/src/arrow/io/memory.h> -#include <contrib/libs/apache/arrow/cpp/src/arrow/ipc/reader.h> -#include <contrib/libs/apache/arrow/cpp/src/arrow/compute/api.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/array/array_primitive.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/array/builder_primitive.h> +#include <contrib/libs/apache/arrow/cpp/src/arrow/compute/api.h> +#include <contrib/libs/apache/arrow/cpp/src/arrow/io/memory.h> +#include <contrib/libs/apache/arrow/cpp/src/arrow/ipc/reader.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/type_traits.h> #include <library/cpp/containers/stack_vector/stack_vec.h> -#include <ydb/library/actors/core/log.h> +#include <util/string/join.h> +#include <util/system/yassert.h> + #include <memory> #define Y_VERIFY_OK(status) Y_ABORT_UNLESS(status.ok(), "%s", status.ToString().c_str()) @@ -82,7 +84,8 @@ arrow::Result<std::shared_ptr<arrow::DataType>> GetCSVArrowType(NScheme::TTypeIn } } -arrow::Result<arrow::FieldVector> MakeArrowFields(const std::vector<std::pair<TString, NScheme::TTypeInfo>>& columns, const std::set<std::string>& notNullColumns) { +arrow::Result<arrow::FieldVector> MakeArrowFields( + const std::vector<std::pair<TString, NScheme::TTypeInfo>>& columns, const std::set<std::string>& notNullColumns) { std::vector<std::shared_ptr<arrow::Field>> fields; fields.reserve(columns.size()); TVector<TString> errors; @@ -101,7 +104,8 @@ arrow::Result<arrow::FieldVector> MakeArrowFields(const std::vector<std::pair<TS return arrow::Status::TypeError(JoinSeq(", ", errors)); } -arrow::Result<std::shared_ptr<arrow::Schema>> MakeArrowSchema(const std::vector<std::pair<TString, NScheme::TTypeInfo>>& ydbColumns, const std::set<std::string>& notNullColumns) { +arrow::Result<std::shared_ptr<arrow::Schema>> MakeArrowSchema( + const std::vector<std::pair<TString, NScheme::TTypeInfo>>& ydbColumns, const std::set<std::string>& notNullColumns) { const auto fields = MakeArrowFields(ydbColumns, notNullColumns); if (fields.ok()) { return std::make_shared<arrow::Schema>(fields.ValueUnsafe()); @@ -130,21 +134,19 @@ TString SerializeBatchNoCompression(const std::shared_ptr<arrow::RecordBatch>& b return SerializeBatch(batch, writeOptions); } -std::shared_ptr<arrow::RecordBatch> DeserializeBatch(const TString& blob, const std::shared_ptr<arrow::Schema>& schema) -{ +std::shared_ptr<arrow::RecordBatch> DeserializeBatch(const TString& blob, const std::shared_ptr<arrow::Schema>& schema) { auto result = NSerialization::TNativeSerializer().Deserialize(blob, schema); if (result.ok()) { return *result; } else { - AFL_ERROR(NKikimrServices::ARROW_HELPER)("event", "cannot_parse")("message", result.status().ToString()) - ("schema_columns_count", schema->num_fields())("schema_columns", JoinSeq(",", schema->field_names())); + AFL_ERROR(NKikimrServices::ARROW_HELPER)("event", "cannot_parse")("message", result.status().ToString())( + "schema_columns_count", schema->num_fields())("schema_columns", JoinSeq(",", schema->field_names())); return nullptr; } } -void DedupSortedBatch(const std::shared_ptr<arrow::RecordBatch>& batch, - const std::shared_ptr<arrow::Schema>& sortingKey, - std::vector<std::shared_ptr<arrow::RecordBatch>>& out) { +void DedupSortedBatch(const std::shared_ptr<arrow::RecordBatch>& batch, const std::shared_ptr<arrow::Schema>& sortingKey, + std::vector<std::shared_ptr<arrow::RecordBatch>>& out) { if (batch->num_rows() < 2) { out.push_back(batch); return; @@ -179,8 +181,7 @@ void DedupSortedBatch(const std::shared_ptr<arrow::RecordBatch>& batch, Y_DEBUG_ABORT_UNLESS(NArrow::IsSortedAndUnique(out.back(), sortingKey)); } -bool IsSorted(const std::shared_ptr<arrow::RecordBatch>& batch, - const std::shared_ptr<arrow::Schema>& sortingKey, bool desc) { +bool IsSorted(const std::shared_ptr<arrow::RecordBatch>& batch, const std::shared_ptr<arrow::Schema>& sortingKey, bool desc) { auto keyBatch = TColumnOperator().Adapt(batch, sortingKey).DetachResult(); if (desc) { return IsSelfSorted<true, false>(keyBatch); @@ -189,8 +190,7 @@ bool IsSorted(const std::shared_ptr<arrow::RecordBatch>& batch, } } -bool IsSortedAndUnique(const std::shared_ptr<arrow::RecordBatch>& batch, - const std::shared_ptr<arrow::Schema>& sortingKey, bool desc) { +bool IsSortedAndUnique(const std::shared_ptr<arrow::RecordBatch>& batch, const std::shared_ptr<arrow::Schema>& sortingKey, bool desc) { auto keyBatch = TColumnOperator().Adapt(batch, sortingKey).DetachResult(); if (desc) { return IsSelfSorted<true, true>(keyBatch); @@ -209,8 +209,8 @@ std::shared_ptr<arrow::RecordBatch> SortBatch( } } -std::shared_ptr<arrow::RecordBatch> SortBatch(const std::shared_ptr<arrow::RecordBatch>& batch, const std::shared_ptr<arrow::Schema>& sortingKey, - const bool andUnique) { +std::shared_ptr<arrow::RecordBatch> SortBatch( + const std::shared_ptr<arrow::RecordBatch>& batch, const std::shared_ptr<arrow::Schema>& sortingKey, const bool andUnique) { auto sortPermutation = MakeSortPermutation(batch, sortingKey, andUnique); if (sortPermutation) { return Reorder(batch, sortPermutation, andUnique); @@ -240,4 +240,4 @@ std::shared_ptr<arrow::Table> ReallocateBatch(const std::shared_ptr<arrow::Table return NArrow::TStatusValidator::GetValid(arrow::Table::FromRecordBatches(batches)); } -} +} // namespace NKikimr::NArrow diff --git a/ydb/core/formats/arrow/common/adapter.h b/ydb/core/formats/arrow/common/adapter.h index 3e6acbbeb9e..104f7debe96 100644 --- a/ydb/core/formats/arrow/common/adapter.h +++ b/ydb/core/formats/arrow/common/adapter.h @@ -5,7 +5,7 @@ #include <ydb/core/formats/arrow/arrow_filter.h> #include <ydb/library/formats/arrow/arrow_helpers.h> -#include <ydb/library/formats/arrow/common/validation.h> +#include <ydb/library/formats/arrow/validation/validation.h> #include <ydb/library/yverify_stream/yverify_stream.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/array/array_base.h> diff --git a/ydb/core/formats/arrow/common/container.h b/ydb/core/formats/arrow/common/container.h index 5e4705f2014..72dccce50b1 100644 --- a/ydb/core/formats/arrow/common/container.h +++ b/ydb/core/formats/arrow/common/container.h @@ -1,22 +1,23 @@ #pragma once +#include <ydb/core/formats/arrow/accessor/abstract/accessor.h> + #include <ydb/library/accessor/accessor.h> #include <ydb/library/conclusion/result.h> #include <ydb/library/conclusion/status.h> #include <ydb/library/formats/arrow/modifier/schema.h> -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> -#include <contrib/libs/apache/arrow/cpp/src/arrow/type.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/table.h> - -#include <util/system/types.h> +#include <contrib/libs/apache/arrow/cpp/src/arrow/type.h> #include <util/string/builder.h> +#include <util/system/types.h> namespace NKikimr::NArrow { class IFieldsConstructor { private: virtual std::shared_ptr<arrow::Scalar> DoGetDefaultColumnElementValue(const std::string& fieldName) const = 0; + public: TConclusion<std::shared_ptr<arrow::Scalar>> GetDefaultColumnElementValue(const std::shared_ptr<arrow::Field>& field, const bool force) const; }; @@ -27,6 +28,7 @@ private: YDB_READONLY_DEF(std::shared_ptr<NModifier::TSchema>, Schema); YDB_READONLY_DEF(std::vector<std::shared_ptr<NAccessor::IChunkedArray>>, Columns); void Initialize(); + public: TGeneralContainer(const ui32 recordsCount); @@ -57,8 +59,8 @@ public: return DebugJson(withData).GetStringRobust(); } - [[nodiscard]] TConclusionStatus SyncSchemaTo(const std::shared_ptr<arrow::Schema>& schema, - const IFieldsConstructor* defaultFieldsConstructor, const bool forceDefaults); + [[nodiscard]] TConclusionStatus SyncSchemaTo( + const std::shared_ptr<arrow::Schema>& schema, const IFieldsConstructor* defaultFieldsConstructor, const bool forceDefaults); bool HasColumn(const std::string& name) { return Schema->HasField(name); @@ -120,7 +122,8 @@ public: TGeneralContainer(const std::shared_ptr<arrow::RecordBatch>& table); TGeneralContainer(const std::shared_ptr<arrow::Schema>& schema, std::vector<std::shared_ptr<NAccessor::IChunkedArray>>&& columns); TGeneralContainer(const std::shared_ptr<NModifier::TSchema>& schema, std::vector<std::shared_ptr<NAccessor::IChunkedArray>>&& columns); - TGeneralContainer(const std::vector<std::shared_ptr<arrow::Field>>& fields, std::vector<std::shared_ptr<NAccessor::IChunkedArray>>&& columns); + TGeneralContainer( + const std::vector<std::shared_ptr<arrow::Field>>& fields, std::vector<std::shared_ptr<NAccessor::IChunkedArray>>&& columns); arrow::Status ValidateFull() const { return arrow::Status::OK(); @@ -130,4 +133,4 @@ public: std::shared_ptr<NAccessor::IChunkedArray> GetAccessorByNameVerified(const std::string& fieldId) const; }; -} +} // namespace NKikimr::NArrow diff --git a/ydb/core/formats/arrow/dictionary/object.cpp b/ydb/core/formats/arrow/dictionary/object.cpp index 36c9fe3fc27..ca7ba9855f3 100644 --- a/ydb/core/formats/arrow/dictionary/object.cpp +++ b/ydb/core/formats/arrow/dictionary/object.cpp @@ -1,6 +1,9 @@ #include "object.h" + #include <ydb/core/formats/arrow/transformer/dictionary.h> -#include <ydb/library/formats/arrow/common/validation.h> + +#include <ydb/library/formats/arrow/validation/validation.h> + #include <util/string/builder.h> namespace NKikimr::NArrow::NDictionary { @@ -40,4 +43,4 @@ NTransformation::ITransformer::TPtr TEncodingSettings::BuildDecoder() const { } } -} +} // namespace NKikimr::NArrow::NDictionary diff --git a/ydb/core/formats/arrow/hash/calcer.h b/ydb/core/formats/arrow/hash/calcer.h index 490a0e05e36..16e390552f2 100644 --- a/ydb/core/formats/arrow/hash/calcer.h +++ b/ydb/core/formats/arrow/hash/calcer.h @@ -3,19 +3,18 @@ #include <ydb/core/formats/arrow/reader/position.h> #include <ydb/library/actors/core/log.h> -#include <ydb/library/services/services.pb.h> #include <ydb/library/formats/arrow/hash/xx_hash.h> -#include <ydb/library/formats/arrow/common/validation.h> +#include <ydb/library/formats/arrow/validation/validation.h> +#include <ydb/library/services/services.pb.h> -#include <contrib/libs/apache/arrow/cpp/src/arrow/record_batch.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/array/array_base.h> - -#include <util/system/types.h> -#include <util/string/join.h> +#include <contrib/libs/apache/arrow/cpp/src/arrow/record_batch.h> #include <util/generic/string.h> +#include <util/string/join.h> +#include <util/system/types.h> -#include <vector> #include <optional> +#include <vector> namespace NKikimr::NArrow::NHash { @@ -26,13 +25,15 @@ public: Verify, ReturnEmpty }; + private: ui64 Seed = 0; const std::vector<TString> ColumnNames; const ENoColumnPolicy NoColumnPolicy; template <class TDataContainer> - std::vector<std::shared_ptr<typename NAdapter::TDataBuilderPolicy<TDataContainer>::TColumn>> GetColumns(const std::shared_ptr<TDataContainer>& batch) const { + std::vector<std::shared_ptr<typename NAdapter::TDataBuilderPolicy<TDataContainer>::TColumn>> GetColumns( + const std::shared_ptr<TDataContainer>& batch) const { std::vector<std::shared_ptr<typename NAdapter::TDataBuilderPolicy<TDataContainer>::TColumn>> columns; columns.reserve(ColumnNames.size()); for (auto& colName : ColumnNames) { @@ -51,8 +52,8 @@ private: } } if (columns.empty()) { - AFL_WARN(NKikimrServices::ARROW_HELPER)("event", "cannot_read_all_columns")("reason", "fields_not_found") - ("field_names", JoinSeq(",", ColumnNames))("batch_fields", JoinSeq(",", batch->schema()->field_names())); + AFL_WARN(NKikimrServices::ARROW_HELPER)("event", "cannot_read_all_columns")("reason", "fields_not_found")( + "field_names", JoinSeq(",", ColumnNames))("batch_fields", JoinSeq(",", batch->schema()->field_names())); } return columns; } @@ -81,10 +82,10 @@ public: std::vector<NAccessor::IChunkedArray::TReader> columnScanners; for (auto&& i : columns) { - columnScanners.emplace_back(NAccessor::IChunkedArray::TReader(std::make_shared<typename NAdapter::TDataBuilderPolicy<TDataContainer>::TAccessor>(i))); + columnScanners.emplace_back( + NAccessor::IChunkedArray::TReader(std::make_shared<typename NAdapter::TDataBuilderPolicy<TDataContainer>::TAccessor>(i))); } - { NXX64::TStreamStringHashCalcer hashCalcer(Seed); for (int row = 0; row < batch->num_rows(); ++row) { @@ -128,7 +129,6 @@ public: AFL_VERIFY(ExecuteToArrayImpl(batch, acceptor)); return result; } - }; -} +} // namespace NKikimr::NArrow::NHash diff --git a/ydb/core/formats/arrow/permutations.cpp b/ydb/core/formats/arrow/permutations.cpp index 2a7804c918e..980c8aecf75 100644 --- a/ydb/core/formats/arrow/permutations.cpp +++ b/ydb/core/formats/arrow/permutations.cpp @@ -1,14 +1,13 @@ -#include "permutations.h" - #include "arrow_helpers.h" +#include "permutations.h" #include "size_calcer.h" -#include "hash/calcer.h" -#include <ydb/library/services/services.pb.h> +#include "hash/calcer.h" -#include <ydb/library/formats/arrow/common/validation.h> -#include <ydb/library/formats/arrow/replace_key.h> #include <ydb/library/actors/core/log.h> +#include <ydb/library/formats/arrow/replace_key.h> +#include <ydb/library/formats/arrow/validation/validation.h> +#include <ydb/library/services/services.pb.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/array/builder_primitive.h> @@ -86,8 +85,8 @@ std::shared_ptr<arrow::UInt64Array> MakeSortPermutation(const std::vector<std::s return out; } -std::shared_ptr<arrow::UInt64Array> MakeSortPermutation(const std::shared_ptr<arrow::RecordBatch>& batch, - const std::shared_ptr<arrow::Schema>& sortingKey, const bool andUnique) { +std::shared_ptr<arrow::UInt64Array> MakeSortPermutation( + const std::shared_ptr<arrow::RecordBatch>& batch, const std::shared_ptr<arrow::Schema>& sortingKey, const bool andUnique) { auto keyBatch = TColumnOperator().VerifyIfAbsent().Adapt(batch, sortingKey).DetachResult(); return MakeSortPermutation(keyBatch->columns(), andUnique); } @@ -107,24 +106,30 @@ bool BuildHashUI64Impl(std::shared_ptr<TDataContainer>& batch, const std::vector return false; } Y_ABORT_UNLESS(column); - if (column->type()->id() == arrow::Type::UINT64 || column->type()->id() == arrow::Type::UINT32 || column->type()->id() == arrow::Type::INT64 || column->type()->id() == arrow::Type::INT32) { - batch = TStatusValidator::GetValid(batch->AddColumn(batch->num_columns(), std::make_shared<arrow::Field>(hashFieldName, column->type()), column)); + if (column->type()->id() == arrow::Type::UINT64 || column->type()->id() == arrow::Type::UINT32 || + column->type()->id() == arrow::Type::INT64 || column->type()->id() == arrow::Type::INT32) { + batch = TStatusValidator::GetValid( + batch->AddColumn(batch->num_columns(), std::make_shared<arrow::Field>(hashFieldName, column->type()), column)); return true; } } - std::shared_ptr<arrow::Array> hashColumn = NArrow::NHash::TXX64(fieldNames, NArrow::NHash::TXX64::ENoColumnPolicy::Verify, 34323543).ExecuteToArray(batch); - batch = NAdapter::TDataBuilderPolicy<TDataContainer>::AddColumn(batch, std::make_shared<arrow::Field>(hashFieldName, hashColumn->type()), hashColumn); + std::shared_ptr<arrow::Array> hashColumn = + NArrow::NHash::TXX64(fieldNames, NArrow::NHash::TXX64::ENoColumnPolicy::Verify, 34323543).ExecuteToArray(batch); + batch = NAdapter::TDataBuilderPolicy<TDataContainer>::AddColumn( + batch, std::make_shared<arrow::Field>(hashFieldName, hashColumn->type()), hashColumn); return true; } -} +} // namespace -bool THashConstructor::BuildHashUI64(std::shared_ptr<arrow::Table>& batch, const std::vector<std::string>& fieldNames, const std::string& hashFieldName) { +bool THashConstructor::BuildHashUI64( + std::shared_ptr<arrow::Table>& batch, const std::vector<std::string>& fieldNames, const std::string& hashFieldName) { return BuildHashUI64Impl(batch, fieldNames, hashFieldName); } -bool THashConstructor::BuildHashUI64(std::shared_ptr<arrow::RecordBatch>& batch, const std::vector<std::string>& fieldNames, const std::string& hashFieldName) { +bool THashConstructor::BuildHashUI64( + std::shared_ptr<arrow::RecordBatch>& batch, const std::vector<std::string>& fieldNames, const std::string& hashFieldName) { return BuildHashUI64Impl(batch, fieldNames, hashFieldName); } -} +} // namespace NKikimr::NArrow diff --git a/ydb/core/formats/arrow/program/abstract.h b/ydb/core/formats/arrow/program/abstract.h index 5a64b6ce371..ebf886d7935 100644 --- a/ydb/core/formats/arrow/program/abstract.h +++ b/ydb/core/formats/arrow/program/abstract.h @@ -1,8 +1,9 @@ #pragma once +#include <ydb/core/formats/arrow/accessor/abstract/accessor.h> + #include <ydb/library/accessor/accessor.h> #include <ydb/library/conclusion/result.h> #include <ydb/library/conclusion/status.h> -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> #include <util/generic/string.h> diff --git a/ydb/core/formats/arrow/program/collection.cpp b/ydb/core/formats/arrow/program/collection.cpp index c11fb24d0d3..e85cf2756bc 100644 --- a/ydb/core/formats/arrow/program/collection.cpp +++ b/ydb/core/formats/arrow/program/collection.cpp @@ -20,7 +20,7 @@ void TAccessorsCollection::AddVerified(const ui32 columnId, const TAccessorColle AFL_VERIFY(!data.GetItWasScalar()); } if (UseFilter && withFilter && !Filter->IsTotalAllowFilter()) { - auto filtered = data->ApplyFilter(*Filter); + auto filtered = Filter->Apply(data.GetData()); RecordsCountActual = filtered->GetRecordsCount(); AFL_VERIFY(Accessors.emplace(columnId, filtered).second); } else { diff --git a/ydb/core/formats/arrow/program/collection.h b/ydb/core/formats/arrow/program/collection.h index 1a69b9e4244..9e7434e863c 100644 --- a/ydb/core/formats/arrow/program/collection.h +++ b/ydb/core/formats/arrow/program/collection.h @@ -2,10 +2,10 @@ #include "abstract.h" +#include <ydb/core/formats/arrow/accessor/abstract/accessor.h> #include <ydb/core/formats/arrow/arrow_filter.h> #include <ydb/core/formats/arrow/common/container.h> -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> #include <ydb/library/formats/arrow/validation/validation.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/datum.h> @@ -417,7 +417,7 @@ public: } else { *Filter = Filter->CombineSequentialAnd(filter); for (auto&& i : Accessors) { - i.second = TAccessorCollectedContainer(i.second.GetData()->ApplyFilter(filter)); + i.second = TAccessorCollectedContainer(i.second.GetData()->ApplyFilter(filter, i.second.GetData())); } } RecordsCountActual = Filter->GetFilteredCount(); diff --git a/ydb/core/formats/arrow/program/kernel_logic.cpp b/ydb/core/formats/arrow/program/kernel_logic.cpp index e8a4e8911d5..7822ea65c6e 100644 --- a/ydb/core/formats/arrow/program/kernel_logic.cpp +++ b/ydb/core/formats/arrow/program/kernel_logic.cpp @@ -1,9 +1,8 @@ #include "kernel_logic.h" +#include <ydb/core/formats/arrow/accessor/composite/accessor.h> #include <ydb/core/formats/arrow/accessor/sub_columns/accessor.h> -#include <ydb/library/formats/arrow/accessor/composite/accessor.h> - namespace NKikimr::NArrow::NSSA { TConclusion<bool> TGetJsonPath::DoExecute(const std::vector<TColumnChainInfo>& input, const std::vector<TColumnChainInfo>& output, @@ -44,7 +43,7 @@ TConclusion<bool> TGetJsonPath::DoExecute(const std::vector<TColumnChainInfo>& i resources->AddVerified(output.front().GetColumnId(), builder.Finish()); return true; } - if (accJson->GetType() == IChunkedArray::EType::SubColumnsArray) { + if (accJson->GetType() != IChunkedArray::EType::SubColumnsArray) { return false; } resources->AddVerified(output.front().GetColumnId(), ExtractArray(accJson, svPath)); diff --git a/ydb/core/formats/arrow/reader/merger.cpp b/ydb/core/formats/arrow/reader/merger.cpp index 06b5d2be4b2..b6c56ba2318 100644 --- a/ydb/core/formats/arrow/reader/merger.cpp +++ b/ydb/core/formats/arrow/reader/merger.cpp @@ -105,7 +105,7 @@ std::shared_ptr<arrow::Table> TMergePartialStream::SingleSourceDrain(const TSort *lastResultPosition = TCursor(keys, 0, SortSchema->field_names()); } if (SortHeap.Current().GetFilter()) { - SortHeap.Current().GetFilter()->Apply(result, TColumnFilter::TApplyContext(pos.GetPosition() + (include ? 0 : 1), resultSize)); + AFL_VERIFY(SortHeap.Current().GetFilter()->Apply(result, TColumnFilter::TApplyContext(pos.GetPosition() + (include ? 0 : 1), resultSize))); } } else { result = SortHeap.Current().GetKeyColumns().SliceData(startPos, resultSize); @@ -114,7 +114,7 @@ std::shared_ptr<arrow::Table> TMergePartialStream::SingleSourceDrain(const TSort *lastResultPosition = TCursor(keys, keys->num_rows() - 1, SortSchema->field_names()); } if (SortHeap.Current().GetFilter()) { - SortHeap.Current().GetFilter()->Apply(result, TColumnFilter::TApplyContext(startPos, resultSize)); + AFL_VERIFY(SortHeap.Current().GetFilter()->Apply(result, TColumnFilter::TApplyContext(startPos, resultSize))); } } if (!result || !result->num_rows()) { diff --git a/ydb/core/formats/arrow/reader/position.h b/ydb/core/formats/arrow/reader/position.h index 34d4d2f7c72..f403dd7fe3a 100644 --- a/ydb/core/formats/arrow/reader/position.h +++ b/ydb/core/formats/arrow/reader/position.h @@ -1,16 +1,16 @@ #pragma once +#include <ydb/core/formats/arrow/accessor/abstract/accessor.h> +#include <ydb/core/formats/arrow/common/container.h> #include <ydb/core/formats/arrow/permutations.h> #include <ydb/core/formats/arrow/switch/switch_type.h> -#include <ydb/core/formats/arrow/common/container.h> -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> #include <ydb/library/accessor/accessor.h> #include <ydb/library/actors/core/log.h> -#include <library/cpp/json/writer/json_value.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/array/array_base.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/record_batch.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/type.h> +#include <library/cpp/json/writer/json_value.h> #include <util/system/types.h> namespace NKikimr::NArrow::NMerger { @@ -22,15 +22,14 @@ class TCursor { private: YDB_READONLY(ui64, Position, 0); std::vector<NAccessor::IChunkedArray::TFullDataAddress> PositionAddress; + public: TCursor() = default; TCursor(const std::shared_ptr<arrow::Table>& table, const ui64 position, const std::vector<std::string>& columns); TCursor(const ui64 position, const std::vector<NAccessor::IChunkedArray::TFullDataAddress>& addresses) : Position(position) - , PositionAddress(addresses) - { - + , PositionAddress(addresses) { } NJson::TJsonValue DebugJson() const { @@ -57,7 +56,6 @@ public: std::partial_ordering Compare(const TSortableScanData& item, const ui64 itemPosition) const; std::partial_ordering Compare(const TCursor& item) const; - }; class TSortableScanData { @@ -74,16 +72,17 @@ private: bool Contains(const ui64 position) const { return StartPosition <= position && position < FinishPosition; } + public: TSortableScanData(const ui64 position, const std::shared_ptr<arrow::RecordBatch>& batch); TSortableScanData(const ui64 position, const std::shared_ptr<arrow::RecordBatch>& batch, const std::vector<std::string>& columns); TSortableScanData(const ui64 position, const std::shared_ptr<arrow::Table>& batch, const std::vector<std::string>& columns); TSortableScanData(const ui64 position, const std::shared_ptr<TGeneralContainer>& batch, const std::vector<std::string>& columns); - TSortableScanData(const ui64 position, const ui64 recordsCount, const std::vector<std::shared_ptr<NAccessor::IChunkedArray>>& columns, const std::vector<std::shared_ptr<arrow::Field>>& fields) + TSortableScanData(const ui64 position, const ui64 recordsCount, const std::vector<std::shared_ptr<NAccessor::IChunkedArray>>& columns, + const std::vector<std::shared_ptr<arrow::Field>>& fields) : RecordsCount(recordsCount) , Columns(columns) - , Fields(fields) - { + , Fields(fields) { BuildPosition(position); } @@ -154,13 +153,13 @@ public: [[nodiscard]] bool InitPosition(const ui64 position); - std::shared_ptr<arrow::Table> Slice(const ui64 offset, const ui64 count) const { - std::vector<std::shared_ptr<arrow::ChunkedArray>> slicedArrays; - for (auto&& i : Columns) { - slicedArrays.emplace_back(i->Slice(offset, count)); - } - return arrow::Table::Make(std::make_shared<arrow::Schema>(Fields), slicedArrays, count); - } + std::shared_ptr<arrow::Table> Slice(const ui64 offset, const ui64 count) const { + std::vector<std::shared_ptr<arrow::ChunkedArray>> slicedArrays; + for (auto&& i : Columns) { + slicedArrays.emplace_back(i->Slice(offset, count)); + } + return arrow::Table::Make(std::make_shared<arrow::Schema>(Fields), slicedArrays, count); + } bool IsSameSchema(const std::shared_ptr<arrow::Schema>& schema) const { if (Fields.size() != (size_t)schema->num_fields()) { @@ -254,9 +253,9 @@ public: TSortableBatchPosition(TRWSortableBatchPosition& source) = delete; TSortableBatchPosition(TRWSortableBatchPosition&& source) = delete; - TSortableBatchPosition operator= (const TRWSortableBatchPosition& source) = delete; - TSortableBatchPosition operator= (TRWSortableBatchPosition& source) = delete; - TSortableBatchPosition operator= (TRWSortableBatchPosition&& source) = delete; + TSortableBatchPosition operator=(const TRWSortableBatchPosition& source) = delete; + TSortableBatchPosition operator=(TRWSortableBatchPosition& source) = delete; + TSortableBatchPosition operator=(TRWSortableBatchPosition&& source) = delete; TRWSortableBatchPosition BuildRWPosition(const bool needData, const bool deepCopy) const; @@ -281,6 +280,7 @@ public: explicit TFoundPosition(const ui32 pos) : Position(pos) { } + public: TString DebugString() const { TStringBuilder result; @@ -322,7 +322,8 @@ public: static std::optional<TFoundPosition> FindPosition(const std::shared_ptr<arrow::RecordBatch>& batch, const TSortableBatchPosition& forFound, const bool needGreater, const std::optional<ui32> includedStartPosition); - static std::optional<TSortableBatchPosition::TFoundPosition> FindPosition(TRWSortableBatchPosition& position, const ui64 posStart, const ui64 posFinish, const TSortableBatchPosition& forFound, const bool greater); + static std::optional<TSortableBatchPosition::TFoundPosition> FindPosition(TRWSortableBatchPosition& position, const ui64 posStart, + const ui64 posFinish, const TSortableBatchPosition& forFound, const bool greater); const TSortableScanData& GetData() const { AFL_VERIFY(!!Data); @@ -419,13 +420,13 @@ public: bool operator!=(const TSortableBatchPosition& item) const { return Compare(item) != std::partial_ordering::equivalent; } - }; class TIntervalPosition { private: TSortableBatchPosition Position; bool LeftIntervalInclude; + public: const TSortableBatchPosition& GetPosition() const { return Position; @@ -436,20 +437,18 @@ public: TIntervalPosition(TSortableBatchPosition&& position, const bool leftIntervalInclude) : Position(std::move(position)) , LeftIntervalInclude(leftIntervalInclude) { - } TIntervalPosition(const TSortableBatchPosition& position, const bool leftIntervalInclude) : Position(position) , LeftIntervalInclude(leftIntervalInclude) { - } bool operator<(const TIntervalPosition& item) const { std::partial_ordering cmp = Position.Compare(item.Position); if (cmp == std::partial_ordering::equivalent) { return (LeftIntervalInclude ? 1 : 0) < (item.LeftIntervalInclude ? 1 : 0); - } + } return cmp == std::partial_ordering::less; } @@ -464,6 +463,7 @@ public: class TIntervalPositions { private: std::vector<TIntervalPosition> Positions; + public: using const_iterator = std::vector<TIntervalPosition>::const_iterator; @@ -554,6 +554,7 @@ public: class TRWSortableBatchPosition: public TSortableBatchPosition, public TMoveOnly { private: using TBase = TSortableBatchPosition; + public: using TBase::TBase; @@ -576,10 +577,10 @@ public: class TAsymmetricPositionGuard: TNonCopyable { private: TRWSortableBatchPosition& Owner; + public: TAsymmetricPositionGuard(TRWSortableBatchPosition& owner) - : Owner(owner) - { + : Owner(owner) { } [[nodiscard]] bool InitSortingPosition(const i64 position) { @@ -608,8 +609,8 @@ public: // (-inf, it1), [it1, it2), [it2, it3), ..., [itLast, +inf) template <class TBordersIterator> - static std::vector<std::shared_ptr<arrow::RecordBatch>> SplitByBorders(const std::shared_ptr<arrow::RecordBatch>& batch, - const std::vector<std::string>& columnNames, TBordersIterator& it) { + static std::vector<std::shared_ptr<arrow::RecordBatch>> SplitByBorders( + const std::shared_ptr<arrow::RecordBatch>& batch, const std::vector<std::string>& columnNames, TBordersIterator& it) { std::vector<std::shared_ptr<arrow::RecordBatch>> result; if (!batch || batch->num_rows() == 0) { while (it.IsValid()) { @@ -662,6 +663,7 @@ public: private: typename TContainer::const_iterator Current; typename TContainer::const_iterator End; + public: TAssociatedContainerIterator(const TContainer& container) : Current(container.begin()) @@ -686,7 +688,8 @@ public: }; template <class TContainer> - static std::vector<std::shared_ptr<arrow::RecordBatch>> SplitByBordersInAssociativeContainer(const std::shared_ptr<arrow::RecordBatch>& batch, const std::vector<std::string>& columnNames, const TContainer& container) { + static std::vector<std::shared_ptr<arrow::RecordBatch>> SplitByBordersInAssociativeContainer( + const std::shared_ptr<arrow::RecordBatch>& batch, const std::vector<std::string>& columnNames, const TContainer& container) { TAssociatedContainerIterator<TContainer> it(container); return SplitByBorders(batch, columnNames, it); } @@ -696,6 +699,7 @@ public: private: typename TContainer::const_iterator Current; typename TContainer::const_iterator End; + public: TSequentialContainerIterator(const TContainer& container) : Current(container.begin()) @@ -720,7 +724,8 @@ public: }; template <class TContainer> - static std::vector<std::shared_ptr<arrow::RecordBatch>> SplitByBordersInSequentialContainer(const std::shared_ptr<arrow::RecordBatch>& batch, const std::vector<std::string>& columnNames, const TContainer& container) { + static std::vector<std::shared_ptr<arrow::RecordBatch>> SplitByBordersInSequentialContainer( + const std::shared_ptr<arrow::RecordBatch>& batch, const std::vector<std::string>& columnNames, const TContainer& container) { TSequentialContainerIterator<TContainer> it(container); return SplitByBorders(batch, columnNames, it); } @@ -755,7 +760,6 @@ public: bool operator()(const TSortableBatchPosition& value, const TIntervalPosition& pos) const { return value < pos.GetPosition(); } - }; void SkipToUpper(const TSortableBatchPosition& toPos) { @@ -770,4 +774,4 @@ public: } }; -} +} // namespace NKikimr::NArrow::NMerger diff --git a/ydb/core/formats/arrow/reader/result_builder.cpp b/ydb/core/formats/arrow/reader/result_builder.cpp index eed162c76d9..795e693fe2a 100644 --- a/ydb/core/formats/arrow/reader/result_builder.cpp +++ b/ydb/core/formats/arrow/reader/result_builder.cpp @@ -1,13 +1,12 @@ +#include "position.h" #include "result_builder.h" #include <ydb/library/actors/core/log.h> +#include <ydb/library/formats/arrow/validation/validation.h> #include <ydb/library/services/services.pb.h> -#include <ydb/library/formats/arrow/common/validation.h> #include <util/string/builder.h> -#include "position.h" - namespace NKikimr::NArrow::NMerger { void TRecordBatchBuilder::ValidateDataSchema(const std::shared_ptr<arrow::Schema>& schema) { @@ -15,8 +14,8 @@ void TRecordBatchBuilder::ValidateDataSchema(const std::shared_ptr<arrow::Schema } void TRecordBatchBuilder::AddRecord(const TCursor& position) { -// AFL_VERIFY_DEBUG(IsSameFieldsSequence(position.GetData().GetFields(), Fields)); -// AFL_TRACE(NKikimrServices::TX_COLUMNSHARD)("event", "record_add_on_read")("record", position.DebugJson()); + // AFL_VERIFY_DEBUG(IsSameFieldsSequence(position.GetData().GetFields(), Fields)); + // AFL_TRACE(NKikimrServices::TX_COLUMNSHARD)("event", "record_add_on_read")("record", position.DebugJson()); position.AppendPositionTo(Builders, MemoryBufferLimit ? &CurrentBytesUsed : nullptr); ++RecordsCount; } @@ -29,7 +28,8 @@ void TRecordBatchBuilder::AddRecord(const TRWSortableBatchPosition& position) { ++RecordsCount; } -bool TRecordBatchBuilder::IsSameFieldsSequence(const std::vector<std::shared_ptr<arrow::Field>>& f1, const std::vector<std::shared_ptr<arrow::Field>>& f2) { +bool TRecordBatchBuilder::IsSameFieldsSequence( + const std::vector<std::shared_ptr<arrow::Field>>& f1, const std::vector<std::shared_ptr<arrow::Field>>& f2) { if (f1.size() != f2.size()) { return false; } @@ -44,9 +44,9 @@ bool TRecordBatchBuilder::IsSameFieldsSequence(const std::vector<std::shared_ptr return true; } -TRecordBatchBuilder::TRecordBatchBuilder(const std::vector<std::shared_ptr<arrow::Field>>& fields, const std::optional<ui32> rowsCountExpectation /*= {}*/, const THashMap<std::string, ui64>& fieldDataSizePreallocated /*= {}*/) - : Fields(fields) -{ +TRecordBatchBuilder::TRecordBatchBuilder(const std::vector<std::shared_ptr<arrow::Field>>& fields, + const std::optional<ui32> rowsCountExpectation /*= {}*/, const THashMap<std::string, ui64>& fieldDataSizePreallocated /*= {}*/) + : Fields(fields) { AFL_VERIFY(Fields.size()); for (auto&& f : fields) { Builders.emplace_back(NArrow::MakeBuilder(f)); @@ -81,4 +81,4 @@ TString TRecordBatchBuilder::GetColumnNames() const { return result; } -} +} // namespace NKikimr::NArrow::NMerger diff --git a/ydb/core/formats/arrow/save_load/loader.cpp b/ydb/core/formats/arrow/save_load/loader.cpp index 0b87f220cd4..93a9243185a 100644 --- a/ydb/core/formats/arrow/save_load/loader.cpp +++ b/ydb/core/formats/arrow/save_load/loader.cpp @@ -1,6 +1,6 @@ #include "loader.h" -#include <ydb/library/formats/arrow/common/validation.h> +#include <ydb/library/formats/arrow/validation/validation.h> namespace NKikimr::NArrow::NAccessor { @@ -12,9 +12,8 @@ TString TColumnLoader::DebugString() const { return result; } -TColumnLoader::TColumnLoader(const NSerialization::TSerializerContainer& serializer, - const TConstructorContainer& accessorConstructor, const std::shared_ptr<arrow::Field>& resultField, - const std::shared_ptr<arrow::Scalar>& defaultValue, const ui32 columnId) +TColumnLoader::TColumnLoader(const NSerialization::TSerializerContainer& serializer, const TConstructorContainer& accessorConstructor, + const std::shared_ptr<arrow::Field>& resultField, const std::shared_ptr<arrow::Scalar>& defaultValue, const ui32 columnId) : Serializer(serializer) , AccessorConstructor(accessorConstructor) , ResultField(resultField) @@ -41,7 +40,8 @@ std::shared_ptr<IChunkedArray> TColumnLoader::ApplyVerified(const TString& dataS return BuildAccessor(dataStr, BuildAccessorContext(recordsCount)).DetachResult(); } -TConclusion<std::shared_ptr<IChunkedArray>> TColumnLoader::BuildAccessor(const TString& originalData, const TChunkConstructionData& chunkData) const { +TConclusion<std::shared_ptr<IChunkedArray>> TColumnLoader::BuildAccessor( + const TString& originalData, const TChunkConstructionData& chunkData) const { return AccessorConstructor->DeserializeFromString(originalData, chunkData); } diff --git a/ydb/core/formats/arrow/serializer/abstract.h b/ydb/core/formats/arrow/serializer/abstract.h index 9811aaaf0f2..c9e558f4ef0 100644 --- a/ydb/core/formats/arrow/serializer/abstract.h +++ b/ydb/core/formats/arrow/serializer/abstract.h @@ -1,15 +1,15 @@ #pragma once #include <ydb/core/protos/flat_scheme_op.pb.h> -#include <ydb/library/conclusion/status.h> -#include <ydb/services/metadata/abstract/request_features.h> -#include <ydb/services/bg_tasks/abstract/interface.h> #include <ydb/library/conclusion/result.h> -#include <ydb/library/formats/arrow/common/validation.h> +#include <ydb/library/conclusion/status.h> +#include <ydb/library/formats/arrow/validation/validation.h> +#include <ydb/services/bg_tasks/abstract/interface.h> +#include <ydb/services/metadata/abstract/request_features.h> -#include <contrib/libs/apache/arrow/cpp/src/arrow/status.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/record_batch.h> +#include <contrib/libs/apache/arrow/cpp/src/arrow/status.h> #include <util/generic/string.h> #include <util/string/builder.h> @@ -20,7 +20,8 @@ protected: virtual TString DoSerializeFull(const std::shared_ptr<arrow::RecordBatch>& batch) const = 0; virtual TString DoSerializePayload(const std::shared_ptr<arrow::RecordBatch>& batch) const = 0; virtual arrow::Result<std::shared_ptr<arrow::RecordBatch>> DoDeserialize(const TString& data) const = 0; - virtual arrow::Result<std::shared_ptr<arrow::RecordBatch>> DoDeserialize(const TString& data, const std::shared_ptr<arrow::Schema>& schema) const = 0; + virtual arrow::Result<std::shared_ptr<arrow::RecordBatch>> DoDeserialize( + const TString& data, const std::shared_ptr<arrow::Schema>& schema) const = 0; virtual TString DoDebugString() const { return ""; } @@ -29,6 +30,7 @@ protected: virtual TConclusionStatus DoDeserializeFromProto(const NKikimrSchemeOp::TOlapColumn::TSerializer& proto) = 0; virtual void DoSerializeToProto(NKikimrSchemeOp::TOlapColumn::TSerializer& proto) const = 0; + public: using TPtr = std::shared_ptr<ISerializer>; using TFactory = NObjectFactory::TObjectFactory<ISerializer, TString>; @@ -101,13 +103,12 @@ public: class TSerializerContainer: public NBackgroundTasks::TInterfaceProtoContainer<ISerializer> { private: using TBase = NBackgroundTasks::TInterfaceProtoContainer<ISerializer>; + public: using TBase::TBase; TSerializerContainer(const std::shared_ptr<ISerializer>& object) - : TBase(object) - { - + : TBase(object) { } bool IsCompatibleForExchange(const TSerializerContainer& item) const { @@ -170,4 +171,4 @@ public: } }; -} +} // namespace NKikimr::NArrow::NSerialization diff --git a/ydb/core/formats/arrow/serializer/native.cpp b/ydb/core/formats/arrow/serializer/native.cpp index 0676d5652bc..2c4f598089d 100644 --- a/ydb/core/formats/arrow/serializer/native.cpp +++ b/ydb/core/formats/arrow/serializer/native.cpp @@ -1,15 +1,16 @@ #include "native.h" -#include "stream.h" #include "parsing.h" +#include "stream.h" + #include <ydb/core/formats/arrow/dictionary/conversion.h> -#include <ydb/library/services/services.pb.h> #include <ydb/library/actors/core/log.h> -#include <ydb/library/formats/arrow/common/validation.h> +#include <ydb/library/formats/arrow/validation/validation.h> +#include <ydb/library/services/services.pb.h> -#include <contrib/libs/apache/arrow/cpp/src/arrow/ipc/dictionary.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/buffer.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/io/memory.h> +#include <contrib/libs/apache/arrow/cpp/src/arrow/ipc/dictionary.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/ipc/reader.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/ipc/writer.h> @@ -58,7 +59,8 @@ TString TNativeSerializer::DoSerializeFull(const std::shared_ptr<arrow::RecordBa return result; } -arrow::Result<std::shared_ptr<arrow::RecordBatch>> TNativeSerializer::DoDeserialize(const TString& data, const std::shared_ptr<arrow::Schema>& schema) const { +arrow::Result<std::shared_ptr<arrow::RecordBatch>> TNativeSerializer::DoDeserialize( + const TString& data, const std::shared_ptr<arrow::Schema>& schema) const { arrow::ipc::DictionaryMemo dictMemo; auto options = arrow::ipc::IpcReadOptions::Defaults(); options.use_threads = false; @@ -109,7 +111,8 @@ TString TNativeSerializer::DoSerializePayload(const std::shared_ptr<arrow::Recor return str; } -NKikimr::TConclusion<std::shared_ptr<arrow::util::Codec>> TNativeSerializer::BuildCodec(const arrow::Compression::type& cType, const std::optional<ui32> level) const { +NKikimr::TConclusion<std::shared_ptr<arrow::util::Codec>> TNativeSerializer::BuildCodec( + const arrow::Compression::type& cType, const std::optional<ui32> level) const { auto codec = NArrow::TStatusValidator::GetValid(arrow::util::Codec::Create(cType)); if (!codec) { return std::shared_ptr<arrow::util::Codec>(); @@ -194,4 +197,4 @@ void TNativeSerializer::DoSerializeToProto(NKikimrSchemeOp::TOlapColumn::TSerial } } -} +} // namespace NKikimr::NArrow::NSerialization diff --git a/ydb/core/formats/arrow/splitter/simple.cpp b/ydb/core/formats/arrow/splitter/simple.cpp index 1a3bb840c3b..5fac9ea2381 100644 --- a/ydb/core/formats/arrow/splitter/simple.cpp +++ b/ydb/core/formats/arrow/splitter/simple.cpp @@ -2,8 +2,8 @@ #include <ydb/core/formats/arrow/size_calcer.h> -#include <ydb/library/formats/arrow/common/validation.h> #include <ydb/library/formats/arrow/splitter/similar_packer.h> +#include <ydb/library/formats/arrow/validation/validation.h> #include <util/string/join.h> diff --git a/ydb/core/formats/arrow/ut/ya.make b/ydb/core/formats/arrow/ut/ya.make index 87c8e341530..6249b422456 100644 --- a/ydb/core/formats/arrow/ut/ya.make +++ b/ydb/core/formats/arrow/ut/ya.make @@ -8,6 +8,7 @@ PEERDIR( ydb/library/formats/arrow/simple_builder ydb/core/formats/arrow/program ydb/core/base + ydb/library/formats/arrow # for NYql::NUdf alloc stuff used in binary_json yql/essentials/public/udf/service/exception_policy diff --git a/ydb/core/kqp/compute_actor/kqp_compute_events.h b/ydb/core/kqp/compute_actor/kqp_compute_events.h index 28baedc4b34..4b9d3f29242 100644 --- a/ydb/core/kqp/compute_actor/kqp_compute_events.h +++ b/ydb/core/kqp/compute_actor/kqp_compute_events.h @@ -1,7 +1,7 @@ #pragma once #include <ydb/core/formats/arrow/arrow_helpers.h> -#include <ydb/library/formats/arrow/common/validation.h> +#include <ydb/library/formats/arrow/validation/validation.h> #include <ydb/core/kqp/common/kqp.h> #include <ydb/core/protos/tx_datashard.pb.h> #include <ydb/core/protos/data_events.pb.h> diff --git a/ydb/core/kqp/query_compiler/kqp_olap_compiler.cpp b/ydb/core/kqp/query_compiler/kqp_olap_compiler.cpp index 7d1e36ec225..df17a471197 100644 --- a/ydb/core/kqp/query_compiler/kqp_olap_compiler.cpp +++ b/ydb/core/kqp/query_compiler/kqp_olap_compiler.cpp @@ -426,6 +426,7 @@ const TTypedColumn ConvertJsonValueToColumn(const TKqpOlapJsonValue& jsonValueCa jsonValueCallable.Column(), jsonValueCallable.Path(), type); + jsonValueFunc->SetKernelName("JsonValue"); jsonValueFunc->SetKernelIdx(idx); return {command->GetColumn().GetId(), ctx.ConvertToBlockType(type)}; diff --git a/ydb/core/kqp/runtime/kqp_scan_data.h b/ydb/core/kqp/runtime/kqp_scan_data.h index b547fe3ccd0..592025dc592 100644 --- a/ydb/core/kqp/runtime/kqp_scan_data.h +++ b/ydb/core/kqp/runtime/kqp_scan_data.h @@ -14,6 +14,7 @@ #include <ydb/core/tablet_flat/flat_database.h> #include <ydb/library/yql/dq/actors/protos/dq_stats.pb.h> +#include <ydb/library/formats/arrow/validation/validation.h> #include <yql/essentials/minikql/computation/mkql_computation_node_holders.h> #include <ydb/library/actors/core/log.h> diff --git a/ydb/core/kqp/ut/olap/json_ut.cpp b/ydb/core/kqp/ut/olap/json_ut.cpp index afd0e48d554..c13df8ee725 100644 --- a/ydb/core/kqp/ut/olap/json_ut.cpp +++ b/ydb/core/kqp/ut/olap/json_ut.cpp @@ -368,6 +368,35 @@ Y_UNIT_TEST_SUITE(KqpOlapJson) { TScriptVariator(script).Execute(); } + Y_UNIT_TEST(FilterVariantsCount) { + TString script = 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_ACCESSOR_CONSTRUCTOR.CLASS_NAME`=`SUB_COLUMNS`, + `COLUMNS_LIMIT`=`$$1024|0|1$$`, `SPARSED_DETECTOR_KFF`=`$$0|10|1000$$`, `MEM_LIMIT_CHUNK`=`$$0|100|1000000$$`, `OTHERS_ALLOWED_FRACTION`=`$$0|0.5$$`) + ------ + DATA: + REPLACE INTO `/Root/ColumnTable` (Col1, Col2) VALUES(1u, JsonDocument('{"a" : "a1", "b" : "b1", "c" : "c1"}')), (2u, JsonDocument('{"a" : "a2"}')), + (3u, JsonDocument('{"b" : "b3", "d" : "d3"}')), (4u, JsonDocument('{"b" : "b4asdsasdaa", "a" : "a4"}')) + ------ + READ: SELECT COUNT(*) FROM `/Root/ColumnTable` WHERE JSON_VALUE(Col2, "$.a") = "a2"; + EXPECTED: [[1u]] + + )"; + TScriptVariator(script).Execute(); + } + Y_UNIT_TEST(SimpleVariants) { TString script = R"( SCHEMA: diff --git a/ydb/core/tx/columnshard/common/scalars.cpp b/ydb/core/tx/columnshard/common/scalars.cpp index d85622edeee..0a3ae473f56 100644 --- a/ydb/core/tx/columnshard/common/scalars.cpp +++ b/ydb/core/tx/columnshard/common/scalars.cpp @@ -1,7 +1,8 @@ #include "scalars.h" -#include <ydb/library/formats/arrow/switch_type.h> +#include <ydb/library/formats/arrow/switch/switch_type.h> #include <ydb/library/yverify_stream/yverify_stream.h> + #include <util/system/unaligned_mem.h> namespace NKikimr::NOlap { @@ -64,27 +65,27 @@ void ScalarToConstant(const arrow::Scalar& scalar, NKikimrSSA::TProgram_TConstan break; case arrow::Type::STRING: { auto& buffer = static_cast<const arrow::StringScalar&>(scalar).value; - value.SetText(TString(reinterpret_cast<const char *>(buffer->data()), buffer->size())); + value.SetText(TString(reinterpret_cast<const char*>(buffer->data()), buffer->size())); break; } case arrow::Type::LARGE_STRING: { auto& buffer = static_cast<const arrow::LargeStringScalar&>(scalar).value; - value.SetText(TString(reinterpret_cast<const char *>(buffer->data()), buffer->size())); + value.SetText(TString(reinterpret_cast<const char*>(buffer->data()), buffer->size())); break; } case arrow::Type::BINARY: { auto& buffer = static_cast<const arrow::BinaryScalar&>(scalar).value; - value.SetBytes(TString(reinterpret_cast<const char *>(buffer->data()), buffer->size())); + value.SetBytes(TString(reinterpret_cast<const char*>(buffer->data()), buffer->size())); break; } case arrow::Type::LARGE_BINARY: { auto& buffer = static_cast<const arrow::LargeBinaryScalar&>(scalar).value; - value.SetBytes(TString(reinterpret_cast<const char *>(buffer->data()), buffer->size())); + value.SetBytes(TString(reinterpret_cast<const char*>(buffer->data()), buffer->size())); break; } case arrow::Type::FIXED_SIZE_BINARY: { auto& buffer = static_cast<const arrow::FixedSizeBinaryScalar&>(scalar).value; - value.SetBytes(TString(reinterpret_cast<const char *>(buffer->data()), buffer->size())); + value.SetBytes(TString(reinterpret_cast<const char*>(buffer->data()), buffer->size())); break; } default: @@ -92,8 +93,7 @@ void ScalarToConstant(const arrow::Scalar& scalar, NKikimrSSA::TProgram_TConstan } } -std::shared_ptr<arrow::Scalar> ConstantToScalar(const NKikimrSSA::TProgram_TConstant& value, - const std::shared_ptr<arrow::DataType>& type) { +std::shared_ptr<arrow::Scalar> ConstantToScalar(const NKikimrSSA::TProgram_TConstant& value, const std::shared_ptr<arrow::DataType>& type) { switch (type->id()) { case arrow::Type::BOOL: return std::make_shared<arrow::BooleanScalar>(value.GetBool()); @@ -166,10 +166,8 @@ TString SerializeKeyScalar(const std::shared_ptr<arrow::Scalar>& key) { using T = typename TWrap::T; using TScalar = typename arrow::TypeTraits<T>::ScalarType; - if constexpr (std::is_same_v<T, arrow::StringType> || - std::is_same_v<T, arrow::BinaryType> || - std::is_same_v<T, arrow::LargeStringType> || - std::is_same_v<T, arrow::LargeBinaryType> || + if constexpr (std::is_same_v<T, arrow::StringType> || std::is_same_v<T, arrow::BinaryType> || + std::is_same_v<T, arrow::LargeStringType> || std::is_same_v<T, arrow::LargeBinaryType> || std::is_same_v<T, arrow::FixedSizeBinaryType>) { auto& buffer = static_cast<const TScalar&>(*key).value; out = buffer->ToString(); @@ -197,10 +195,8 @@ std::shared_ptr<arrow::Scalar> DeserializeKeyScalar(const TString& key, const st using T = typename TWrap::T; using TScalar = typename arrow::TypeTraits<T>::ScalarType; - if constexpr (std::is_same_v<T, arrow::StringType> || - std::is_same_v<T, arrow::BinaryType> || - std::is_same_v<T, arrow::LargeStringType> || - std::is_same_v<T, arrow::LargeBinaryType> || + if constexpr (std::is_same_v<T, arrow::StringType> || std::is_same_v<T, arrow::BinaryType> || + std::is_same_v<T, arrow::LargeStringType> || std::is_same_v<T, arrow::LargeBinaryType> || std::is_same_v<T, arrow::FixedSizeBinaryType>) { out = std::make_shared<TScalar>(arrow::Buffer::FromString(key), type); } else if constexpr (std::is_same_v<T, arrow::HalfFloatType>) { @@ -228,4 +224,4 @@ std::shared_ptr<arrow::Scalar> DeserializeKeyScalar(const TString& key, const st return out; } -} +} // namespace NKikimr::NOlap diff --git a/ydb/core/tx/columnshard/engines/changes/compaction/plain/column_cursor.cpp b/ydb/core/tx/columnshard/engines/changes/compaction/plain/column_cursor.cpp index 9fd0c4d301e..3467eb9da56 100644 --- a/ydb/core/tx/columnshard/engines/changes/compaction/plain/column_cursor.cpp +++ b/ydb/core/tx/columnshard/engines/changes/compaction/plain/column_cursor.cpp @@ -1,5 +1,5 @@ #include "column_cursor.h" -#include <ydb/library/formats/arrow/common/validation.h> +#include <ydb/library/formats/arrow/validation/validation.h> namespace NKikimr::NOlap::NCompaction { diff --git a/ydb/core/tx/columnshard/engines/changes/compaction/plain/column_portion_chunk.cpp b/ydb/core/tx/columnshard/engines/changes/compaction/plain/column_portion_chunk.cpp index 3db4127653b..8ba384ba2f7 100644 --- a/ydb/core/tx/columnshard/engines/changes/compaction/plain/column_portion_chunk.cpp +++ b/ydb/core/tx/columnshard/engines/changes/compaction/plain/column_portion_chunk.cpp @@ -1,7 +1,7 @@ #include "column_portion_chunk.h" #include <ydb/core/formats/arrow/accessor/plain/accessor.h> -#include <ydb/library/formats/arrow/common/validation.h> +#include <ydb/library/formats/arrow/validation/validation.h> #include <ydb/core/tx/columnshard/engines/changes/counters/general.h> #include <ydb/core/tx/columnshard/engines/storage/chunks/column.h> diff --git a/ydb/core/tx/columnshard/engines/changes/compaction/plain/logic.h b/ydb/core/tx/columnshard/engines/changes/compaction/plain/logic.h index 5b3c53f2eec..9e3ec9a7c18 100644 --- a/ydb/core/tx/columnshard/engines/changes/compaction/plain/logic.h +++ b/ydb/core/tx/columnshard/engines/changes/compaction/plain/logic.h @@ -1,8 +1,8 @@ #pragma once #include "column_cursor.h" -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> -#include <ydb/library/formats/arrow/accessor/common/const.h> +#include <ydb/core/formats/arrow/accessor/abstract/accessor.h> +#include <ydb/core/formats/arrow/accessor/common/const.h> #include <ydb/core/tx/columnshard/engines/changes/compaction/abstract/merger.h> namespace NKikimr::NOlap::NCompaction { diff --git a/ydb/core/tx/columnshard/engines/changes/compaction/sparsed/logic.h b/ydb/core/tx/columnshard/engines/changes/compaction/sparsed/logic.h index 9fc64606a09..a66e5e6e714 100644 --- a/ydb/core/tx/columnshard/engines/changes/compaction/sparsed/logic.h +++ b/ydb/core/tx/columnshard/engines/changes/compaction/sparsed/logic.h @@ -1,6 +1,6 @@ #pragma once -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> -#include <ydb/library/formats/arrow/accessor/common/const.h> +#include <ydb/core/formats/arrow/accessor/abstract/accessor.h> +#include <ydb/core/formats/arrow/accessor/common/const.h> #include <ydb/core/formats/arrow/accessor/sparsed/accessor.h> #include <ydb/core/tx/columnshard/engines/changes/compaction/abstract/merger.h> @@ -114,11 +114,10 @@ private: AFL_VERIFY(!Chunk || CurrentOwnedArray->GetAddress().GetGlobalStartPosition() + Chunk->GetFinishPosition() <= position); Chunk = &CurrentSparsedArray->GetSparsedChunk(CurrentOwnedArray->GetAddress().GetLocalIndex(position)); AFL_VERIFY(Chunk->GetRecordsCount()); - AFL_VERIFY(CurrentOwnedArray->GetAddress().GetGlobalStartPosition() + Chunk->GetStartPosition() <= position && + AFL_VERIFY(CurrentOwnedArray->GetAddress().GetGlobalStartPosition() <= position && position < CurrentOwnedArray->GetAddress().GetGlobalStartPosition() + Chunk->GetFinishPosition()) - ("pos", position)("start", Chunk->GetStartPosition())("finish", Chunk->GetFinishPosition())( - "shift", CurrentOwnedArray->GetAddress().GetGlobalStartPosition()); - ChunkStartGlobalPosition = CurrentOwnedArray->GetAddress().GetGlobalStartPosition() + Chunk->GetStartPosition(); + ("pos", position)("finish", Chunk->GetFinishPosition())("shift", CurrentOwnedArray->GetAddress().GetGlobalStartPosition()); + ChunkStartGlobalPosition = CurrentOwnedArray->GetAddress().GetGlobalStartPosition(); NextGlobalPosition = CurrentOwnedArray->GetAddress().GetGlobalStartPosition() + Chunk->GetFirstIndexNotDefault(); NextLocalPosition = 0; FinishGlobalPosition = CurrentOwnedArray->GetAddress().GetGlobalStartPosition() + Chunk->GetFinishPosition(); @@ -274,8 +273,7 @@ private: std::deque<TCursor> Cursors; std::list<TCursorPosition> CursorPositions; - virtual void DoStart( - const std::vector<std::shared_ptr<NArrow::NAccessor::IChunkedArray>>& input, TMergingContext& mergeContext) override; + virtual void DoStart(const std::vector<std::shared_ptr<NArrow::NAccessor::IChunkedArray>>& input, TMergingContext& mergeContext) override; virtual std::vector<TColumnPortionResult> DoExecute(const TChunkMergeContext& context, TMergingContext& mergeContext) override; diff --git a/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/builder.h b/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/builder.h index dcea4eb4b54..4ff263cdb3f 100644 --- a/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/builder.h +++ b/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/builder.h @@ -8,9 +8,6 @@ #include <ydb/core/tx/columnshard/engines/changes/compaction/abstract/merger.h> #include <ydb/core/tx/columnshard/engines/storage/chunks/column.h> -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> -#include <ydb/library/formats/arrow/accessor/common/const.h> - namespace NKikimr::NOlap::NCompaction::NSubColumns { class TMergedBuilder { diff --git a/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/iterator.h b/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/iterator.h index 042d7a88f81..d4995e7da5e 100644 --- a/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/iterator.h +++ b/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/iterator.h @@ -8,9 +8,6 @@ #include <ydb/core/tx/columnshard/engines/changes/compaction/abstract/merger.h> #include <ydb/core/tx/columnshard/engines/storage/chunks/column.h> -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> -#include <ydb/library/formats/arrow/accessor/common/const.h> - namespace NKikimr::NOlap::NCompaction::NSubColumns { class TChunksIterator { diff --git a/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/logic.h b/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/logic.h index d5d84781146..5d440f24e41 100644 --- a/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/logic.h +++ b/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/logic.h @@ -1,6 +1,7 @@ #pragma once #include "iterator.h" +#include <ydb/core/formats/arrow/accessor/common/const.h> #include <ydb/core/formats/arrow/accessor/plain/accessor.h> #include <ydb/core/formats/arrow/accessor/sparsed/accessor.h> #include <ydb/core/formats/arrow/accessor/sub_columns/accessor.h> @@ -8,9 +9,6 @@ #include <ydb/core/tx/columnshard/engines/changes/compaction/abstract/merger.h> #include <ydb/core/tx/columnshard/engines/storage/chunks/column.h> -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> -#include <ydb/library/formats/arrow/accessor/common/const.h> - namespace NKikimr::NOlap::NCompaction { class TSubColumnsMerger: public IColumnMerger { diff --git a/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/remap.h b/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/remap.h index de17796f146..15b631323c6 100644 --- a/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/remap.h +++ b/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/remap.h @@ -6,9 +6,6 @@ #include <ydb/core/tx/columnshard/engines/changes/compaction/abstract/merger.h> #include <ydb/core/tx/columnshard/engines/storage/chunks/column.h> -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> -#include <ydb/library/formats/arrow/accessor/common/const.h> - namespace NKikimr::NOlap::NCompaction::NSubColumns { class TRemapColumns { @@ -66,7 +63,8 @@ private: public: TRemapColumns() = default; - TOthersData::TFinishContext BuildRemapInfo(const std::vector<TDictStats::TRTStatsValue>& statsByKeyIndex, const TSettings& settings, const ui32 recordsCount) const; + TOthersData::TFinishContext BuildRemapInfo( + const std::vector<TDictStats::TRTStatsValue>& statsByKeyIndex, const TSettings& settings, const ui32 recordsCount) const; void RegisterColumnStats(const TDictStats& resultColumnStats) { ResultColumnStats = &resultColumnStats; diff --git a/ydb/core/tx/columnshard/engines/portions/column_record.h b/ydb/core/tx/columnshard/engines/portions/column_record.h index a063d8b621c..55ae8f164d6 100644 --- a/ydb/core/tx/columnshard/engines/portions/column_record.h +++ b/ydb/core/tx/columnshard/engines/portions/column_record.h @@ -9,7 +9,6 @@ #include <ydb/core/tx/columnshard/splitter/chunks.h> #include <ydb/library/accessor/accessor.h> -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> #include <ydb/library/formats/arrow/splitter/stats.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/array/array_base.h> diff --git a/ydb/core/tx/columnshard/engines/portions/data_accessor.cpp b/ydb/core/tx/columnshard/engines/portions/data_accessor.cpp index a34f6020281..73377ae1028 100644 --- a/ydb/core/tx/columnshard/engines/portions/data_accessor.cpp +++ b/ydb/core/tx/columnshard/engines/portions/data_accessor.cpp @@ -1,6 +1,8 @@ #include "constructor_meta.h" #include "data_accessor.h" +#include <ydb/core/formats/arrow/accessor/abstract/accessor.h> +#include <ydb/core/formats/arrow/accessor/composite/accessor.h> #include <ydb/core/formats/arrow/accessor/plain/accessor.h> #include <ydb/core/formats/arrow/common/container.h> #include <ydb/core/tx/columnshard/blobs_reader/task.h> @@ -10,8 +12,6 @@ #include <ydb/core/tx/columnshard/engines/storage/chunks/column.h> #include <ydb/core/tx/columnshard/engines/storage/chunks/data.h> -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> -#include <ydb/library/formats/arrow/accessor/composite/accessor.h> #include <ydb/library/formats/arrow/simple_arrays_cache.h> namespace NKikimr::NOlap { diff --git a/ydb/core/tx/columnshard/engines/predicate/filter.cpp b/ydb/core/tx/columnshard/engines/predicate/filter.cpp index 2b2e55cec36..1dc7b77bff5 100644 --- a/ydb/core/tx/columnshard/engines/predicate/filter.cpp +++ b/ydb/core/tx/columnshard/engines/predicate/filter.cpp @@ -3,6 +3,7 @@ #include <ydb/core/formats/arrow/serializer/native.h> #include <ydb/library/actors/core/log.h> +#include <ydb/library/formats/arrow/switch/switch_type.h> namespace NKikimr::NOlap { diff --git a/ydb/core/tx/columnshard/engines/predicate/predicate.cpp b/ydb/core/tx/columnshard/engines/predicate/predicate.cpp index 94ebaf4b978..e00ce54385d 100644 --- a/ydb/core/tx/columnshard/engines/predicate/predicate.cpp +++ b/ydb/core/tx/columnshard/engines/predicate/predicate.cpp @@ -6,7 +6,7 @@ #include <ydb/library/actors/core/log.h> #include <ydb/library/formats/arrow/arrow_helpers.h> -#include <ydb/library/formats/arrow/switch_type.h> +#include <ydb/library/formats/arrow/switch/switch_type.h> namespace NKikimr::NOlap { diff --git a/ydb/core/tx/columnshard/engines/reader/common_reader/iterator/fetched_data.cpp b/ydb/core/tx/columnshard/engines/reader/common_reader/iterator/fetched_data.cpp index 5d29a4d4c94..8e141196f06 100644 --- a/ydb/core/tx/columnshard/engines/reader/common_reader/iterator/fetched_data.cpp +++ b/ydb/core/tx/columnshard/engines/reader/common_reader/iterator/fetched_data.cpp @@ -2,7 +2,7 @@ #include <ydb/core/formats/arrow/accessor/plain/accessor.h> -#include <ydb/library/formats/arrow/common/validation.h> +#include <ydb/library/formats/arrow/validation/validation.h> #include <ydb/library/formats/arrow/simple_arrays_cache.h> namespace NKikimr::NOlap::NReader::NCommon { diff --git a/ydb/core/tx/columnshard/engines/reader/common_reader/iterator/fetching.cpp b/ydb/core/tx/columnshard/engines/reader/common_reader/iterator/fetching.cpp index 292b7758ce4..71b31f7ca96 100644 --- a/ydb/core/tx/columnshard/engines/reader/common_reader/iterator/fetching.cpp +++ b/ydb/core/tx/columnshard/engines/reader/common_reader/iterator/fetching.cpp @@ -168,6 +168,7 @@ TConclusion<bool> TProgramStep::DoExecuteInplace(const std::shared_ptr<IDataSour if (result.IsFail()) { return result; } + source->GetStageData().GetTable()->Remove(Step.GetColumnsToDrop()); return true; } diff --git a/ydb/core/tx/columnshard/engines/reader/simple_reader/iterator/fetching.cpp b/ydb/core/tx/columnshard/engines/reader/simple_reader/iterator/fetching.cpp index c3680f76618..8a6e9428e85 100644 --- a/ydb/core/tx/columnshard/engines/reader/simple_reader/iterator/fetching.cpp +++ b/ydb/core/tx/columnshard/engines/reader/simple_reader/iterator/fetching.cpp @@ -158,7 +158,7 @@ TConclusion<bool> TBuildResultStep::DoExecuteInplace(const std::shared_ptr<IData if (!source->GetStageResult().IsEmpty()) { resultBatch = source->GetStageResult().GetBatch()->BuildTableVerified(contextTableConstruct); if (auto filter = source->GetStageResult().GetNotAppliedFilter()) { - filter->Apply(resultBatch, NArrow::TColumnFilter::TApplyContext(StartIndex, RecordsCount).SetTrySlices(true)); + AFL_VERIFY(filter->Apply(resultBatch, NArrow::TColumnFilter::TApplyContext(StartIndex, RecordsCount).SetTrySlices(true))); } } NActors::TActivationContext::AsActorContext().Send(context->GetCommonContext()->GetScanActorId(), diff --git a/ydb/core/tx/columnshard/engines/reader/sys_view/abstract/iterator.cpp b/ydb/core/tx/columnshard/engines/reader/sys_view/abstract/iterator.cpp index 70355350f84..81ee1b3281a 100644 --- a/ydb/core/tx/columnshard/engines/reader/sys_view/abstract/iterator.cpp +++ b/ydb/core/tx/columnshard/engines/reader/sys_view/abstract/iterator.cpp @@ -45,7 +45,7 @@ TConclusion<std::shared_ptr<TPartialReadResult>> TStatsIteratorBase::GetBatch() { NArrow::TColumnFilter filter = ReadMetadata->GetPKRangesFilter().BuildFilter(originalBatch); - filter.Apply(originalBatch); + AFL_VERIFY(filter.Apply(originalBatch)); } // Leave only requested columns diff --git a/ydb/core/tx/columnshard/engines/scheme/column/info.h b/ydb/core/tx/columnshard/engines/scheme/column/info.h index a7eaea7933a..a081754ef9b 100644 --- a/ydb/core/tx/columnshard/engines/scheme/column/info.h +++ b/ydb/core/tx/columnshard/engines/scheme/column/info.h @@ -1,6 +1,6 @@ #pragma once #include <ydb/core/formats/arrow/accessor/abstract/constructor.h> -#include <ydb/library/formats/arrow/common/validation.h> +#include <ydb/library/formats/arrow/validation/validation.h> #include <ydb/core/formats/arrow/dictionary/object.h> #include <ydb/core/formats/arrow/save_load/loader.h> #include <ydb/core/formats/arrow/save_load/saver.h> diff --git a/ydb/core/tx/columnshard/engines/scheme/column_features.h b/ydb/core/tx/columnshard/engines/scheme/column_features.h index d1c9507fa19..d988202ebc5 100644 --- a/ydb/core/tx/columnshard/engines/scheme/column_features.h +++ b/ydb/core/tx/columnshard/engines/scheme/column_features.h @@ -7,7 +7,7 @@ #include <ydb/core/tx/columnshard/blobs_action/abstract/storage.h> #include <ydb/core/tx/columnshard/blobs_action/abstract/storages_manager.h> #include <ydb/core/tx/columnshard/splitter/abstract/chunks.h> -#include <ydb/library/formats/arrow/common/validation.h> +#include <ydb/library/formats/arrow/validation/validation.h> #include <ydb/core/tx/columnshard/engines/scheme/abstract/index_info.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/type.h> diff --git a/ydb/core/tx/columnshard/engines/scheme/defaults/common/scalar.cpp b/ydb/core/tx/columnshard/engines/scheme/defaults/common/scalar.cpp index 7131d946904..8146ac1d38c 100644 --- a/ydb/core/tx/columnshard/engines/scheme/defaults/common/scalar.cpp +++ b/ydb/core/tx/columnshard/engines/scheme/defaults/common/scalar.cpp @@ -1,8 +1,12 @@ #include "scalar.h" + #include <ydb/core/formats/arrow/arrow_helpers.h> + #include <ydb/library/actors/core/log.h> -#include <util/string/cast.h> +#include <ydb/library/formats/arrow/switch/switch_type.h> + #include <util/string/builder.h> +#include <util/string/cast.h> namespace NKikimr::NOlap { @@ -47,8 +51,7 @@ NKikimrColumnShardColumnDefaults::TColumnDefault TColumnDefaultScalarValue::Seri case arrow::Type::FLOAT: resultScalar.SetFloat(static_cast<const arrow::FloatScalar&>(scalar).value); break; - case arrow::Type::TIMESTAMP: - { + case arrow::Type::TIMESTAMP: { auto* ts = resultScalar.MutableTimestamp(); ts->SetValue(static_cast<const arrow::TimestampScalar&>(scalar).value); ts->SetUnit(static_cast<const arrow::TimestampType&>(*scalar.type).unit()); @@ -142,4 +145,4 @@ TString TColumnDefaultScalarValue::DebugString() const { return TStringBuilder() << Value->ToString(); } -}
\ No newline at end of file +} // namespace NKikimr::NOlap diff --git a/ydb/core/tx/columnshard/engines/scheme/tiering/tier_info.h b/ydb/core/tx/columnshard/engines/scheme/tiering/tier_info.h index 806459c93eb..2bf46e91da5 100644 --- a/ydb/core/tx/columnshard/engines/scheme/tiering/tier_info.h +++ b/ydb/core/tx/columnshard/engines/scheme/tiering/tier_info.h @@ -7,7 +7,7 @@ #include <ydb/core/tx/columnshard/common/scalars.h> #include <ydb/core/tx/tiering/tier/identifier.h> -#include <ydb/library/formats/arrow/common/validation.h> +#include <ydb/library/formats/arrow/validation/validation.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/util/compression.h> #include <util/generic/hash_set.h> diff --git a/ydb/core/tx/columnshard/engines/storage/indexes/bloom/checker.cpp b/ydb/core/tx/columnshard/engines/storage/indexes/bloom/checker.cpp index 1613bd10e7d..1afa71281e5 100644 --- a/ydb/core/tx/columnshard/engines/storage/indexes/bloom/checker.cpp +++ b/ydb/core/tx/columnshard/engines/storage/indexes/bloom/checker.cpp @@ -1,6 +1,6 @@ #include "checker.h" #include <ydb/core/formats/arrow/serializer/abstract.h> -#include <ydb/library/formats/arrow/common/validation.h> +#include <ydb/library/formats/arrow/validation/validation.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/array/array_primitive.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/record_batch.h> diff --git a/ydb/core/tx/columnshard/engines/storage/indexes/bloom_ngramm/checker.cpp b/ydb/core/tx/columnshard/engines/storage/indexes/bloom_ngramm/checker.cpp index 9eed9642375..0161bde3e35 100644 --- a/ydb/core/tx/columnshard/engines/storage/indexes/bloom_ngramm/checker.cpp +++ b/ydb/core/tx/columnshard/engines/storage/indexes/bloom_ngramm/checker.cpp @@ -3,7 +3,7 @@ #include <ydb/core/formats/arrow/serializer/abstract.h> #include <ydb/core/tx/columnshard/engines/storage/indexes/bloom/checker.h> -#include <ydb/library/formats/arrow/common/validation.h> +#include <ydb/library/formats/arrow/validation/validation.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/array/array_primitive.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/record_batch.h> diff --git a/ydb/core/tx/columnshard/engines/storage/indexes/count_min_sketch/checker.cpp b/ydb/core/tx/columnshard/engines/storage/indexes/count_min_sketch/checker.cpp index aa40668897d..496656d0c57 100644 --- a/ydb/core/tx/columnshard/engines/storage/indexes/count_min_sketch/checker.cpp +++ b/ydb/core/tx/columnshard/engines/storage/indexes/count_min_sketch/checker.cpp @@ -1,6 +1,6 @@ #include "checker.h" #include <ydb/core/formats/arrow/serializer/abstract.h> -#include <ydb/library/formats/arrow/common/validation.h> +#include <ydb/library/formats/arrow/validation/validation.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/array/array_primitive.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/record_batch.h> diff --git a/ydb/core/tx/columnshard/splitter/abstract/chunk_meta.cpp b/ydb/core/tx/columnshard/splitter/abstract/chunk_meta.cpp index b0f368afb1a..f3564c963b8 100644 --- a/ydb/core/tx/columnshard/splitter/abstract/chunk_meta.cpp +++ b/ydb/core/tx/columnshard/splitter/abstract/chunk_meta.cpp @@ -1,15 +1,16 @@ #include "chunk_meta.h" + +#include <ydb/core/formats/arrow/accessor/abstract/accessor.h> #include <ydb/core/formats/arrow/arrow_helpers.h> #include <ydb/core/formats/arrow/size_calcer.h> namespace NKikimr::NOlap { -TSimpleChunkMeta::TSimpleChunkMeta( - const std::shared_ptr<NArrow::NAccessor::IChunkedArray>& column) { +TSimpleChunkMeta::TSimpleChunkMeta(const std::shared_ptr<NArrow::NAccessor::IChunkedArray>& column) { Y_ABORT_UNLESS(column); Y_ABORT_UNLESS(column->GetRecordsCount()); RecordsCount = column->GetRecordsCount(); RawBytes = column->GetRawSizeVerified(); } -} +} // namespace NKikimr::NOlap diff --git a/ydb/core/tx/columnshard/splitter/abstract/chunk_meta.h b/ydb/core/tx/columnshard/splitter/abstract/chunk_meta.h index 6b9964a5d91..c7011287a49 100644 --- a/ydb/core/tx/columnshard/splitter/abstract/chunk_meta.h +++ b/ydb/core/tx/columnshard/splitter/abstract/chunk_meta.h @@ -1,14 +1,15 @@ #pragma once -#include <ydb/library/formats/arrow/accessor/abstract/accessor.h> - -#include <contrib/libs/apache/arrow/cpp/src/arrow/scalar.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/array/array_base.h> - +#include <contrib/libs/apache/arrow/cpp/src/arrow/scalar.h> #include <util/system/types.h> #include <util/system/yassert.h> -#include <optional> #include <memory> +#include <optional> + +namespace NKikimr::NArrow::NAccessor { +class IChunkedArray; +} namespace NKikimr::NOlap { @@ -17,6 +18,7 @@ protected: ui32 RecordsCount = 0; ui32 RawBytes = 0; TSimpleChunkMeta() = default; + public: TSimpleChunkMeta(const std::shared_ptr<NArrow::NAccessor::IChunkedArray>& column); @@ -30,6 +32,5 @@ public: ui32 GetRawBytes() const { return RawBytes; } - }; -} +} // namespace NKikimr::NOlap diff --git a/ydb/core/tx/columnshard/test_helper/columnshard_ut_common.h b/ydb/core/tx/columnshard/test_helper/columnshard_ut_common.h index 31e2f28869b..1be3d35827e 100644 --- a/ydb/core/tx/columnshard/test_helper/columnshard_ut_common.h +++ b/ydb/core/tx/columnshard/test_helper/columnshard_ut_common.h @@ -13,10 +13,11 @@ #include <ydb/core/tx/long_tx_service/public/types.h> #include <ydb/core/tx/tiering/manager.h> -#include <ydb-cpp-sdk/client/value/value.h> +#include <ydb/library/formats/arrow/switch/switch_type.h> #include <ydb/services/metadata/abstract/fetcher.h> #include <library/cpp/testing/unittest/registar.h> +#include <ydb-cpp-sdk/client/value/value.h> namespace NKikimr::NOlap { struct TIndexInfo; @@ -34,7 +35,7 @@ inline TPrivateEvent* TryGetPrivateEvent(TAutoPtr<IEventHandle>& ev) { return dynamic_cast<TPrivateEvent*>(ev->StaticCastAsLocal<IEventBase>()); } -class TTester : public TNonCopyable { +class TTester: public TNonCopyable { public: static constexpr const ui64 FAKE_SCHEMESHARD_TABLET_ID = 4200; @@ -57,11 +58,12 @@ struct TTestSchema { NKikimrSchemeOp::TS3Settings S3 = FakeS3(); TStorageTier(const TString& name = {}) - : Name(name) - {} + : Name(name) { + } TString DebugString() const { - return TStringBuilder() << "{Column=" << TtlColumn << ";EvictAfter=" << EvictAfter.value_or(TDuration::Zero()) << ";Name=" << Name << ";Codec=" << Codec << "};"; + return TStringBuilder() << "{Column=" << TtlColumn << ";EvictAfter=" << EvictAfter.value_or(TDuration::Zero()) << ";Name=" << Name + << ";Codec=" << Codec << "};"; } NKikimrSchemeOp::EColumnCodec GetCodecId() const { @@ -118,7 +120,7 @@ struct TTestSchema { } }; - struct TTableSpecials : public TStorageTier { + struct TTableSpecials: public TStorageTier { public: std::vector<TStorageTier> Tiers; bool WaitEmptyAfter = false; @@ -149,7 +151,7 @@ struct TTestSchema { } TString DebugString() const { - auto result = TStringBuilder() << "WaitEmptyAfter=" << WaitEmptyAfter << ";Tiers="; + auto result = TStringBuilder() << "WaitEmptyAfter=" << WaitEmptyAfter << ";Tiers="; for (auto&& tier : Tiers) { result << "{" << tier.DebugString() << "}"; } @@ -166,18 +168,13 @@ struct TTestSchema { }; using TTestColumn = NArrow::NTest::TTestColumn; static auto YdbSchema(const TTestColumn& firstKeyItem = TTestColumn("timestamp", TTypeInfo(NTypeIds::Timestamp))) { - std::vector<TTestColumn> schema = { - // PK - firstKeyItem, - TTestColumn("resource_type", TTypeInfo(NTypeIds::Utf8) ), + std::vector<TTestColumn> schema = { // PK + firstKeyItem, TTestColumn("resource_type", TTypeInfo(NTypeIds::Utf8)), TTestColumn("resource_id", TTypeInfo(NTypeIds::Utf8)).SetAccessorClassName("SPARSED"), - TTestColumn("uid", TTypeInfo(NTypeIds::Utf8) ).SetStorageId("__MEMORY"), - TTestColumn("level", TTypeInfo(NTypeIds::Int32) ), - TTestColumn("message", TTypeInfo(NTypeIds::Utf8) ).SetStorageId("__MEMORY"), - TTestColumn("json_payload", TTypeInfo(NTypeIds::Json) ), - TTestColumn("ingested_at", TTypeInfo(NTypeIds::Timestamp) ), - TTestColumn("saved_at", TTypeInfo(NTypeIds::Timestamp) ), - TTestColumn("request_id", TTypeInfo(NTypeIds::Utf8) ) + TTestColumn("uid", TTypeInfo(NTypeIds::Utf8)).SetStorageId("__MEMORY"), TTestColumn("level", TTypeInfo(NTypeIds::Int32)), + TTestColumn("message", TTypeInfo(NTypeIds::Utf8)).SetStorageId("__MEMORY"), TTestColumn("json_payload", TTypeInfo(NTypeIds::Json)), + TTestColumn("ingested_at", TTypeInfo(NTypeIds::Timestamp)), TTestColumn("saved_at", TTypeInfo(NTypeIds::Timestamp)), + TTestColumn("request_id", TTypeInfo(NTypeIds::Utf8)) }; return schema; }; @@ -199,70 +196,52 @@ struct TTestSchema { } static auto YdbExoticSchema() { - std::vector<TTestColumn> schema = { - // PK - TTestColumn("timestamp", TTypeInfo(NTypeIds::Timestamp) ), + std::vector<TTestColumn> schema = { // PK + TTestColumn("timestamp", TTypeInfo(NTypeIds::Timestamp)), TTestColumn("resource_type", TTypeInfo(NTypeIds::Utf8)).SetAccessorClassName("SPARSED"), - TTestColumn("resource_id", TTypeInfo(NTypeIds::Utf8) ), - TTestColumn("uid", TTypeInfo(NTypeIds::Utf8) ).SetStorageId("__MEMORY"), + TTestColumn("resource_id", TTypeInfo(NTypeIds::Utf8)), TTestColumn("uid", TTypeInfo(NTypeIds::Utf8)).SetStorageId("__MEMORY"), // - TTestColumn("level", TTypeInfo(NTypeIds::Int32) ), - TTestColumn("message", TTypeInfo(NTypeIds::String4k) ).SetStorageId("__MEMORY"), - TTestColumn("json_payload", TTypeInfo(NTypeIds::JsonDocument) ), - TTestColumn("ingested_at", TTypeInfo(NTypeIds::Timestamp) ), - TTestColumn("saved_at", TTypeInfo(NTypeIds::Timestamp) ), + TTestColumn("level", TTypeInfo(NTypeIds::Int32)), TTestColumn("message", TTypeInfo(NTypeIds::String4k)).SetStorageId("__MEMORY"), + TTestColumn("json_payload", TTypeInfo(NTypeIds::JsonDocument)), TTestColumn("ingested_at", TTypeInfo(NTypeIds::Timestamp)), + TTestColumn("saved_at", TTypeInfo(NTypeIds::Timestamp)), TTestColumn("request_id", TTypeInfo(NTypeIds::Yson)).SetAccessorClassName("SPARSED") }; return schema; }; static auto YdbPkSchema() { - std::vector<TTestColumn> schema = { - TTestColumn("timestamp", TTypeInfo(NTypeIds::Timestamp) ), - TTestColumn("resource_type", TTypeInfo(NTypeIds::Utf8) ).SetStorageId("__MEMORY"), + std::vector<TTestColumn> schema = { TTestColumn("timestamp", TTypeInfo(NTypeIds::Timestamp)), + TTestColumn("resource_type", TTypeInfo(NTypeIds::Utf8)).SetStorageId("__MEMORY"), TTestColumn("resource_id", TTypeInfo(NTypeIds::Utf8)).SetAccessorClassName("SPARSED"), - TTestColumn("uid", TTypeInfo(NTypeIds::Utf8) ).SetStorageId("__MEMORY") - }; + TTestColumn("uid", TTypeInfo(NTypeIds::Utf8)).SetStorageId("__MEMORY") }; return schema; } static auto YdbAllTypesSchema() { - std::vector<TTestColumn> schema = { - TTestColumn("ts", TTypeInfo(NTypeIds::Timestamp) ), + std::vector<TTestColumn> schema = { TTestColumn("ts", TTypeInfo(NTypeIds::Timestamp)), - TTestColumn( "i8", TTypeInfo(NTypeIds::Int8) ), - TTestColumn( "i16", TTypeInfo(NTypeIds::Int16) ), - TTestColumn( "i32", TTypeInfo(NTypeIds::Int32) ), - TTestColumn( "i64", TTypeInfo(NTypeIds::Int64) ), - TTestColumn( "u8", TTypeInfo(NTypeIds::Uint8) ), - TTestColumn( "u16", TTypeInfo(NTypeIds::Uint16) ), - TTestColumn( "u32", TTypeInfo(NTypeIds::Uint32) ), - TTestColumn( "u64", TTypeInfo(NTypeIds::Uint64) ), - TTestColumn( "float", TTypeInfo(NTypeIds::Float) ), - TTestColumn( "double", TTypeInfo(NTypeIds::Double) ), + TTestColumn("i8", TTypeInfo(NTypeIds::Int8)), TTestColumn("i16", TTypeInfo(NTypeIds::Int16)), + TTestColumn("i32", TTypeInfo(NTypeIds::Int32)), TTestColumn("i64", TTypeInfo(NTypeIds::Int64)), + TTestColumn("u8", TTypeInfo(NTypeIds::Uint8)), TTestColumn("u16", TTypeInfo(NTypeIds::Uint16)), + TTestColumn("u32", TTypeInfo(NTypeIds::Uint32)), TTestColumn("u64", TTypeInfo(NTypeIds::Uint64)), + TTestColumn("float", TTypeInfo(NTypeIds::Float)), TTestColumn("double", TTypeInfo(NTypeIds::Double)), - TTestColumn("byte", TTypeInfo(NTypeIds::Byte) ), + TTestColumn("byte", TTypeInfo(NTypeIds::Byte)), //{ "bool", TTypeInfo(NTypeIds::Bool) }, //{ "decimal", TTypeInfo(NTypeIds::Decimal) }, //{ "dynum", TTypeInfo(NTypeIds::DyNumber) }, - TTestColumn( "date", TTypeInfo(NTypeIds::Date) ), - TTestColumn( "datetime", TTypeInfo(NTypeIds::Datetime) ), + TTestColumn("date", TTypeInfo(NTypeIds::Date)), TTestColumn("datetime", TTypeInfo(NTypeIds::Datetime)), //{ "interval", TTypeInfo(NTypeIds::Interval) }, - TTestColumn("text", TTypeInfo(NTypeIds::Text) ), - TTestColumn("bytes", TTypeInfo(NTypeIds::Bytes) ), - TTestColumn("yson", TTypeInfo(NTypeIds::Yson) ), - TTestColumn("json", TTypeInfo(NTypeIds::Json) ), - TTestColumn("jsondoc", TTypeInfo(NTypeIds::JsonDocument) ) - }; + TTestColumn("text", TTypeInfo(NTypeIds::Text)), TTestColumn("bytes", TTypeInfo(NTypeIds::Bytes)), + TTestColumn("yson", TTypeInfo(NTypeIds::Yson)), TTestColumn("json", TTypeInfo(NTypeIds::Json)), + TTestColumn("jsondoc", TTypeInfo(NTypeIds::JsonDocument)) }; return schema; }; - static void InitSchema(const std::vector<NArrow::NTest::TTestColumn>& columns, - const std::vector<NArrow::NTest::TTestColumn>& pk, - const TTableSpecials& specials, - NKikimrSchemeOp::TColumnTableSchema* schema); + static void InitSchema(const std::vector<NArrow::NTest::TTestColumn>& columns, const std::vector<NArrow::NTest::TTestColumn>& pk, + const TTableSpecials& specials, NKikimrSchemeOp::TColumnTableSchema* schema); static bool InitTiersAndTtl(const TTableSpecials& specials, NKikimrSchemeOp::TColumnDataLifeCycle* ttlSettings) { ttlSettings->SetVersion(1); @@ -286,16 +265,14 @@ struct TTestSchema { } static TString CreateTableTxBody(ui64 pathId, const std::vector<NArrow::NTest::TTestColumn>& columns, - const std::vector<NArrow::NTest::TTestColumn>& pk, - const TTableSpecials& specialsExt = {}, ui64 generation = 0) - { + const std::vector<NArrow::NTest::TTestColumn>& pk, const TTableSpecials& specialsExt = {}, ui64 generation = 0) { auto specials = specialsExt; NKikimrTxColumnShard::TSchemaTxBody tx; tx.MutableSeqNo()->SetGeneration(generation); auto* table = tx.MutableEnsureTables()->AddTables(); table->SetPathId(pathId); - { // preset + { // preset auto* preset = table->MutableSchemaPreset(); preset->SetId(1); preset->SetName("default"); @@ -314,8 +291,7 @@ struct TTestSchema { } static TString CreateInitShardTxBody(ui64 pathId, const std::vector<NArrow::NTest::TTestColumn>& columns, - const std::vector<NArrow::NTest::TTestColumn>& pk, - const TTableSpecials& specials = {}, const TString& ownerPath = "/Root/olap") { + const std::vector<NArrow::NTest::TTestColumn>& pk, const TTableSpecials& specials = {}, const TString& ownerPath = "/Root/olap") { NKikimrTxColumnShard::TSchemaTxBody tx; auto* table = tx.MutableInitShard()->AddTables(); tx.MutableInitShard()->SetOwnerPath(ownerPath); @@ -332,8 +308,7 @@ struct TTestSchema { } static TString CreateStandaloneTableTxBody(ui64 pathId, const std::vector<NArrow::NTest::TTestColumn>& columns, - const std::vector<NArrow::NTest::TTestColumn>& pk, - const TTableSpecials& specials = {}) { + const std::vector<NArrow::NTest::TTestColumn>& pk, const TTableSpecials& specials = {}) { NKikimrTxColumnShard::TSchemaTxBody tx; auto* table = tx.MutableEnsureTables()->AddTables(); table->SetPathId(pathId); @@ -446,16 +421,16 @@ bool WriteData(TTestBasicRuntime& runtime, TActorId& sender, const ui64 writeId, const std::vector<NArrow::NTest::TTestColumn>& ydbSchema, bool waitResult = true, std::vector<ui64>* writeIds = nullptr, const NEvWrite::EModificationType mType = NEvWrite::EModificationType::Upsert, const ui64 lockId = 1); -std::optional<ui64> WriteData(TTestBasicRuntime& runtime, TActorId& sender, const NLongTxService::TLongTxId& longTxId, - ui64 tableId, const ui64 writePartId, const TString& data, - const std::vector<NArrow::NTest::TTestColumn>& ydbSchema, const NEvWrite::EModificationType mType = NEvWrite::EModificationType::Upsert); +std::optional<ui64> WriteData(TTestBasicRuntime& runtime, TActorId& sender, const NLongTxService::TLongTxId& longTxId, ui64 tableId, + const ui64 writePartId, const TString& data, const std::vector<NArrow::NTest::TTestColumn>& ydbSchema, + const NEvWrite::EModificationType mType = NEvWrite::EModificationType::Upsert); ui32 WaitWriteResult(TTestBasicRuntime& runtime, ui64 shardId, std::vector<ui64>* writeIds = nullptr); -void ScanIndexStats(TTestBasicRuntime& runtime, TActorId& sender, const std::vector<ui64>& pathIds, - NOlap::TSnapshot snap, ui64 scanId = 0); +void ScanIndexStats(TTestBasicRuntime& runtime, TActorId& sender, const std::vector<ui64>& pathIds, NOlap::TSnapshot snap, ui64 scanId = 0); -void ProposeCommit(TTestBasicRuntime& runtime, TActorId& sender, ui64 shardId, ui64 txId, const std::vector<ui64>& writeIds, const ui64 lockId = 1); +void ProposeCommit( + TTestBasicRuntime& runtime, TActorId& sender, ui64 shardId, ui64 txId, const std::vector<ui64>& writeIds, const ui64 lockId = 1); void ProposeCommit(TTestBasicRuntime& runtime, TActorId& sender, ui64 txId, const std::vector<ui64>& writeIds, const ui64 lockId = 1); void PlanCommit(TTestBasicRuntime& runtime, TActorId& sender, ui64 shardId, ui64 planStep, const TSet<ui64>& txIds); @@ -476,131 +451,130 @@ struct TTestBlobOptions { }; TCell MakeTestCell(const TTypeInfo& typeInfo, ui32 value, std::vector<TString>& mem); -TString MakeTestBlob(std::pair<ui64, ui64> range, const std::vector<NArrow::NTest::TTestColumn>& columns, - const TTestBlobOptions& options = {}, const std::set<std::string>& notNullColumns = {}); -TSerializedTableRange MakeTestRange(std::pair<ui64, ui64> range, bool inclusiveFrom, bool inclusiveTo, - const std::vector<NArrow::NTest::TTestColumn>& columns); +TString MakeTestBlob(std::pair<ui64, ui64> range, const std::vector<NArrow::NTest::TTestColumn>& columns, const TTestBlobOptions& options = {}, + const std::set<std::string>& notNullColumns = {}); +TSerializedTableRange MakeTestRange( + std::pair<ui64, ui64> range, bool inclusiveFrom, bool inclusiveTo, const std::vector<NArrow::NTest::TTestColumn>& columns); - -} +} // namespace NKikimr::NTxUT namespace NKikimr::NColumnShard { - class TTableUpdatesBuilder { - std::vector<std::unique_ptr<arrow::ArrayBuilder>> Builders; - std::shared_ptr<arrow::Schema> Schema; - ui32 RowsCount = 0; +class TTableUpdatesBuilder { + std::vector<std::unique_ptr<arrow::ArrayBuilder>> Builders; + std::shared_ptr<arrow::Schema> Schema; + ui32 RowsCount = 0; + +public: + class TRowBuilder { + TTableUpdatesBuilder& Owner; + YDB_READONLY(ui32, Index, 0); + public: - class TRowBuilder { - TTableUpdatesBuilder& Owner; - YDB_READONLY(ui32, Index, 0); - public: - TRowBuilder(ui32 index, TTableUpdatesBuilder& owner) - : Owner(owner) - , Index(index) - {} + TRowBuilder(ui32 index, TTableUpdatesBuilder& owner) + : Owner(owner) + , Index(index) { + } - TRowBuilder Add(const char* data) { - return Add<std::string>(data); - } + TRowBuilder Add(const char* data) { + return Add<std::string>(data); + } - template <class TData> - TRowBuilder Add(const TData& data) { - Y_ABORT_UNLESS(Index < Owner.Builders.size()); - auto& builder = Owner.Builders[Index]; - auto type = builder->type(); + template <class TData> + TRowBuilder Add(const TData& data) { + Y_ABORT_UNLESS(Index < Owner.Builders.size()); + auto& builder = Owner.Builders[Index]; + auto type = builder->type(); - Y_ABORT_UNLESS(NArrow::SwitchType(type->id(), [&](const auto& t) { - using TWrap = std::decay_t<decltype(t)>; - using T = typename TWrap::T; - using TBuilder = typename arrow::TypeTraits<typename TWrap::T>::BuilderType; + Y_ABORT_UNLESS(NArrow::SwitchType(type->id(), [&](const auto& t) { + using TWrap = std::decay_t<decltype(t)>; + using T = typename TWrap::T; + using TBuilder = typename arrow::TypeTraits<typename TWrap::T>::BuilderType; - AFL_NOTICE(NKikimrServices::TX_COLUMNSHARD)("T", typeid(T).name()); + AFL_NOTICE(NKikimrServices::TX_COLUMNSHARD)("T", typeid(T).name()); - auto& typedBuilder = static_cast<TBuilder&>(*builder); - if constexpr (std::is_arithmetic<TData>::value) { - if constexpr (arrow::has_c_type<T>::value) { - using CType = typename T::c_type; - Y_ABORT_UNLESS(typedBuilder.Append((CType)data).ok()); - return true; - } + auto& typedBuilder = static_cast<TBuilder&>(*builder); + if constexpr (std::is_arithmetic<TData>::value) { + if constexpr (arrow::has_c_type<T>::value) { + using CType = typename T::c_type; + Y_ABORT_UNLESS(typedBuilder.Append((CType)data).ok()); + return true; } - if constexpr (std::is_same<TData, std::string>::value) { - if constexpr (arrow::has_string_view<T>::value && arrow::is_parameter_free_type<T>::value) { - Y_ABORT_UNLESS(typedBuilder.Append(data.data(), data.size()).ok()); - return true; - } + } + if constexpr (std::is_same<TData, std::string>::value) { + if constexpr (arrow::has_string_view<T>::value && arrow::is_parameter_free_type<T>::value) { + Y_ABORT_UNLESS(typedBuilder.Append(data.data(), data.size()).ok()); + return true; } + } - if constexpr (std::is_same<TData, NYdb::TDecimalValue>::value) { - if constexpr (arrow::is_decimal128_type<T>::value) { - Y_ABORT_UNLESS(typedBuilder.Append(arrow::Decimal128(data.Hi_, data.Low_)).ok()); - return true; - } + if constexpr (std::is_same<TData, NYdb::TDecimalValue>::value) { + if constexpr (arrow::is_decimal128_type<T>::value) { + Y_ABORT_UNLESS(typedBuilder.Append(arrow::Decimal128(data.Hi_, data.Low_)).ok()); + return true; } - Y_ABORT("Unknown type combination"); - return false; - })); - return TRowBuilder(Index + 1, Owner); - } - - TRowBuilder AddNull() { - Y_ABORT_UNLESS(Index < Owner.Builders.size()); - auto res = Owner.Builders[Index]->AppendNull(); - return TRowBuilder(Index + 1, Owner); - } - }; - - TTableUpdatesBuilder(std::shared_ptr<arrow::Schema> schema) - : Schema(schema) - { - Builders = NArrow::MakeBuilders(schema); - Y_ABORT_UNLESS(Builders.size() == schema->fields().size()); - } - - TTableUpdatesBuilder(arrow::Result<std::shared_ptr<arrow::Schema>> schema) { - UNIT_ASSERT_C(schema.ok(), schema.status().ToString()); - Schema = schema.ValueUnsafe(); - Builders = NArrow::MakeBuilders(Schema); - Y_ABORT_UNLESS(Builders.size() == Schema->fields().size()); - } - - TRowBuilder AddRow() { - ++RowsCount; - return TRowBuilder(0, *this); + } + Y_ABORT("Unknown type combination"); + return false; + })); + return TRowBuilder(Index + 1, Owner); } - std::shared_ptr<arrow::RecordBatch> BuildArrow() { - TVector<std::shared_ptr<arrow::Array>> columns; - columns.reserve(Builders.size()); - for (auto&& builder : Builders) { - auto arrayDataRes = builder->Finish(); - Y_ABORT_UNLESS(arrayDataRes.ok()); - columns.push_back(*arrayDataRes); - } - return arrow::RecordBatch::Make(Schema, RowsCount, columns); + TRowBuilder AddNull() { + Y_ABORT_UNLESS(Index < Owner.Builders.size()); + auto res = Owner.Builders[Index]->AppendNull(); + return TRowBuilder(Index + 1, Owner); } }; - NOlap::TIndexInfo BuildTableInfo(const std::vector<NArrow::NTest::TTestColumn>& ydbSchema, - const std::vector<NArrow::NTest::TTestColumn>& key); + TTableUpdatesBuilder(std::shared_ptr<arrow::Schema> schema) + : Schema(schema) { + Builders = NArrow::MakeBuilders(schema); + Y_ABORT_UNLESS(Builders.size() == schema->fields().size()); + } + TTableUpdatesBuilder(arrow::Result<std::shared_ptr<arrow::Schema>> schema) { + UNIT_ASSERT_C(schema.ok(), schema.status().ToString()); + Schema = schema.ValueUnsafe(); + Builders = NArrow::MakeBuilders(Schema); + Y_ABORT_UNLESS(Builders.size() == Schema->fields().size()); + } - struct TestTableDescription { - std::vector<NArrow::NTest::TTestColumn> Schema = NTxUT::TTestSchema::YdbSchema(); - std::vector<NArrow::NTest::TTestColumn> Pk = NTxUT::TTestSchema::YdbPkSchema(); - bool InStore = true; + TRowBuilder AddRow() { + ++RowsCount; + return TRowBuilder(0, *this); + } - std::vector<ui32> GetColumnIds(const std::vector<TString>& names) const { - return NTxUT::TTestSchema::GetColumnIds(Schema, names); + std::shared_ptr<arrow::RecordBatch> BuildArrow() { + TVector<std::shared_ptr<arrow::Array>> columns; + columns.reserve(Builders.size()); + for (auto&& builder : Builders) { + auto arrayDataRes = builder->Finish(); + Y_ABORT_UNLESS(arrayDataRes.ok()); + columns.push_back(*arrayDataRes); } - }; + return arrow::RecordBatch::Make(Schema, RowsCount, columns); + } +}; - void SetupSchema(TTestBasicRuntime& runtime, TActorId& sender, ui64 pathId, - const TestTableDescription& table = {}, TString codec = "none"); - void SetupSchema(TTestBasicRuntime& runtime, TActorId& sender, const TString& txBody, const NOlap::TSnapshot& snapshot, bool succeed = true); +NOlap::TIndexInfo BuildTableInfo(const std::vector<NArrow::NTest::TTestColumn>& ydbSchema, const std::vector<NArrow::NTest::TTestColumn>& key); - void PrepareTablet(TTestBasicRuntime& runtime, const ui64 tableId, const std::vector<NArrow::NTest::TTestColumn>& schema, const ui32 keySize = 1); - void PrepareTablet(TTestBasicRuntime& runtime, const TString& schemaTxBody, bool succeed); +struct TestTableDescription { + std::vector<NArrow::NTest::TTestColumn> Schema = NTxUT::TTestSchema::YdbSchema(); + std::vector<NArrow::NTest::TTestColumn> Pk = NTxUT::TTestSchema::YdbPkSchema(); + bool InStore = true; - std::shared_ptr<arrow::RecordBatch> ReadAllAsBatch(TTestBasicRuntime& runtime, const ui64 tableId, const NOlap::TSnapshot& snapshot, const std::vector<NArrow::NTest::TTestColumn>& schema); -} + std::vector<ui32> GetColumnIds(const std::vector<TString>& names) const { + return NTxUT::TTestSchema::GetColumnIds(Schema, names); + } +}; + +void SetupSchema(TTestBasicRuntime& runtime, TActorId& sender, ui64 pathId, const TestTableDescription& table = {}, TString codec = "none"); +void SetupSchema(TTestBasicRuntime& runtime, TActorId& sender, const TString& txBody, const NOlap::TSnapshot& snapshot, bool succeed = true); + +void PrepareTablet( + TTestBasicRuntime& runtime, const ui64 tableId, const std::vector<NArrow::NTest::TTestColumn>& schema, const ui32 keySize = 1); +void PrepareTablet(TTestBasicRuntime& runtime, const TString& schemaTxBody, bool succeed); + +std::shared_ptr<arrow::RecordBatch> ReadAllAsBatch( + TTestBasicRuntime& runtime, const ui64 tableId, const NOlap::TSnapshot& snapshot, const std::vector<NArrow::NTest::TTestColumn>& schema); +} // namespace NKikimr::NColumnShard diff --git a/ydb/core/tx/datashard/datashard__conditional_erase_rows.cpp b/ydb/core/tx/datashard/datashard__conditional_erase_rows.cpp index ca82cd1e338..d68bfe97048 100644 --- a/ydb/core/tx/datashard/datashard__conditional_erase_rows.cpp +++ b/ydb/core/tx/datashard/datashard__conditional_erase_rows.cpp @@ -12,6 +12,7 @@ #include <util/generic/vector.h> #include <util/stream/output.h> #include <util/string/builder.h> +#include <yql/essentials/parser/pg_wrapper/postgresql/src/backend/catalog/pg_type_d.h> namespace NKikimr { namespace NDataShard { diff --git a/ydb/core/tx/schemeshard/olap/operations/alter/common/update.h b/ydb/core/tx/schemeshard/olap/operations/alter/common/update.h index fd10245bc28..439fe2fc78a 100644 --- a/ydb/core/tx/schemeshard/olap/operations/alter/common/update.h +++ b/ydb/core/tx/schemeshard/olap/operations/alter/common/update.h @@ -2,7 +2,7 @@ #include <ydb/core/tx/schemeshard/olap/operations/alter/abstract/update.h> #include <ydb/core/tx/schemeshard/olap/operations/alter/abstract/context.h> #include <ydb/core/tx/schemeshard/olap/table/table.h> -#include <ydb/library/formats/arrow/accessor/common/const.h> +#include <ydb/core/formats/arrow/accessor/common/const.h> namespace NKikimr::NSchemeShard::NOlap::NAlter { @@ -68,7 +68,7 @@ protected: bool CheckTargetSchema(const TOlapSchema& targetSchema) { if (!AppData()->FeatureFlags.GetEnableSparsedColumns()) { - for (auto& [_, column]: targetSchema.GetColumns().GetColumns()) { + for (auto& [_, column] : targetSchema.GetColumns().GetColumns()) { if (column.GetDefaultValue().GetValue() || (column.GetAccessorConstructor().GetClassName() == NKikimr::NArrow::NAccessor::TGlobalConst::SparsedDataAccessorName)) { return false; } diff --git a/ydb/core/tx/schemeshard/olap/operations/alter/standalone/update.cpp b/ydb/core/tx/schemeshard/olap/operations/alter/standalone/update.cpp index 1b11cd53ff9..e2307d7afe7 100644 --- a/ydb/core/tx/schemeshard/olap/operations/alter/standalone/update.cpp +++ b/ydb/core/tx/schemeshard/olap/operations/alter/standalone/update.cpp @@ -1,7 +1,6 @@ #include "update.h" #include <ydb/core/tx/schemeshard/olap/operations/alter/abstract/converter.h> #include <ydb/core/tx/schemeshard/olap/common/common.h> -#include <ydb/library/formats/arrow/accessor/common/const.h> namespace NKikimr::NSchemeShard::NOlap::NAlter { diff --git a/ydb/core/tx/schemeshard/olap/operations/alter_store.cpp b/ydb/core/tx/schemeshard/olap/operations/alter_store.cpp index bc29d6c5613..668495bdc08 100644 --- a/ydb/core/tx/schemeshard/olap/operations/alter_store.cpp +++ b/ydb/core/tx/schemeshard/olap/operations/alter_store.cpp @@ -1,7 +1,7 @@ #include <ydb/core/tx/schemeshard/schemeshard__operation_part.h> #include <ydb/core/tx/schemeshard/schemeshard__operation_common.h> #include <ydb/core/tx/schemeshard/schemeshard_impl.h> -#include <ydb/library/formats/arrow/accessor/common/const.h> +#include <ydb/core/formats/arrow/accessor/common/const.h> #include "checks.h" @@ -11,9 +11,8 @@ using namespace NKikimr; using namespace NSchemeShard; TOlapStoreInfo::TPtr ParseParams(const TOlapStoreInfo::TPtr& storeInfo, - const NKikimrSchemeOp::TAlterColumnStore& alter, - IErrorCollector& errors) -{ + const NKikimrSchemeOp::TAlterColumnStore& alter, + IErrorCollector& errors) { if (!alter.GetRemoveSchemaPresets().empty()) { errors.AddError(NKikimrScheme::StatusInvalidParameter, "Removing schema presets is not supported yet"); return nullptr; @@ -73,27 +72,26 @@ private: TString DebugHint() const override { return TStringBuilder() - << "TAlterOlapStore TConfigureParts" - << " operationId# " << OperationId; + << "TAlterOlapStore TConfigureParts" + << " operationId# " << OperationId; } public: TConfigureParts(TOperationId id) - : OperationId(id) - { - IgnoreMessages(DebugHint(), {TEvHive::TEvCreateTabletReply::EventType}); + : OperationId(id) { + IgnoreMessages(DebugHint(), { TEvHive::TEvCreateTabletReply::EventType }); } bool HandleReply(TEvColumnShard::TEvProposeTransactionResult::TPtr& ev, TOperationContext& context) override { - return NTableState::CollectProposeTransactionResults(OperationId, ev, context); + return NTableState::CollectProposeTransactionResults(OperationId, ev, context); } bool ProgressState(TOperationContext& context) override { TTabletId ssId = context.SS->SelfTabletId(); LOG_INFO_S(context.Ctx, NKikimrServices::FLAT_TX_SCHEMESHARD, - DebugHint() << " ProgressState" - << " at tabletId# " << ssId); + DebugHint() << " ProgressState" + << " at tabletId# " << ssId); TTxState* txState = context.SS->FindTxSafe(OperationId, TTxState::TxAlterOlapStore); TOlapStoreInfo::TPtr storeInfo = context.SS->OlapStores[txState->TargetPathId]; @@ -161,9 +159,9 @@ public: } LOG_DEBUG_S(context.Ctx, NKikimrServices::FLAT_TX_SCHEMESHARD, - DebugHint() << " ProgressState" - << " Propose modify scheme on shard" - << " tabletId: " << tabletId); + DebugHint() << " ProgressState" + << " Propose modify scheme on shard" + << " tabletId: " << tabletId); } txState->UpdateShardsInProgress(); @@ -177,17 +175,16 @@ private: TString DebugHint() const override { return TStringBuilder() - << "TAlterOlapStore TPropose" - << " operationId# " << OperationId; + << "TAlterOlapStore TPropose" + << " operationId# " << OperationId; } public: TPropose(TOperationId id) - : OperationId(id) - { + : OperationId(id) { IgnoreMessages(DebugHint(), - {TEvHive::TEvCreateTabletReply::EventType, - TEvColumnShard::TEvProposeTransactionResult::EventType}); + { TEvHive::TEvCreateTabletReply::EventType, + TEvColumnShard::TEvProposeTransactionResult::EventType }); } bool HandleReply(TEvPrivate::TEvOperationPlan::TPtr& ev, TOperationContext& context) override { @@ -195,9 +192,9 @@ public: TTabletId ssId = context.SS->SelfTabletId(); LOG_INFO_S(context.Ctx, NKikimrServices::FLAT_TX_SCHEMESHARD, - DebugHint() << " HandleReply TEvOperationPlan" - << " at tablet: " << ssId - << ", stepId: " << step); + DebugHint() << " HandleReply TEvOperationPlan" + << " at tablet: " << ssId + << ", stepId: " << step); TTxState* txState = context.SS->FindTx(OperationId); Y_ABORT_UNLESS(txState->TxType == TTxState::TxAlterOlapStore); @@ -242,8 +239,8 @@ public: TTabletId ssId = context.SS->SelfTabletId(); LOG_INFO_S(context.Ctx, NKikimrServices::FLAT_TX_SCHEMESHARD, - DebugHint() << " ProgressState" - << " at tablet: " << ssId); + DebugHint() << " ProgressState" + << " at tablet: " << ssId); TTxState* txState = context.SS->FindTx(OperationId); Y_ABORT_UNLESS(txState); @@ -272,18 +269,17 @@ private: TString DebugHint() const override { return TStringBuilder() - << "TAlterOlapStore TProposedWaitParts" - << " operationId# " << OperationId; + << "TAlterOlapStore TProposedWaitParts" + << " operationId# " << OperationId; } public: TProposedWaitParts(TOperationId id) - : OperationId(id) - { + : OperationId(id) { IgnoreMessages(DebugHint(), - {TEvHive::TEvCreateTabletReply::EventType, + { TEvHive::TEvCreateTabletReply::EventType, TEvColumnShard::TEvProposeTransactionResult::EventType, - TEvPrivate::TEvOperationPlan::EventType}); + TEvPrivate::TEvOperationPlan::EventType }); } bool HandleReply(TEvColumnShard::TEvNotifyTxCompletionResult::TPtr& ev, TOperationContext& context) override { @@ -303,8 +299,8 @@ public: TTabletId ssId = context.SS->SelfTabletId(); LOG_INFO_S(context.Ctx, NKikimrServices::FLAT_TX_SCHEMESHARD, - DebugHint() << " ProgressState" - << " at tablet: " << ssId); + DebugHint() << " ProgressState" + << " at tablet: " << ssId); TTxState* txState = context.SS->FindTx(OperationId); Y_ABORT_UNLESS(txState); @@ -329,9 +325,9 @@ public: } LOG_DEBUG_S(context.Ctx, NKikimrServices::FLAT_TX_SCHEMESHARD, - DebugHint() << " ProgressState" - << " wait for NotifyTxCompletionResult" - << " tabletId: " << tabletId); + DebugHint() << " ProgressState" + << " wait for NotifyTxCompletionResult" + << " tabletId: " << tabletId); } MessagesSent = true; @@ -411,29 +407,29 @@ class TAlterOlapStore: public TSubOperation { TTxState::ETxState NextState(TTxState::ETxState state) const override { switch (state) { - case TTxState::ConfigureParts: - return TTxState::Propose; - case TTxState::Propose: - return TTxState::ProposedWaitParts; - case TTxState::ProposedWaitParts: - return TTxState::Done; - default: - return TTxState::Invalid; + case TTxState::ConfigureParts: + return TTxState::Propose; + case TTxState::Propose: + return TTxState::ProposedWaitParts; + case TTxState::ProposedWaitParts: + return TTxState::Done; + default: + return TTxState::Invalid; } } TSubOperationState::TPtr SelectStateFunc(TTxState::ETxState state) override { switch (state) { - case TTxState::ConfigureParts: - return MakeHolder<TConfigureParts>(OperationId); - case TTxState::Propose: - return MakeHolder<TPropose>(OperationId); - case TTxState::ProposedWaitParts: - return MakeHolder<TProposedWaitParts>(OperationId); - case TTxState::Done: - return MakeHolder<TDone>(OperationId); - default: - return nullptr; + case TTxState::ConfigureParts: + return MakeHolder<TConfigureParts>(OperationId); + case TTxState::Propose: + return MakeHolder<TPropose>(OperationId); + case TTxState::ProposedWaitParts: + return MakeHolder<TProposedWaitParts>(OperationId); + case TTxState::Done: + return MakeHolder<TDone>(OperationId); + default: + return nullptr; } } @@ -461,10 +457,10 @@ public: const TString& name = alter.GetName(); LOG_NOTICE_S(context.Ctx, NKikimrServices::FLAT_TX_SCHEMESHARD, - "TAlterOlapStore Propose" - << ", path: " << parentPathStr << "/" << name - << ", opId: " << OperationId - << ", at schemeshard: " << ssId); + "TAlterOlapStore Propose" + << ", path: " << parentPathStr << "/" << name + << ", opId: " << OperationId + << ", at schemeshard: " << ssId); auto result = MakeHolder<TProposeResponse>(NKikimrScheme::StatusAccepted, ui64(OperationId.GetTxId()), ui64(ssId)); @@ -526,7 +522,7 @@ public: return result; } - for (auto&& tPathId: alterData->ColumnTables) { + for (auto&& tPathId : alterData->ColumnTables) { auto table = context.SS->ColumnTables.GetVerifiedPtr(tPathId); if (!table->Description.HasTtlSettings()) { continue; @@ -539,10 +535,10 @@ public: } if (!AppData()->FeatureFlags.GetEnableSparsedColumns()) { - for (auto& [_, preset]: alterData->SchemaPresets) { - for (auto& [_, column]: preset.GetColumns().GetColumns()) { + for (auto& [_, preset] : alterData->SchemaPresets) { + for (auto& [_, column] : preset.GetColumns().GetColumns()) { if (column.GetDefaultValue().GetValue() || (column.GetAccessorConstructor().GetClassName() == NKikimr::NArrow::NAccessor::TGlobalConst::SparsedDataAccessorName)) { - result->SetError(NKikimrScheme::StatusSchemeError,"schema update error: sparsed columns are disabled"); + result->SetError(NKikimrScheme::StatusSchemeError, "schema update error: sparsed columns are disabled"); return result; } } @@ -591,10 +587,10 @@ public: void AbortUnsafe(TTxId forceDropTxId, TOperationContext& context) override { LOG_NOTICE_S(context.Ctx, NKikimrServices::FLAT_TX_SCHEMESHARD, - "TAlterOlapStore AbortUnsafe" - << ", opId: " << OperationId - << ", forceDropId: " << forceDropTxId - << ", at schemeshard: " << context.SS->TabletID()); + "TAlterOlapStore AbortUnsafe" + << ", opId: " << OperationId + << ", forceDropId: " << forceDropTxId + << ", at schemeshard: " << context.SS->TabletID()); context.OnComplete.DoneOperation(OperationId); } diff --git a/ydb/library/formats/arrow/accessor/abstract/ya.make b/ydb/library/formats/arrow/accessor/abstract/ya.make deleted file mode 100644 index 2e3d8103a49..00000000000 --- a/ydb/library/formats/arrow/accessor/abstract/ya.make +++ /dev/null @@ -1,17 +0,0 @@ -LIBRARY(library-formats-arrow-accessor-abstract) - -PEERDIR( - ydb/library/formats/arrow/protos - ydb/library/formats/arrow/accessor/common - contrib/libs/apache/arrow - ydb/library/conclusion - ydb/library/actors/core -) - -SRCS( - accessor.cpp -) - -GENERATE_ENUM_SERIALIZATION(accessor.h) - -END() diff --git a/ydb/library/formats/arrow/accessor/ya.make b/ydb/library/formats/arrow/accessor/ya.make deleted file mode 100644 index 4c6a540c4b0..00000000000 --- a/ydb/library/formats/arrow/accessor/ya.make +++ /dev/null @@ -1,8 +0,0 @@ -LIBRARY(library-formats-arrow-accessor) - -PEERDIR( - ydb/library/formats/arrow/accessor/abstract - ydb/library/formats/arrow/accessor/composite -) - -END() diff --git a/ydb/library/formats/arrow/arrow_helpers.cpp b/ydb/library/formats/arrow/arrow_helpers.cpp index c9744af773e..af786dfd3ab 100644 --- a/ydb/library/formats/arrow/arrow_helpers.cpp +++ b/ydb/library/formats/arrow/arrow_helpers.cpp @@ -1,23 +1,25 @@ #include "arrow_helpers.h" -#include "switch_type.h" -#include "common/validation.h" #include "permutations.h" -#include "simple_arrays_cache.h" #include "replace_key.h" +#include "simple_arrays_cache.h" -#include <ydb/library/yverify_stream/yverify_stream.h> +#include "switch/switch_type.h" +#include "validation/validation.h" + +#include <ydb/library/actors/core/log.h> #include <ydb/library/services/services.pb.h> +#include <ydb/library/yverify_stream/yverify_stream.h> -#include <util/system/yassert.h> -#include <util/string/join.h> -#include <contrib/libs/apache/arrow/cpp/src/arrow/io/memory.h> -#include <contrib/libs/apache/arrow/cpp/src/arrow/ipc/reader.h> -#include <contrib/libs/apache/arrow/cpp/src/arrow/compute/api.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/array/array_primitive.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/array/builder_primitive.h> +#include <contrib/libs/apache/arrow/cpp/src/arrow/compute/api.h> +#include <contrib/libs/apache/arrow/cpp/src/arrow/io/memory.h> +#include <contrib/libs/apache/arrow/cpp/src/arrow/ipc/reader.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/type_traits.h> #include <library/cpp/containers/stack_vector/stack_vec.h> -#include <ydb/library/actors/core/log.h> +#include <util/string/join.h> +#include <util/system/yassert.h> + #include <memory> #define Y_VERIFY_OK(status) Y_ABORT_UNLESS(status.ok(), "%s", status.ToString().c_str()) @@ -80,8 +82,8 @@ bool IsTrivial(const arrow::UInt64Array& permutation, const ui64 originalLength) return true; } -std::shared_ptr<arrow::RecordBatch> Reorder(const std::shared_ptr<arrow::RecordBatch>& batch, - const std::shared_ptr<arrow::UInt64Array>& permutation, const bool canRemove) { +std::shared_ptr<arrow::RecordBatch> Reorder( + const std::shared_ptr<arrow::RecordBatch>& batch, const std::shared_ptr<arrow::UInt64Array>& permutation, const bool canRemove) { Y_ABORT_UNLESS(permutation->length() == batch->num_rows() || canRemove); auto res = IsTrivial(*permutation, batch->num_rows()) ? batch : arrow::compute::Take(batch, permutation); @@ -89,14 +91,15 @@ std::shared_ptr<arrow::RecordBatch> Reorder(const std::shared_ptr<arrow::RecordB return (*res).record_batch(); } -THashMap<ui64, std::shared_ptr<arrow::RecordBatch>> ShardingSplit(const std::shared_ptr<arrow::RecordBatch>& batch, const THashMap<ui64, std::vector<ui32>>& shardRows) { +THashMap<ui64, std::shared_ptr<arrow::RecordBatch>> ShardingSplit( + const std::shared_ptr<arrow::RecordBatch>& batch, const THashMap<ui64, std::vector<ui32>>& shardRows) { AFL_VERIFY(batch); std::shared_ptr<arrow::UInt64Array> permutation; { arrow::UInt64Builder builder; Y_VERIFY_OK(builder.Reserve(batch->num_rows())); - for (auto&& [shardId, rowIdxs]: shardRows) { + for (auto&& [shardId, rowIdxs] : shardRows) { for (auto& row : rowIdxs) { Y_VERIFY_OK(builder.Append(row)); } @@ -121,7 +124,8 @@ THashMap<ui64, std::shared_ptr<arrow::RecordBatch>> ShardingSplit(const std::sha return out; } -std::vector<std::shared_ptr<arrow::RecordBatch>> ShardingSplit(const std::shared_ptr<arrow::RecordBatch>& batch, const std::vector<std::vector<ui32>>& shardRows, const ui32 numShards) { +std::vector<std::shared_ptr<arrow::RecordBatch>> ShardingSplit( + const std::shared_ptr<arrow::RecordBatch>& batch, const std::vector<std::vector<ui32>>& shardRows, const ui32 numShards) { AFL_VERIFY(batch); std::shared_ptr<arrow::UInt64Array> permutation; { @@ -153,8 +157,8 @@ std::vector<std::shared_ptr<arrow::RecordBatch>> ShardingSplit(const std::shared return out; } -std::vector<std::shared_ptr<arrow::RecordBatch>> ShardingSplit(const std::shared_ptr<arrow::RecordBatch>& batch, - const std::vector<ui32>& sharding, ui32 numShards) { +std::vector<std::shared_ptr<arrow::RecordBatch>> ShardingSplit( + const std::shared_ptr<arrow::RecordBatch>& batch, const std::vector<ui32>& sharding, ui32 numShards) { AFL_VERIFY(batch); Y_ABORT_UNLESS((size_t)batch->num_rows() == sharding.size()); @@ -176,8 +180,8 @@ bool HasAllColumns(const std::shared_ptr<arrow::RecordBatch>& batch, const std:: return true; } -std::vector<std::unique_ptr<arrow::ArrayBuilder>> MakeBuilders(const std::shared_ptr<arrow::Schema>& schema, - size_t reserve, const std::map<std::string, ui64>& sizeByColumn) { +std::vector<std::unique_ptr<arrow::ArrayBuilder>> MakeBuilders( + const std::shared_ptr<arrow::Schema>& schema, size_t reserve, const std::map<std::string, ui64>& sizeByColumn) { std::vector<std::unique_ptr<arrow::ArrayBuilder>> builders; builders.reserve(schema->num_fields()); @@ -196,7 +200,6 @@ std::vector<std::unique_ptr<arrow::ArrayBuilder>> MakeBuilders(const std::shared } builders.emplace_back(std::move(builder)); - } return builders; } @@ -257,7 +260,7 @@ std::shared_ptr<arrow::StringArray> MakeStringArray(const TString& value, const std::pair<int, int> FindMinMaxPosition(const std::shared_ptr<arrow::Array>& array) { if (array->length() == 0) { - return {-1, -1}; + return { -1, -1 }; } int minPos = 0; @@ -279,7 +282,7 @@ std::pair<int, int> FindMinMaxPosition(const std::shared_ptr<arrow::Array>& arra } return true; }); - return {minPos, maxPos}; + return { minPos, maxPos }; } std::shared_ptr<arrow::Scalar> MinScalar(const std::shared_ptr<arrow::DataType>& type) { @@ -289,10 +292,8 @@ std::shared_ptr<arrow::Scalar> MinScalar(const std::shared_ptr<arrow::DataType>& using T = typename TWrap::T; using TScalar = typename arrow::TypeTraits<T>::ScalarType; - if constexpr (std::is_same_v<T, arrow::StringType> || - std::is_same_v<T, arrow::BinaryType> || - std::is_same_v<T, arrow::LargeStringType> || - std::is_same_v<T, arrow::LargeBinaryType>) { + if constexpr (std::is_same_v<T, arrow::StringType> || std::is_same_v<T, arrow::BinaryType> || + std::is_same_v<T, arrow::LargeStringType> || std::is_same_v<T, arrow::LargeBinaryType>) { out = std::make_shared<TScalar>(arrow::Buffer::FromString(""), type); } else if constexpr (std::is_same_v<T, arrow::FixedSizeBinaryType>) { std::string s(static_cast<arrow::FixedSizeBinaryType&>(*type).byte_width(), '\0'); @@ -328,7 +329,7 @@ public: static constexpr bool Value = false; }; -} +} // namespace std::shared_ptr<arrow::Scalar> DefaultScalar(const std::shared_ptr<arrow::DataType>& type) { std::shared_ptr<arrow::Scalar> out; @@ -337,10 +338,8 @@ std::shared_ptr<arrow::Scalar> DefaultScalar(const std::shared_ptr<arrow::DataTy using T = typename TWrap::T; using TScalar = typename arrow::TypeTraits<T>::ScalarType; - if constexpr (std::is_same_v<T, arrow::StringType> || - std::is_same_v<T, arrow::BinaryType> || - std::is_same_v<T, arrow::LargeStringType> || - std::is_same_v<T, arrow::LargeBinaryType>) { + if constexpr (std::is_same_v<T, arrow::StringType> || std::is_same_v<T, arrow::BinaryType> || + std::is_same_v<T, arrow::LargeStringType> || std::is_same_v<T, arrow::LargeBinaryType>) { out = std::make_shared<TScalar>(arrow::Buffer::FromString(""), type); } else if constexpr (std::is_same_v<T, arrow::FixedSizeBinaryType>) { std::string s(static_cast<arrow::FixedSizeBinaryType&>(*type).byte_width(), '\0'); @@ -399,11 +398,10 @@ bool ScalarLess(const arrow::Scalar& x, const arrow::Scalar& y) { return ScalarCompare(x, y) < 0; } -bool ColumnEqualsScalar( - const std::shared_ptr<arrow::Array>& c, const ui32 position, const std::shared_ptr<arrow::Scalar>& s) { +bool ColumnEqualsScalar(const std::shared_ptr<arrow::Array>& c, const ui32 position, const std::shared_ptr<arrow::Scalar>& s) { AFL_VERIFY(c); if (!s) { - return c->IsNull(position) ; + return c->IsNull(position); } AFL_VERIFY(c->type()->Equals(s->type))("s", s->type->ToString())("c", c->type()->ToString()); @@ -465,7 +463,7 @@ int ScalarCompare(const arrow::Scalar& x, const arrow::Scalar& y) { return 0; } } - Y_ABORT_UNLESS(false); // TODO: non primitive types + Y_ABORT_UNLESS(false); // TODO: non primitive types return 0; }); } @@ -499,7 +497,6 @@ std::shared_ptr<arrow::Array> BoolVecToArray(const std::vector<bool>& vec) { return out; } - bool ArrayScalarsEqual(const std::shared_ptr<arrow::Array>& lhs, const std::shared_ptr<arrow::Array>& rhs) { bool res = lhs->length() == rhs->length(); for (int64_t i = 0; i < lhs->length() && res; ++i) { @@ -510,9 +507,7 @@ bool ArrayScalarsEqual(const std::shared_ptr<arrow::Array>& lhs, const std::shar bool ReserveData(arrow::ArrayBuilder& builder, const size_t size) { arrow::Status result = arrow::Status::OK(); - if (builder.type()->id() == arrow::Type::BINARY || - builder.type()->id() == arrow::Type::STRING) - { + if (builder.type()->id() == arrow::Type::BINARY || builder.type()->id() == arrow::Type::STRING) { static_assert(std::is_convertible_v<arrow::StringBuilder&, arrow::BaseBinaryBuilder<arrow::BinaryType>&>, "Expected StringBuilder to be BaseBinaryBuilder<BinaryType>"); auto& bBuilder = static_cast<arrow::BaseBinaryBuilder<arrow::BinaryType>&>(builder); @@ -549,7 +544,8 @@ bool MergeBatchColumnsImpl(const std::vector<std::shared_ptr<TData>>& batches, s fields.emplace_back(f); } if (i->num_rows() != batches.front()->num_rows()) { - AFL_ERROR(NKikimrServices::ARROW_HELPER)("event", "inconsistency record sizes")("i", i->num_rows())("front", batches.front()->num_rows()); + AFL_ERROR(NKikimrServices::ARROW_HELPER)("event", "inconsistency record sizes")("i", i->num_rows())( + "front", batches.front()->num_rows()); return false; } for (auto&& c : i->columns()) { @@ -578,23 +574,28 @@ bool MergeBatchColumnsImpl(const std::vector<std::shared_ptr<TData>>& batches, s return true; } -bool MergeBatchColumns(const std::vector<std::shared_ptr<arrow::Table>>& batches, std::shared_ptr<arrow::Table>& result, const std::vector<std::string>& columnsOrder, const bool orderFieldsAreNecessary) { - const auto builder = [](const std::shared_ptr<arrow::Schema>& schema, const ui32 recordsCount, std::vector<std::shared_ptr<arrow::ChunkedArray>>&& columns) { +bool MergeBatchColumns(const std::vector<std::shared_ptr<arrow::Table>>& batches, std::shared_ptr<arrow::Table>& result, + const std::vector<std::string>& columnsOrder, const bool orderFieldsAreNecessary) { + const auto builder = [](const std::shared_ptr<arrow::Schema>& schema, const ui32 recordsCount, + std::vector<std::shared_ptr<arrow::ChunkedArray>>&& columns) { return arrow::Table::Make(schema, columns, recordsCount); }; return MergeBatchColumnsImpl<arrow::Table, arrow::ChunkedArray>(batches, result, columnsOrder, orderFieldsAreNecessary, builder); } -bool MergeBatchColumns(const std::vector<std::shared_ptr<arrow::RecordBatch>>& batches, std::shared_ptr<arrow::RecordBatch>& result, const std::vector<std::string>& columnsOrder, const bool orderFieldsAreNecessary) { - const auto builder = [](const std::shared_ptr<arrow::Schema>& schema, const ui32 recordsCount, std::vector<std::shared_ptr<arrow::Array>>&& columns) { +bool MergeBatchColumns(const std::vector<std::shared_ptr<arrow::RecordBatch>>& batches, std::shared_ptr<arrow::RecordBatch>& result, + const std::vector<std::string>& columnsOrder, const bool orderFieldsAreNecessary) { + const auto builder = [](const std::shared_ptr<arrow::Schema>& schema, const ui32 recordsCount, + std::vector<std::shared_ptr<arrow::Array>>&& columns) { return arrow::RecordBatch::Make(schema, recordsCount, columns); }; return MergeBatchColumnsImpl<arrow::RecordBatch, arrow::Array>(batches, result, columnsOrder, orderFieldsAreNecessary, builder); } -std::partial_ordering ColumnsCompare(const std::vector<std::shared_ptr<arrow::Array>>& x, const ui32 xRow, const std::vector<std::shared_ptr<arrow::Array>>& y, const ui32 yRow) { +std::partial_ordering ColumnsCompare( + const std::vector<std::shared_ptr<arrow::Array>>& x, const ui32 xRow, const std::vector<std::shared_ptr<arrow::Array>>& y, const ui32 yRow) { return TRawReplaceKey(&x, xRow).CompareNotNull(TRawReplaceKey(&y, yRow)); } @@ -681,7 +682,7 @@ NJson::TJsonValue DebugJson(std::shared_ptr<arrow::Array> array, const ui32 head } } return true; - }); + }); return resultFull; } @@ -766,7 +767,8 @@ std::vector<std::shared_ptr<arrow::RecordBatch>> SliceToRecordBatches(const std: AFL_VERIFY(it != i->chunks().end()); AFL_VERIFY(positions[idx + 1] - currentPosition <= length)("length", length)("idx+1", positions[idx + 1])("pos", currentPosition); auto chunk = (*it)->Slice(positions[idx] - currentPosition, positions[idx + 1] - positions[idx]); - AFL_VERIFY_DEBUG(chunk->length() == positions[idx + 1] - positions[idx])("length", chunk->length())("expect", positions[idx + 1] - positions[idx]); + AFL_VERIFY_DEBUG(chunk->length() == positions[idx + 1] - positions[idx]) + ("length", chunk->length())("expect", positions[idx + 1] - positions[idx]); if (positions[idx + 1] - currentPosition == length) { ++it; initializeIt(); @@ -784,7 +786,7 @@ std::vector<std::shared_ptr<arrow::RecordBatch>> SliceToRecordBatches(const std: count += result.back()->num_rows(); } AFL_VERIFY(count == t->num_rows())("count", count)("t", t->num_rows())("sd_size", slicedData.size())("columns", t->num_columns())( - "schema", t->schema()->ToString()); + "schema", t->schema()->ToString()); return result; } @@ -792,7 +794,7 @@ std::shared_ptr<arrow::Table> ToTable(const std::shared_ptr<arrow::RecordBatch>& if (!batch) { return nullptr; } - return TStatusValidator::GetValid(arrow::Table::FromRecordBatches(batch->schema(), {batch})); + return TStatusValidator::GetValid(arrow::Table::FromRecordBatches(batch->schema(), { batch })); } bool HasNulls(const std::shared_ptr<arrow::Array>& column) { @@ -819,7 +821,7 @@ std::vector<std::string> ConvertStrings(const std::vector<TString>& input) { std::shared_ptr<arrow::Table> DeepCopy(const std::shared_ptr<arrow::Table>& table, arrow::MemoryPool* pool) { arrow::ArrayVector arrays; - for (const auto& column: table->columns()) { + for (const auto& column : table->columns()) { auto&& array = TStatusValidator::GetValid(arrow::Concatenate(column->chunks(), pool)); arrays.push_back(std::move(array)); } @@ -827,4 +829,4 @@ std::shared_ptr<arrow::Table> DeepCopy(const std::shared_ptr<arrow::Table>& tabl return arrow::Table::Make(table->schema(), arrays); } -} // namespace NKikimr::NArrow +} // namespace NKikimr::NArrow diff --git a/ydb/library/formats/arrow/arrow_helpers.h b/ydb/library/formats/arrow/arrow_helpers.h index de507558028..c73610d2140 100644 --- a/ydb/library/formats/arrow/arrow_helpers.h +++ b/ydb/library/formats/arrow/arrow_helpers.h @@ -1,5 +1,4 @@ #pragma once -#include "switch_type.h" #include <library/cpp/json/writer/json_value.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/api.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/type_traits.h> diff --git a/ydb/library/formats/arrow/common/validation.h b/ydb/library/formats/arrow/common/validation.h deleted file mode 100644 index 171b50041db..00000000000 --- a/ydb/library/formats/arrow/common/validation.h +++ /dev/null @@ -1,3 +0,0 @@ -#pragma once - -#include <ydb/library/formats/arrow/validation/validation.h> diff --git a/ydb/library/formats/arrow/permutations.cpp b/ydb/library/formats/arrow/permutations.cpp index 8eb270c2a42..edd98870bfa 100644 --- a/ydb/library/formats/arrow/permutations.cpp +++ b/ydb/library/formats/arrow/permutations.cpp @@ -1,13 +1,13 @@ -#include "permutations.h" - #include "arrow_helpers.h" +#include "permutations.h" #include "replace_key.h" #include "size_calcer.h" -#include <ydb/library/formats/arrow/common/validation.h> -#include <ydb/library/services/services.pb.h> +#include "switch/switch_type.h" +#include "validation/validation.h" #include <ydb/library/actors/core/log.h> +#include <ydb/library/services/services.pb.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/array/builder_primitive.h> #include <contrib/libs/xxhash/xxhash.h> @@ -120,20 +120,17 @@ ui64 TShardedRecordBatch::GetMemorySize() const { TShardedRecordBatch::TShardedRecordBatch(const std::shared_ptr<arrow::RecordBatch>& batch) { AFL_VERIFY(batch); - RecordBatch = TStatusValidator::GetValid(arrow::Table::FromRecordBatches(batch->schema(), {batch})); + RecordBatch = TStatusValidator::GetValid(arrow::Table::FromRecordBatches(batch->schema(), { batch })); } - TShardedRecordBatch::TShardedRecordBatch(const std::shared_ptr<arrow::Table>& batch) - : RecordBatch(batch) -{ + : RecordBatch(batch) { AFL_VERIFY(RecordBatch); } TShardedRecordBatch::TShardedRecordBatch(const std::shared_ptr<arrow::Table>& batch, std::vector<std::vector<ui32>>&& splittedByShards) : RecordBatch(batch) - , SplittedByShards(std::move(splittedByShards)) -{ + , SplittedByShards(std::move(splittedByShards)) { AFL_VERIFY(RecordBatch); AFL_VERIFY(SplittedByShards.size()); } @@ -154,7 +151,8 @@ std::vector<std::shared_ptr<arrow::Table>> TShardingSplitIndex::Apply(const std: return result; } -NKikimr::NArrow::TShardedRecordBatch TShardingSplitIndex::Apply(const ui32 shardsCount, const std::shared_ptr<arrow::Table>& input, const std::string& hashColumnName) { +NKikimr::NArrow::TShardedRecordBatch TShardingSplitIndex::Apply( + const ui32 shardsCount, const std::shared_ptr<arrow::Table>& input, const std::string& hashColumnName) { AFL_VERIFY(input); if (shardsCount == 1) { return TShardedRecordBatch(input); @@ -179,9 +177,9 @@ NKikimr::NArrow::TShardedRecordBatch TShardingSplitIndex::Apply(const ui32 shard return TShardedRecordBatch(resultBatch, splitter->DetachRemapping()); } -TShardedRecordBatch TShardingSplitIndex::Apply(const ui32 shardsCount, const std::shared_ptr<arrow::RecordBatch>& input, const std::string& hashColumnName) { - return Apply(shardsCount, TStatusValidator::GetValid(arrow::Table::FromRecordBatches(input->schema(), {input})) - , hashColumnName); +TShardedRecordBatch TShardingSplitIndex::Apply( + const ui32 shardsCount, const std::shared_ptr<arrow::RecordBatch>& input, const std::string& hashColumnName) { + return Apply(shardsCount, TStatusValidator::GetValid(arrow::Table::FromRecordBatches(input->schema(), { input })), hashColumnName); } std::shared_ptr<arrow::UInt64Array> TShardingSplitIndex::BuildPermutation() const { @@ -211,4 +209,4 @@ std::shared_ptr<arrow::Table> ReverseRecords(const std::shared_ptr<arrow::Table> return NArrow::TStatusValidator::GetValid(arrow::compute::Take(batch, permutation)).table(); } -} +} // namespace NKikimr::NArrow diff --git a/ydb/library/formats/arrow/replace_key.h b/ydb/library/formats/arrow/replace_key.h index 159915f944a..a0e3a26b27e 100644 --- a/ydb/library/formats/arrow/replace_key.h +++ b/ydb/library/formats/arrow/replace_key.h @@ -2,7 +2,6 @@ #include "arrow_helpers.h" #include "permutations.h" -#include "common/validation.h" #include "switch/compare.h" #include <ydb/library/actors/core/log.h> diff --git a/ydb/library/formats/arrow/simple_arrays_cache.cpp b/ydb/library/formats/arrow/simple_arrays_cache.cpp index e963f50607c..35ef20a8d09 100644 --- a/ydb/library/formats/arrow/simple_arrays_cache.cpp +++ b/ydb/library/formats/arrow/simple_arrays_cache.cpp @@ -1,5 +1,6 @@ #include "simple_arrays_cache.h" -#include "common/validation.h" + +#include "validation/validation.h" #include <arrow/array/util.h> @@ -9,12 +10,13 @@ std::shared_ptr<arrow::Array> TThreadSimpleArraysCache::GetNullImpl(const std::s AFL_VERIFY(type); const TString key = type->ToString(); const auto initializer = [type](const ui32 recordsCount) { - return NArrow::TStatusValidator::GetValid(arrow::MakeArrayOfNull(type, recordsCount)); + return TStatusValidator::GetValid(arrow::MakeArrayOfNull(type, recordsCount)); }; return InitializePosition(type->ToString(), recordsCount, initializer); } -std::shared_ptr<arrow::Array> TThreadSimpleArraysCache::GetConstImpl(const std::shared_ptr<arrow::DataType>& type, const std::shared_ptr<arrow::Scalar>& scalar, const ui32 recordsCount) { +std::shared_ptr<arrow::Array> TThreadSimpleArraysCache::GetConstImpl( + const std::shared_ptr<arrow::DataType>& type, const std::shared_ptr<arrow::Scalar>& scalar, const ui32 recordsCount) { AFL_VERIFY(type); AFL_VERIFY(scalar); AFL_VERIFY(scalar->type->id() == type->id())("scalar", scalar->type->ToString())("field", type->ToString()); @@ -32,11 +34,13 @@ std::shared_ptr<arrow::Array> TThreadSimpleArraysCache::GetNull(const std::share return SimpleArraysCache.GetNullImpl(type, recordsCount); } -std::shared_ptr<arrow::Array> TThreadSimpleArraysCache::GetConst(const std::shared_ptr<arrow::DataType>& type, const std::shared_ptr<arrow::Scalar>& scalar, const ui32 recordsCount) { +std::shared_ptr<arrow::Array> TThreadSimpleArraysCache::GetConst( + const std::shared_ptr<arrow::DataType>& type, const std::shared_ptr<arrow::Scalar>& scalar, const ui32 recordsCount) { return SimpleArraysCache.GetConstImpl(type, scalar, recordsCount); } -std::shared_ptr<arrow::Array> TThreadSimpleArraysCache::Get(const std::shared_ptr<arrow::DataType>& type, const std::shared_ptr<arrow::Scalar>& scalar, const ui32 recordsCount) { +std::shared_ptr<arrow::Array> TThreadSimpleArraysCache::Get( + const std::shared_ptr<arrow::DataType>& type, const std::shared_ptr<arrow::Scalar>& scalar, const ui32 recordsCount) { if (scalar) { return GetConst(type, scalar, recordsCount); } else { @@ -44,5 +48,4 @@ std::shared_ptr<arrow::Array> TThreadSimpleArraysCache::Get(const std::shared_pt } } -} - +} // namespace NKikimr::NArrow diff --git a/ydb/library/formats/arrow/switch/switch_type.h b/ydb/library/formats/arrow/switch/switch_type.h index 54ea8f19636..ac988ae2417 100644 --- a/ydb/library/formats/arrow/switch/switch_type.h +++ b/ydb/library/formats/arrow/switch/switch_type.h @@ -1,9 +1,9 @@ #pragma once -#include <ydb/library/formats/arrow/common/validation.h> -#include <yql/essentials/parser/pg_wrapper/interface/type_desc.h> +#include <ydb/library/formats/arrow/validation/validation.h> #include <contrib/libs/apache/arrow/cpp/src/arrow/api.h> #include <util/system/yassert.h> +#include <yql/essentials/parser/pg_wrapper/interface/type_desc.h> extern "C" { #include <yql/essentials/parser/pg_wrapper/postgresql/src/include/catalog/pg_type_d.h> @@ -12,8 +12,7 @@ extern "C" { namespace NKikimr::NArrow { template <typename TType> -struct TTypeWrapper -{ +struct TTypeWrapper { using T = TType; }; @@ -181,4 +180,4 @@ template <typename T> }); } -} +} // namespace NKikimr::NArrow diff --git a/ydb/library/formats/arrow/switch/ya.make b/ydb/library/formats/arrow/switch/ya.make index 1b6d894de7d..a56125663fc 100644 --- a/ydb/library/formats/arrow/switch/ya.make +++ b/ydb/library/formats/arrow/switch/ya.make @@ -3,6 +3,7 @@ LIBRARY(library-formats-arrow-switch) PEERDIR( contrib/libs/apache/arrow ydb/library/actors/core + ydb/library/formats/arrow/validation ) SRCS( diff --git a/ydb/library/formats/arrow/switch_type.h b/ydb/library/formats/arrow/switch_type.h deleted file mode 100644 index 1acf90b1bb2..00000000000 --- a/ydb/library/formats/arrow/switch_type.h +++ /dev/null @@ -1,2 +0,0 @@ -#pragma once -#include "switch/switch_type.h" diff --git a/ydb/library/formats/arrow/ya.make b/ydb/library/formats/arrow/ya.make index 316ee9f4e5c..dff92da2abf 100644 --- a/ydb/library/formats/arrow/ya.make +++ b/ydb/library/formats/arrow/ya.make @@ -3,7 +3,6 @@ RECURSE_FOR_TESTS( ) RECURSE( - accessor common switch csv @@ -20,7 +19,6 @@ LIBRARY() PEERDIR( contrib/libs/apache/arrow - ydb/library/formats/arrow/accessor ydb/library/formats/arrow/simple_builder ydb/library/formats/arrow/transformer ydb/library/formats/arrow/splitter diff --git a/ydb/services/ydb/ydb_common_ut.h b/ydb/services/ydb/ydb_common_ut.h index 1352f62a458..5a008348551 100644 --- a/ydb/services/ydb/ydb_common_ut.h +++ b/ydb/services/ydb/ydb_common_ut.h @@ -4,6 +4,7 @@ #include <library/cpp/testing/unittest/registar.h> #include <ydb/core/testlib/test_client.h> #include <ydb/core/formats/arrow/arrow_helpers.h> +#include <ydb/library/formats/arrow/switch/switch_type.h> #include <ydb/core/security/certificate_check/cert_auth_utils.h> #include <ydb/services/ydb/ydb_dummy.h> #include <ydb-cpp-sdk/client/value/value.h> |
