diff options
| author | risenberg <[email protected]> | 2026-07-11 10:22:48 +0300 |
|---|---|---|
| committer | GitHub <[email protected]> | 2026-07-11 10:22:48 +0300 |
| commit | ebff6a5cc9608f384a445589efbaa55099729cf5 (patch) | |
| tree | de1c7d575141156c1e23632c79aa244994a3f40b | |
| parent | 13a3a42e8d6f4e7c7926ed27e620425b467f4483 (diff) | |
Store same-type subcolumns as native scalars (#45969)
Adds an opt-in encoding for the SUB_COLUMNS JSON accessor: when every present value of a
separated sub-column is a scalar of one JSON type, store it in the matching native Arrow array
(float64 / boolean / raw-string binary) instead of a per-value BinaryJson blob. Gated behind a
new default-off setting ENABLE_NATIVE_COLUMNS; behavior is unchanged unless it's turned
on.
29 files changed, 936 insertions, 150 deletions
diff --git a/ydb/core/formats/arrow/accessor/plain/accessor.h b/ydb/core/formats/arrow/accessor/plain/accessor.h index ab5d122c359..5a3e7d36a77 100644 --- a/ydb/core/formats/arrow/accessor/plain/accessor.h +++ b/ydb/core/formats/arrow/accessor/plain/accessor.h @@ -82,6 +82,11 @@ public: } void AddRecord(const ui32 recordIndex, const std::string_view value) { + AddValue(recordIndex, arrow::util::string_view(value.data(), value.size())); + } + + template <class TValue> + void AddValue(const ui32 recordIndex, const TValue& value) { if (LastRecordIndex) { AFL_VERIFY(*LastRecordIndex < recordIndex)("last", LastRecordIndex)("index", recordIndex); TStatusValidator::Validate(Builder->AppendNulls(recordIndex - *LastRecordIndex - 1)); @@ -89,7 +94,7 @@ public: TStatusValidator::Validate(Builder->AppendNulls(recordIndex)); } LastRecordIndex = recordIndex; - AFL_VERIFY(NArrow::Append<TArrowDataType>(*Builder, arrow::util::string_view(value.data(), value.size()))); + AFL_VERIFY(NArrow::Append<TArrowDataType>(*Builder, value)); } void AddNull(const ui32 recordIndex) { diff --git a/ydb/core/formats/arrow/accessor/sub_columns/accessor.cpp b/ydb/core/formats/arrow/accessor/sub_columns/accessor.cpp index 2a8be9ad3eb..58894ffd10f 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/accessor.cpp +++ b/ydb/core/formats/arrow/accessor/sub_columns/accessor.cpp @@ -74,7 +74,8 @@ TString TSubColumnsArray::SerializeToString(const TChunkConstructionData& extern ui32 columnIdx = 0; TMonotonic pred = TMonotonic::Now(); for (auto&& i : ColumnsData.GetRecords()->GetColumns()) { - TChunkConstructionData cData(GetRecordsCount(), nullptr, arrow::binary(), externalInfo.GetDefaultSerializer()); + TChunkConstructionData cData( + GetRecordsCount(), nullptr, ColumnsData.GetStats().GetField(columnIdx)->type(), externalInfo.GetDefaultSerializer()); auto* cInfo = proto.AddKeyColumns(); if (ColumnsData.GetStats().GetAccessorType(columnIdx) == IChunkedArray::EType::Dictionary) { // Dictionary columns produce [dictionary blob][positions blob]; the split is not 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 f49bacd4ccc..46d2fbd4fe6 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/columns_storage.cpp +++ b/ydb/core/formats/arrow/accessor/sub_columns/columns_storage.cpp @@ -15,7 +15,8 @@ TColumnsData TColumnsData::Slice(const ui32 offset, const ui32 count) const { 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()); + builder.Add(Stats.GetColumnName(idx), i->GetRecordsCount() - i->GetNullsCount(), i->GetValueRawBytes(), i->GetType(), + Stats.GetValueType(idx)); } else { indexesToRemove.emplace_back(idx); } @@ -40,7 +41,8 @@ TColumnsData TColumnsData::ApplyFilter(const TColumnFilter& filter) const { 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()); + builder.Add(Stats.GetColumnName(idx), i->GetRecordsCount() - i->GetNullsCount(), i->GetValueRawBytes(), i->GetType(), + Stats.GetValueType(idx)); } else { indexesToRemove.emplace_back(idx); } @@ -62,9 +64,8 @@ void TColumnsData::TIterator::InitArrays() { } const ui32 localIndex = FullArrayAddress->GetAddress().GetLocalIndex(CurrentIndex); ChunkAddress = FullArrayAddress->GetArray()->GetChunk(ChunkAddress, localIndex); - AFL_VERIFY(ChunkAddress->GetArray()->type()->id() == arrow::binary()->id()); - CurrentArrayData = static_cast<const arrow::BinaryArray*>(ChunkAddress->GetArray().get()); - // Dictionary columns materialize (decode) to a dense binary array, so they are + CurrentArrayData = ChunkAddress->GetArray().get(); + // Dictionary columns materialize (decode) to a dense array, so they are // read exactly like a plain Array here. if (FullArrayAddress->GetArray()->GetType() == IChunkedArray::EType::Array || FullArrayAddress->GetArray()->GetType() == IChunkedArray::EType::Dictionary) { @@ -75,8 +76,8 @@ void TColumnsData::TIterator::InitArrays() { } else if (FullArrayAddress->GetArray()->GetType() == IChunkedArray::EType::SparsedArray) { AFL_VERIFY(localIndex < CurrentArrayData->length()) ("localIndex", localIndex) - ("CurrentArrayData->length()", CurrentArrayData->length()) - ("CurrentArrayData", CurrentArrayData->ToString()); + ("CurrentArray->length()", CurrentArrayData->length()) + ("CurrentArray", CurrentArrayData->ToString()); if (CurrentArrayData->IsNull(localIndex) && std::static_pointer_cast<TSparsedArray>(FullArrayAddress->GetArray())->GetDefaultValue() == nullptr) { CurrentIndex = ChunkAddress->GetAddress().GetGlobalFinishPosition(); @@ -90,9 +91,4 @@ void TColumnsData::TIterator::InitArrays() { AFL_VERIFY(CurrentIndex <= GlobalChunkedArray->GetRecordsCount())("index", CurrentIndex)("count", GlobalChunkedArray->GetRecordsCount()); } -NArrow::NAccessor::TBinaryJsonValueView TColumnsData::TIterator::GetValue() const { - auto view = CurrentArrayData->GetView(ChunkAddress->GetAddress().GetLocalIndex(CurrentIndex)); - return NArrow::NAccessor::TBinaryJsonValueView(TStringBuf(view.data(), view.size())); -} - } // namespace NKikimr::NArrow::NAccessor::NSubColumns 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 c8821f77d10..0947bad085a 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/columns_storage.h +++ b/ydb/core/formats/arrow/accessor/sub_columns/columns_storage.h @@ -51,7 +51,7 @@ public: private: ui32 KeyIndex; std::shared_ptr<IChunkedArray> GlobalChunkedArray; - const arrow::BinaryArray* CurrentArrayData; + const arrow::Array* CurrentArrayData; std::optional<IChunkedArray::TFullChunkedArrayAddress> FullArrayAddress; std::optional<IChunkedArray::TFullDataAddress> ChunkAddress; ui32 CurrentIndex = 0; @@ -73,12 +73,19 @@ public: return KeyIndex; } - std::string_view GetRawValue() const { - auto view = CurrentArrayData->GetView(ChunkAddress->GetAddress().GetLocalIndex(CurrentIndex)); - return std::string_view(view.data(), view.size()); + // Current value is exposed as (array, local index); the reader interprets it per the + // column's value type. + const arrow::Array& GetArray() const { + return *CurrentArrayData; + } + i64 GetLocalIndex() const { + return ChunkAddress->GetAddress().GetLocalIndex(CurrentIndex); } - NArrow::NAccessor::TBinaryJsonValueView GetValue() const; + NArrow::NAccessor::TBinaryJsonValueView GetValue() const { + auto view = static_cast<const arrow::BinaryArray&>(*CurrentArrayData).GetView(GetLocalIndex()); + return NArrow::NAccessor::TBinaryJsonValueView(TStringBuf(view.data(), view.size())); + } bool HasValue() const { return !CurrentArrayData->IsNull(ChunkAddress->GetAddress().GetLocalIndex(CurrentIndex)); @@ -132,8 +139,9 @@ public: : Stats(dict) , Records(data) { AFL_VERIFY(Records->num_columns() == Stats.GetColumnsCount())("records", Records->num_columns())("stats", Stats.GetColumnsCount()); - for (auto&& i : Records->GetColumns()) { - AFL_VERIFY(i->GetDataType()->id() == arrow::binary()->id()); + for (ui32 i = 0; i < (ui32)Records->num_columns(); ++i) { + AFL_VERIFY(Records->GetColumnVerified(i)->GetDataType()->id() == Stats.GetField(i)->type()->id())( + "column", Records->GetColumnVerified(i)->GetDataType()->ToString())("stats", Stats.GetField(i)->type()->ToString()); } } }; diff --git a/ydb/core/formats/arrow/accessor/sub_columns/direct_builder.cpp b/ydb/core/formats/arrow/accessor/sub_columns/direct_builder.cpp index e8188811bd1..a67b4ef3867 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/direct_builder.cpp +++ b/ydb/core/formats/arrow/accessor/sub_columns/direct_builder.cpp @@ -13,30 +13,66 @@ namespace NKikimr::NArrow::NAccessor::NSubColumns { -void TColumnElements::BuildSparsedAccessor(const ui32 recordsCount) { +namespace { + +std::string_view MakeStoredBytesView(const NBinaryJson::TBinaryJson& rec, const EValueType valueType) { + if (valueType == EValueType::BinaryJson) { + return std::string_view(rec.Data(), rec.Size()); + } + AFL_VERIFY(valueType == EValueType::String)("value_type", (ui32)valueType); + const auto scalar = ExtractStringScalar(rec); + return std::string_view(scalar.data(), scalar.size()); +} + +template <class TArrow, class TExtractor> +std::shared_ptr<IChunkedArray> BuildTypedPlain(const std::deque<NBinaryJson::TBinaryJson>& values, + const std::vector<ui32>& recordIndexes, const ui32 recordsCount, const ui32 reserveData, const TExtractor& extract) { + TTrivialArray::TPlainBuilder<TArrow> builder(recordsCount, reserveData); + for (ui32 i = 0; i < recordIndexes.size(); ++i) { + builder.AddValue(recordIndexes[i], extract(values[i])); + } + return builder.Finish(recordsCount); +} +} // namespace + +void TColumnElements::BuildSparsedAccessor(const ui32 recordsCount, const EValueType valueType) { AFL_VERIFY(!Accessor); + // Columns with non-binary encoding are not supported. + AFL_VERIFY(valueType == EValueType::BinaryJson || valueType == EValueType::String)("value_type", (ui32)valueType); auto recordsBuilder = TSparsedArray::MakeBuilderBinary(RecordIndexes.size(), DataSize); for (ui32 idx = 0; idx < RecordIndexes.size(); ++idx) { - const auto& rec = Values[idx]; - recordsBuilder.AddRecord(RecordIndexes[idx], std::string_view(rec.Data(), rec.Size())); + recordsBuilder.AddRecord(RecordIndexes[idx], MakeStoredBytesView(Values[idx], valueType)); } Accessor = recordsBuilder.Finish(recordsCount); } -void TColumnElements::BuildPlainAccessor(const ui32 recordsCount) { +void TColumnElements::BuildPlainAccessor(const ui32 recordsCount, const EValueType valueType) { AFL_VERIFY(!Accessor); - auto builder = TTrivialArray::MakeBuilderBinary(recordsCount, DataSize); - for (auto it = RecordIndexes.begin(); it != RecordIndexes.end(); ++it) { - const auto& rec = Values[it - RecordIndexes.begin()]; - builder.AddRecord(*it, std::string_view(rec.Data(), rec.Size())); + switch (valueType) { + case EValueType::BinaryJson: + case EValueType::String: + Accessor = BuildTypedPlain<arrow::BinaryType>(Values, RecordIndexes, recordsCount, DataSize, + [valueType](const NBinaryJson::TBinaryJson& rec) { + const auto sv = MakeStoredBytesView(rec, valueType); + return arrow::util::string_view(sv.data(), sv.size()); + }); + break; + case EValueType::Double: + // No need to reserve by size for fixed-size arrays + Accessor = BuildTypedPlain<arrow::DoubleType>(Values, RecordIndexes, recordsCount, 0, + &ExtractDoubleScalar); + break; + case EValueType::Bool: + Accessor = BuildTypedPlain<arrow::BooleanType>(Values, RecordIndexes, recordsCount, 0, + &ExtractBoolScalar); + break; } - Accessor = builder.Finish(recordsCount); } -void TColumnElements::BuildDictionaryAccessor(const ui32 recordsCount) { - BuildPlainAccessor(recordsCount); +void TColumnElements::BuildDictionaryAccessor(const ui32 recordsCount, const EValueType valueType) { + BuildPlainAccessor(recordsCount, valueType); const TChunkConstructionData cData( - recordsCount, nullptr, arrow::binary(), NSerialization::TSerializerContainer::GetDefaultSerializer()); + recordsCount, nullptr, GetArrowTypeForValueType(valueType), NSerialization::TSerializerContainer::GetDefaultSerializer()); Accessor = NDictionary::TConstructor().Construct(Accessor, cData).DetachResult(); } @@ -67,19 +103,20 @@ std::shared_ptr<TSubColumnsArray> TDataBuilder::Finish() { }; std::sort(columnElements.begin(), columnElements.end(), predSortElements); std::sort(otherElements.begin(), otherElements.end(), predSortElements); - TDictStats columnStats = BuildStats(columnElements, Settings, CurrentRecordIndex, /*allowDictionary*/ true); + TDictStats columnStats = BuildStats(columnElements, Settings, CurrentRecordIndex, true); { ui32 columnIdx = 0; for (auto&& i : columnElements) { + const EValueType valueType = columnStats.GetValueType(columnIdx); switch (columnStats.GetAccessorType(columnIdx)) { case IChunkedArray::EType::Array: - i->BuildPlainAccessor(CurrentRecordIndex); + i->BuildPlainAccessor(CurrentRecordIndex, valueType); break; case IChunkedArray::EType::SparsedArray: - i->BuildSparsedAccessor(CurrentRecordIndex); + i->BuildSparsedAccessor(CurrentRecordIndex, valueType); break; case IChunkedArray::EType::Dictionary: - i->BuildDictionaryAccessor(CurrentRecordIndex); + i->BuildDictionaryAccessor(CurrentRecordIndex, valueType); break; case IChunkedArray::EType::Undefined: case IChunkedArray::EType::SerializedChunkedArray: @@ -96,8 +133,8 @@ std::shared_ptr<TSubColumnsArray> TDataBuilder::Finish() { TOthersData rbOthers = MergeOthers(otherElements, CurrentRecordIndex); auto records = std::make_shared<TGeneralContainer>(CurrentRecordIndex); - for (auto&& i : columnElements) { - records->AddField(std::make_shared<arrow::Field>(std::string(i->GetKeyName()), arrow::binary()), i->GetAccessorVerified()).Validate(); + for (size_t idx = 0; idx < columnElements.size(); ++idx) { + records->AddField(columnStats.GetField(idx), columnElements[idx]->GetAccessorVerified()).Validate(); } TColumnsData cData(std::move(columnStats), std::move(records)); return std::make_shared<TSubColumnsArray>(std::move(cData), std::move(rbOthers), Type, CurrentRecordIndex, Settings); @@ -125,7 +162,7 @@ TOthersData TDataBuilder::MergeOthers(const std::vector<TColumnElements*>& other std::push_heap(heap.begin(), heap.end()); } } - return othersBuilder->Finish(TOthersData::TFinishContext(BuildStats(otherKeys, Settings, recordsCount, /*allowDictionary*/ false))); + return othersBuilder->Finish(TOthersData::TFinishContext(BuildStats(otherKeys, Settings, recordsCount, false))); } std::string BuildString(const TStringBuf currentPrefix, const TStringBuf key) { 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 4607373d599..5c2148566b0 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/direct_builder.h +++ b/ydb/core/formats/arrow/accessor/sub_columns/direct_builder.h @@ -1,4 +1,5 @@ #pragma once +#include "types.h" #include "others_storage.h" #include "settings.h" #include "stats.h" @@ -36,9 +37,9 @@ public: return Accessor; } - void BuildSparsedAccessor(const ui32 recordsCount); - void BuildPlainAccessor(const ui32 recordsCount); - void BuildDictionaryAccessor(const ui32 recordsCount); + void BuildSparsedAccessor(const ui32 recordsCount, const EValueType valueType); + void BuildPlainAccessor(const ui32 recordsCount, const EValueType valueType); + void BuildDictionaryAccessor(const ui32 recordsCount, const EValueType valueType); TColumnElements(const TStringBuf key) : KeyName(key) { @@ -173,9 +174,8 @@ public: } }; - // allowDictionary is set only for the separated columns (the Others store is always plain). TDictStats BuildStats( - const std::vector<TColumnElements*>& keys, const TSettings& settings, const ui32 recordsCount, const bool allowDictionary) const { + const std::vector<TColumnElements*>& keys, const TSettings& settings, const ui32 recordsCount, const bool separateColumns) const { auto builder = TDictStats::MakeBuilder(); for (auto&& i : keys) { const ui32 presentCount = i->GetRecordIndexes().size(); @@ -187,13 +187,19 @@ public: } } }; + // Native scalar storage applies only to separated columns, the Others store is always BinaryJson. + EValueType valueType = EValueType::BinaryJson; + if (separateColumns && settings.GetEnableNativeColumnsResolved()) { + valueType = DetectValueTypeForArray(i->GetValues()); + } IChunkedArray::EType accessorType = IChunkedArray::EType::Array; if (settings.IsSparsed(presentCount, recordsCount)) { accessorType = IChunkedArray::EType::SparsedArray; - } else if (allowDictionary && settings.IsDictionary(presentCount, enumerateValues)) { + } else if (separateColumns && DictionaryApplicableForValueType(valueType) && + settings.IsDictionary(presentCount, enumerateValues)) { accessorType = IChunkedArray::EType::Dictionary; } - builder.Add(i->GetKeyName(), presentCount, i->GetDataSize(), accessorType); + builder.Add(i->GetKeyName(), presentCount, i->GetDataSize(), accessorType, valueType); } return builder.Finish(); } diff --git a/ydb/core/formats/arrow/accessor/sub_columns/header.cpp b/ydb/core/formats/arrow/accessor/sub_columns/header.cpp index f132a1be828..29bc1053248 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/header.cpp +++ b/ydb/core/formats/arrow/accessor/sub_columns/header.cpp @@ -18,9 +18,7 @@ TConclusion<TSubColumnsHeader> TSubColumnsHeader::ReadHeader(const TString& orig currentIndex += protoSize; TDictStats columnStats = [&]() { if (proto.GetColumnStatsSize()) { - std::shared_ptr<arrow::RecordBatch> rbColumnStats = - NArrow::DeserializeBatch(TString(originalData.data() + currentIndex, proto.GetColumnStatsSize()), TDictStats::GetStatsSchema()); - return TDictStats(rbColumnStats); + return TDictStats::DeserializeFromBlob(TString(originalData.data() + currentIndex, proto.GetColumnStatsSize())); } else { return TDictStats::BuildEmpty(); } @@ -28,9 +26,7 @@ TConclusion<TSubColumnsHeader> TSubColumnsHeader::ReadHeader(const TString& orig currentIndex += proto.GetColumnStatsSize(); TDictStats otherStats = [&]() { if (proto.GetOtherStatsSize()) { - std::shared_ptr<arrow::RecordBatch> rbOtherStats = - NArrow::DeserializeBatch(TString(originalData.data() + currentIndex, proto.GetOtherStatsSize()), TDictStats::GetStatsSchema()); - return TDictStats(rbOtherStats); + return TDictStats::DeserializeFromBlob(TString(originalData.data() + currentIndex, proto.GetOtherStatsSize())); } else { return TDictStats::BuildEmpty(); } diff --git a/ydb/core/formats/arrow/accessor/sub_columns/iterators.cpp b/ydb/core/formats/arrow/accessor/sub_columns/iterators.cpp index 19bc62d3d12..27b59bb364c 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/iterators.cpp +++ b/ydb/core/formats/arrow/accessor/sub_columns/iterators.cpp @@ -1,17 +1,11 @@ #include "iterators.h" -#include <yql/essentials/types/binary_json/read.h> +#include "types.h" namespace NKikimr::NArrow::NAccessor::NSubColumns { NJson::TJsonValue TGeneralIterator::GetValue() const { AFL_VERIFY(IsValidFlag); - if (RawValue.empty()) { - return NJson::TJsonValue(NJson::JSON_UNDEFINED); - } - auto data = NBinaryJson::SerializeToJson(TStringBuf(RawValue.data(), RawValue.size())); - NJson::TJsonValue res; - AFL_VERIFY(NJson::ReadJsonTree(data, &res)); - return res; + return ArrayElementToJsonValue(*CurrentArray, LocalIndex, ValueType); } } // namespace NKikimr::NArrow::NAccessor::NSubColumns diff --git a/ydb/core/formats/arrow/accessor/sub_columns/iterators.h b/ydb/core/formats/arrow/accessor/sub_columns/iterators.h index d4173391028..ebd89cddfe4 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/iterators.h +++ b/ydb/core/formats/arrow/accessor/sub_columns/iterators.h @@ -1,5 +1,6 @@ #pragma once #include "columns_storage.h" +#include "types.h" #include "others_storage.h" namespace NKikimr::NArrow::NAccessor::NSubColumns { @@ -13,15 +14,19 @@ private: ui32 KeyIndex = 0; bool IsValidFlag = false; bool HasValueFlag = false; - std::string_view RawValue; bool IsColumnKeyFlag = false; + EValueType ValueType = EValueType::BinaryJson; + // Current value as (array, local index); the reader interprets it per ValueType. + const arrow::Array* CurrentArray = nullptr; + i64 LocalIndex = 0; void InitFromIterator(const TColumnsData::TIterator& iterator) { RecordIndex = iterator.GetCurrentRecordIndex(); KeyIndex = RemappedKey.value_or(iterator.GetKeyIndex()); IsValidFlag = true; HasValueFlag = iterator.HasValue(); - RawValue = iterator.GetRawValue(); + CurrentArray = &iterator.GetArray(); + LocalIndex = iterator.GetLocalIndex(); } void InitFromIterator(const TOthersData::TIterator& iterator) { @@ -29,7 +34,8 @@ private: KeyIndex = RemapKeys.size() ? RemapKeys[iterator.GetKeyIndex()] : iterator.GetKeyIndex(); IsValidFlag = true; HasValueFlag = iterator.HasValue(); - RawValue = iterator.GetRawValue(); + CurrentArray = &iterator.GetArray(); + LocalIndex = iterator.GetLocalIndex(); } bool Initialize() { @@ -63,9 +69,10 @@ private: } public: - TGeneralIterator(TColumnsData::TIterator&& iterator, const std::optional<ui32> remappedKey = {}) + TGeneralIterator(TColumnsData::TIterator&& iterator, const EValueType valueType, const std::optional<ui32> remappedKey = {}) : Iterator(iterator) - , RemappedKey(remappedKey) { + , RemappedKey(remappedKey) + , ValueType(valueType) { Initialize(); } TGeneralIterator(TOthersData::TIterator&& iterator, const std::vector<ui32>& remapKeys = {}) @@ -149,9 +156,11 @@ public: return KeyIndex; } - std::string_view GetRawValue() const { + // The ordered iterator (compaction's reader) expects BinaryJson. BinaryJson columns pass their + // bytes through directly; native columns are re-encoded into an owned buffer. + NBinaryJson::TBinaryJson GetValueAsBinaryJson() { AFL_VERIFY(IsValidFlag); - return RawValue; + return ArrayElementToBinaryJson(*CurrentArray, LocalIndex, ValueType); } NJson::TJsonValue GetValue() const; @@ -181,7 +190,7 @@ public: : ColumnsData(columnsData) , OthersData(othersData) { for (ui32 i = 0; i < ColumnsData.GetStats().GetColumnsCount(); ++i) { - Iterators.emplace_back(ColumnsData.BuildIterator(i)); + Iterators.emplace_back(ColumnsData.BuildIterator(i), ColumnsData.GetStats().GetValueType(i)); } Iterators.emplace_back(OthersData.BuildIterator()); for (auto&& i : Iterators) { @@ -294,7 +303,7 @@ public: } } for (ui32 i = 0; i < ColumnsData.GetStats().GetColumnsCount(); ++i) { - Iterators.emplace_back(ColumnsData.BuildIterator(i), remapColumns[i]); + Iterators.emplace_back(ColumnsData.BuildIterator(i), ColumnsData.GetStats().GetValueType(i), remapColumns[i]); } Iterators.emplace_back(OthersData.BuildIterator(), remapOthers); for (auto&& i : Iterators) { @@ -325,7 +334,7 @@ public: while (SortedIterators.size() && SortedIterators.front()->GetRecordIndex() == recordIndex) { std::pop_heap(SortedIterators.begin(), SortedIterators.end(), TIteratorsComparator()); auto& itColumn = *SortedIterators.back(); - kvActor(Addresses[itColumn.GetKeyIndex()].GetOriginalIndex(), itColumn.GetRawValue(), itColumn.IsColumnKey()); + kvActor(Addresses[itColumn.GetKeyIndex()].GetOriginalIndex(), itColumn.GetValueAsBinaryJson(), itColumn.IsColumnKey()); if (!itColumn.Next()) { SortedIterators.pop_back(); } else { 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 9fd3e7323c8..d73ec7c023a 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/others_storage.cpp +++ b/ydb/core/formats/arrow/accessor/sub_columns/others_storage.cpp @@ -120,8 +120,8 @@ public: 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)); + statBuilder.Add(i.second.GetKeyName(), i.second.GetRecordsCount(), i.second.GetDataSize(), + i.second.GetAccessorType(settings, recordsCount), i.second.GetValueType()); } return statBuilder.Finish(); } 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 f2dfbd5ab53..d760289d04c 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/others_storage.h +++ b/ydb/core/formats/arrow/accessor/sub_columns/others_storage.h @@ -117,6 +117,17 @@ public: return res; } + // The Others store is always BinaryJson; exposed as (array, local index) like the column + // iterator so the general iterator reads both uniformly. + const arrow::Array& GetArray() const { + AFL_VERIFY(IsValid()); + return *Values; + } + i64 GetLocalIndex() const { + AFL_VERIFY(IsValid()); + return CurrentIndex; + } + NArrow::NAccessor::TBinaryJsonValueView GetValue() const; bool HasValue() const { diff --git a/ydb/core/formats/arrow/accessor/sub_columns/request.cpp b/ydb/core/formats/arrow/accessor/sub_columns/request.cpp index 9c5a0acbe47..af340d1d85a 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/request.cpp +++ b/ydb/core/formats/arrow/accessor/sub_columns/request.cpp @@ -25,7 +25,10 @@ TConclusionStatus TRequestedConstuctor::DoDeserializeFromRequest(NYql::TFeatures if (*fraction < 0 || *fraction > 1) { return TConclusionStatus::Fail("DICTIONARY_UNIQUE_FRACTION must be in [0, 1] interval"); } - Settings.SetDictionaryUniqueFraction(*fraction); + Settings.SetDictionaryUniqueFraction(fraction); + } + if (auto enable = features.Extract<bool>("ENABLE_NATIVE_COLUMNS")) { + Settings.SetEnableNativeColumns(enable); } THolder<IDataAdapter> extractor; if (auto dataExtractorClassName = features.Extract<TString>("DATA_EXTRACTOR_CLASS_NAME")) { diff --git a/ydb/core/formats/arrow/accessor/sub_columns/settings.h b/ydb/core/formats/arrow/accessor/sub_columns/settings.h index aa05eb6a9b9..764cbff2e6d 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/settings.h +++ b/ydb/core/formats/arrow/accessor/sub_columns/settings.h @@ -19,7 +19,8 @@ private: YDB_ACCESSOR(ui32, ColumnsLimit, 1024); YDB_ACCESSOR(ui32, ChunkMemoryLimit, 50 * 1024 * 1024); YDB_READONLY(double, OthersAllowedFraction, 0.05); - YDB_ACCESSOR(double, DictionaryUniqueFraction, 0); + YDB_ACCESSOR(std::optional<double>, DictionaryUniqueFraction, std::nullopt); + YDB_ACCESSOR(std::optional<bool>, EnableNativeColumns, std::nullopt); YDB_ACCESSOR_DEF(TDataAdapterContainer, DataExtractor); public: @@ -76,7 +77,12 @@ public: result.InsertValue("columns_limit", ColumnsLimit); result.InsertValue("memory_limit", ChunkMemoryLimit); result.InsertValue("others_allowed_fraction", OthersAllowedFraction); - result.InsertValue("dictionary_unique_fraction", DictionaryUniqueFraction); + if (DictionaryUniqueFraction) { + result.InsertValue("dictionary_unique_fraction", *DictionaryUniqueFraction); + } + if (EnableNativeColumns) { + result.InsertValue("enable_native_columns", *EnableNativeColumns); + } result.InsertValue("data_extractor", DataExtractor->DebugJson()); return result; } @@ -113,13 +119,14 @@ public: // `enumerate` feeds the values counted toward the distinct set. Short-circuits as soon as the verdict is clear. template <class TEnumerator> bool IsDictionary(const ui32 presentCount, const TEnumerator& enumerate) const { - if (DictionaryUniqueFraction == 0) { + double fraction = GetDictionaryUniqueFractionResolved(); + if (fraction == 0) { return false; } - if (DictionaryUniqueFraction == 1) { + if (fraction == 1) { return true; } - const ui32 tooManyUnique = static_cast<ui32>(DictionaryUniqueFraction * presentCount) + 1; + const ui32 tooManyUnique = static_cast<ui32>(fraction * presentCount) + 1; return !HasAtLeastUniqueValues(tooManyUnique, enumerate); } @@ -129,7 +136,14 @@ public: result.SetColumnsLimit(ColumnsLimit); result.SetChunkMemoryLimit(ChunkMemoryLimit); result.SetOthersAllowedFraction(OthersAllowedFraction); - result.SetDictionaryUniqueFraction(DictionaryUniqueFraction); + // These are set only when non-default, so an unset option (equivalent to its default) + // is not persisted and not surfaced back by SHOW CREATE. + if (DictionaryUniqueFraction) { + result.SetDictionaryUniqueFraction(*DictionaryUniqueFraction); + } + if (EnableNativeColumns) { + result.SetEnableNativeColumns(*EnableNativeColumns); + } DataExtractor.SerializeToProto(*result.MutableDataExtractor()); } @@ -139,7 +153,12 @@ public: ColumnsLimit = proto.GetColumnsLimit(); ChunkMemoryLimit = proto.GetChunkMemoryLimit(); OthersAllowedFraction = proto.GetOthersAllowedFraction(); - DictionaryUniqueFraction = proto.GetDictionaryUniqueFraction(); + if (proto.HasDictionaryUniqueFraction()) { + DictionaryUniqueFraction = proto.GetDictionaryUniqueFraction(); + } + if (proto.HasEnableNativeColumns()) { + EnableNativeColumns = proto.GetEnableNativeColumns(); + } if (!proto.HasDataExtractor()) { AFL_VERIFY(DataExtractor.Initialize(TJsonScanExtractor::GetClassNameStatic())); } else if (!DataExtractor.DeserializeFromProto(proto.GetDataExtractor())) { @@ -148,6 +167,14 @@ public: return true; } + double GetDictionaryUniqueFractionResolved() const { + return DictionaryUniqueFraction.value_or(0); + } + + bool GetEnableNativeColumnsResolved() const { + return EnableNativeColumns.value_or(false); + } + NKikimrArrowAccessorProto::TConstructor::TSubColumns::TSettings SerializeToProto() const { NKikimrArrowAccessorProto::TConstructor::TSubColumns::TSettings result; SerializeToProtoImpl(result); diff --git a/ydb/core/formats/arrow/accessor/sub_columns/stats.cpp b/ydb/core/formats/arrow/accessor/sub_columns/stats.cpp index 3673b6308c4..b08040273ae 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/stats.cpp +++ b/ydb/core/formats/arrow/accessor/sub_columns/stats.cpp @@ -5,10 +5,13 @@ #include <ydb/core/formats/arrow/accessor/plain/constructor.h> #include <ydb/core/formats/arrow/accessor/sparsed/constructor.h> #include <ydb/core/formats/arrow/serializer/abstract.h> +#include <ydb/core/formats/arrow/serializer/native.h> #include <ydb/library/formats/arrow/arrow_helpers.h> +#include <ydb/library/formats/arrow/simple_arrays_cache.h> namespace NKikimr::NArrow::NAccessor::NSubColumns { + TSplittedColumns TDictStats::SplitByVolume(const TSettings& settings, const ui32 recordsCount) const { std::map<ui64, std::vector<TRTStats>> bySize; ui64 sumSize = 0; @@ -36,10 +39,10 @@ TSplittedColumns TDictStats::SplitByVolume(const TSettings& settings, const ui32 auto columnsBuilder = MakeBuilder(); auto othersBuilder = MakeBuilder(); for (auto&& i : columnStats) { - columnsBuilder.Add(i.GetKeyName(), i.GetRecordsCount(), i.GetDataSize(), i.GetAccessorType(settings, recordsCount)); + columnsBuilder.Add(i.GetKeyName(), i.GetRecordsCount(), i.GetDataSize(), i.GetAccessorType(settings, recordsCount), i.GetValueType()); } for (auto&& i : otherStats) { - othersBuilder.Add(i.GetKeyName(), i.GetRecordsCount(), i.GetDataSize(), i.GetAccessorType(settings, recordsCount)); + othersBuilder.Add(i.GetKeyName(), i.GetRecordsCount(), i.GetDataSize(), i.GetAccessorType(settings, recordsCount), i.GetValueType()); } return TSplittedColumns(columnsBuilder.Finish(), othersBuilder.Finish()); } @@ -57,7 +60,8 @@ TDictStats TDictStats::Merge(const std::vector<const TDictStats*>& stats, const } auto builder = MakeBuilder(); for (auto&& i : resultMap) { - builder.Add(i.second.GetKeyName(), i.second.GetRecordsCount(), i.second.GetDataSize(), i.second.GetAccessorType(settings, recordsCount)); + builder.Add(i.second.GetKeyName(), i.second.GetRecordsCount(), i.second.GetDataSize(), + i.second.GetAccessorType(settings, recordsCount), i.second.GetValueType()); } return builder.Finish(); } @@ -80,15 +84,24 @@ std::string_view TDictStats::GetColumnName(const ui32 index) const { TDictStats::TDictStats(const std::shared_ptr<arrow::RecordBatch>& original) : Original(original) { - AFL_VERIFY(Original->num_columns() == 4)("count", Original->num_columns()); + AFL_VERIFY(Original->num_columns() == 4 || Original->num_columns() == 5)("count", Original->num_columns()); AFL_VERIFY(Original->column(0)->type()->id() == arrow::binary()->id()); AFL_VERIFY(Original->column(1)->type()->id() == arrow::uint32()->id()); AFL_VERIFY(Original->column(2)->type()->id() == arrow::uint32()->id()); AFL_VERIFY(Original->column(3)->type()->id() == arrow::uint8()->id()); - DataNames = std::static_pointer_cast<arrow::StringArray>(Original->column(0)); + if (Original->num_columns() == 4) { + // Legacy stats (pre native scalar columns): synthesize an all-BinaryJson value_type column. + auto valueTypeArray = NArrow::TThreadSimpleArraysCache::Get( + arrow::uint8(), std::make_shared<arrow::UInt8Scalar>((ui8)EValueType::BinaryJson), Original->num_rows()); + Original = arrow::RecordBatch::Make(GetStatsSchema(), Original->num_rows(), + { Original->column(0), Original->column(1), Original->column(2), Original->column(3), valueTypeArray }); + } + AFL_VERIFY(Original->column(4)->type()->id() == arrow::uint8()->id()); + DataNames = std::static_pointer_cast<arrow::BinaryArray>(Original->column(0)); DataRecordsCount = std::static_pointer_cast<arrow::UInt32Array>(Original->column(1)); DataSize = std::static_pointer_cast<arrow::UInt32Array>(Original->column(2)); AccessorType = std::static_pointer_cast<arrow::UInt8Array>(Original->column(3)); + ValueType = std::static_pointer_cast<arrow::UInt8Array>(Original->column(4)); } TConstructorContainer TDictStats::GetAccessorConstructor(const ui32 columnIndex) const { @@ -115,6 +128,17 @@ TDictStats TDictStats::BuildEmpty() { return result; } +TDictStats TDictStats::DeserializeFromBlob(const TString& blob) { + NSerialization::TNativeSerializer serializer; + auto result = serializer.Deserialize(blob, GetStatsSchema()); + if (result.ok()) { + return TDictStats(*result); + } + auto legacy = serializer.Deserialize(blob, GetStatsSchemaLegacy()); + AFL_VERIFY(legacy.ok())("error", legacy.status().ToString()); + return TDictStats(*legacy); +} + TString TDictStats::SerializeAsString(const std::shared_ptr<NSerialization::ISerializer>& serializer) const { if (serializer) { AFL_VERIFY(serializer); @@ -129,20 +153,28 @@ IChunkedArray::EType TDictStats::GetAccessorType(const ui32 columnIndex) const { return (IChunkedArray::EType)AccessorType->Value(columnIndex); } +EValueType TDictStats::GetValueType(const ui32 columnIndex) const { + AFL_VERIFY(columnIndex < ValueType->length()); + return (EValueType)ValueType->Value(columnIndex); +} + TDictStats::TBuilder::TBuilder() { Builders = NArrow::MakeBuilders(GetStatsSchema()); - AFL_VERIFY(Builders.size() == 4); + AFL_VERIFY(Builders.size() == 5); AFL_VERIFY(Builders[0]->type()->id() == arrow::binary()->id()); AFL_VERIFY(Builders[1]->type()->id() == arrow::uint32()->id()); AFL_VERIFY(Builders[2]->type()->id() == arrow::uint32()->id()); AFL_VERIFY(Builders[3]->type()->id() == arrow::uint8()->id()); - Names = static_cast<arrow::StringBuilder*>(Builders[0].get()); + AFL_VERIFY(Builders[4]->type()->id() == arrow::uint8()->id()); + Names = static_cast<arrow::BinaryBuilder*>(Builders[0].get()); Records = static_cast<arrow::UInt32Builder*>(Builders[1].get()); DataSize = static_cast<arrow::UInt32Builder*>(Builders[2].get()); AccessorType = static_cast<arrow::UInt8Builder*>(Builders[3].get()); + ValueType = static_cast<arrow::UInt8Builder*>(Builders[4].get()); } -void TDictStats::TBuilder::Add(const TString& name, const ui32 recordsCount, const ui32 dataSize, const IChunkedArray::EType accessorType) { +void TDictStats::TBuilder::Add(const TString& name, const ui32 recordsCount, const ui32 dataSize, const IChunkedArray::EType accessorType, + const EValueType valueType) { AFL_VERIFY(Builders.size()); if (!LastKeyName) { LastKeyName = name; @@ -156,12 +188,13 @@ void TDictStats::TBuilder::Add(const TString& name, const ui32 recordsCount, con TStatusValidator::Validate(Records->Append(recordsCount)); TStatusValidator::Validate(DataSize->Append(dataSize)); TStatusValidator::Validate(AccessorType->Append((ui8)accessorType)); + TStatusValidator::Validate(ValueType->Append((ui8)valueType)); ++RecordsCount; } void TDictStats::TBuilder::Add( - const std::string_view name, const ui32 recordsCount, const ui32 dataSize, const IChunkedArray::EType accessorType) { - Add(TString(name.data(), name.size()), recordsCount, dataSize, accessorType); + const std::string_view name, const ui32 recordsCount, const ui32 dataSize, const IChunkedArray::EType accessorType, const EValueType valueType) { + Add(TString(name.data(), name.size()), recordsCount, dataSize, accessorType, valueType); } TDictStats TDictStats::TBuilder::Finish() { diff --git a/ydb/core/formats/arrow/accessor/sub_columns/stats.h b/ydb/core/formats/arrow/accessor/sub_columns/stats.h index 3bd31b0da03..7522e4c51d6 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/stats.h +++ b/ydb/core/formats/arrow/accessor/sub_columns/stats.h @@ -1,5 +1,6 @@ #pragma once #include "settings.h" +#include "types.h" #include <ydb/core/formats/arrow/accessor/abstract/constructor.h> #include <ydb/core/formats/arrow/accessor/sub_columns/json_value_path.h> @@ -25,6 +26,7 @@ private: std::shared_ptr<arrow::UInt32Array> DataRecordsCount; std::shared_ptr<arrow::UInt32Array> DataSize; std::shared_ptr<arrow::UInt8Array> AccessorType; + std::shared_ptr<arrow::UInt8Array> ValueType; TJsonPathAccessorTriePtr CachedJsonPathAccessorTrie; TJsonPathAccessorTriePtr GenerateJsonPathAccessorTrie() const { @@ -55,11 +57,16 @@ public: result.InsertValue("records", NArrow::DebugJson(DataRecordsCount, 1000000, 1000000)["data"]); result.InsertValue("size", NArrow::DebugJson(DataSize, 1000000, 1000000)["data"]); result.InsertValue("accessor", NArrow::DebugJson(AccessorType, 1000000, 1000000)["data"]); + result.InsertValue("value_type", NArrow::DebugJson(ValueType, 1000000, 1000000)["data"]); return result; } static TDictStats BuildEmpty(); TString SerializeAsString(const std::shared_ptr<NSerialization::ISerializer>& serializer) const; + // Deserialize a stats blob, transparently accepting both the current (5-column, with value_type) + // and the legacy (4-column) layout via try-decode-fallback. + static TDictStats DeserializeFromBlob(const TString& blob); + void CreateJsonPathAccessorTrieCache() { CachedJsonPathAccessorTrie = GenerateJsonPathAccessorTrie(); } @@ -88,12 +95,18 @@ public: private: YDB_READONLY(ui32, RecordsCount, 0); YDB_READONLY(ui32, DataSize, 0); + std::optional<EValueType> DeducedValueType; public: TRTStatsValue() = default; - TRTStatsValue(const ui32 recordsCount, const ui32 dataSize) + TRTStatsValue(const ui32 recordsCount, const ui32 dataSize, const std::optional<EValueType>& valueType) : RecordsCount(recordsCount) - , DataSize(dataSize) { + , DataSize(dataSize) + , DeducedValueType(valueType) { + } + + EValueType GetValueType() const { + return DeducedValueType.value_or(EValueType::BinaryJson); } void AddValue(const std::string_view str) { @@ -104,6 +117,7 @@ public: void Add(const TDictStats& stats, const ui32 idx) { RecordsCount += stats.GetColumnRecordsCount(idx); DataSize += stats.GetColumnSize(idx); + DeducedValueType = MergeValueTypes(DeducedValueType, stats.GetValueType(idx)); } // Decides only the Array-vs-Sparsed axis and never returns Dictionary, @@ -122,16 +136,16 @@ public: TRTStats(const TString& keyName) : KeyName(keyName) { } - TRTStats(const TString& keyName, const ui32 recordsCount, const ui32 dataSize) - : TBase(recordsCount, dataSize) + TRTStats(const TString& keyName, const ui32 recordsCount, const ui32 dataSize, const std::optional<EValueType>& valueType) + : TBase(recordsCount, dataSize, valueType) , KeyName(keyName) { } TRTStats(const std::string_view keyName) : KeyName(keyName.data(), keyName.size()) { } - TRTStats(const std::string_view keyName, const ui32 recordsCount, const ui32 dataSize) - : TBase(recordsCount, dataSize) + TRTStats(const std::string_view keyName, const ui32 recordsCount, const ui32 dataSize, const std::optional<EValueType>& valueType) + : TBase(recordsCount, dataSize, valueType) , KeyName(keyName.data(), keyName.size()) { } @@ -151,14 +165,17 @@ public: arrow::UInt32Builder* Records; arrow::UInt32Builder* DataSize; arrow::UInt8Builder* AccessorType; + arrow::UInt8Builder* ValueType; std::optional<TString> LastKeyName; ui32 RecordsCount = 0; public: TBuilder(); - void Add(const TString& name, const ui32 recordsCount, const ui32 dataSize, const IChunkedArray::EType accessorType); - void Add(const std::string_view name, const ui32 recordsCount, const ui32 dataSize, const IChunkedArray::EType accessorType); + void Add(const TString& name, const ui32 recordsCount, const ui32 dataSize, const IChunkedArray::EType accessorType, + const EValueType valueType); + void Add(const std::string_view name, const ui32 recordsCount, const ui32 dataSize, const IChunkedArray::EType accessorType, + const EValueType valueType); TDictStats Finish(); }; @@ -169,8 +186,7 @@ public: std::shared_ptr<arrow::Schema> BuildColumnsSchema() const { arrow::FieldVector fields; for (ui32 i = 0; i < DataNames->length(); ++i) { - const auto view = DataNames->GetView(i); - fields.emplace_back(std::make_shared<arrow::Field>(std::string(view.data(), view.size()), arrow::binary())); + fields.emplace_back(GetField(i)); } return std::make_shared<arrow::Schema>(fields); } @@ -178,12 +194,12 @@ public: std::shared_ptr<arrow::Field> GetField(const ui32 index) const { AFL_VERIFY(index < DataNames->length()); auto name = DataNames->GetView(index); - return std::make_shared<arrow::Field>(std::string(name.data(), name.size()), arrow::binary()); + return std::make_shared<arrow::Field>(std::string(name.data(), name.size()), GetArrowTypeForValueType(GetValueType(index))); } TRTStats GetRTStats(const ui32 index) const { auto view = GetColumnName(index); - return TRTStats(TString(view.data(), view.size()), GetColumnRecordsCount(index), GetColumnSize(index)); + return TRTStats(TString(view.data(), view.size()), GetColumnRecordsCount(index), GetColumnSize(index), GetValueType(index)); } ui32 GetDataNamesCount() const { @@ -196,6 +212,7 @@ public: TConstructorContainer GetAccessorConstructor(const ui32 columnIndex) const; IChunkedArray::EType GetAccessorType(const ui32 columnIndex) const; + EValueType GetValueType(const ui32 columnIndex) const; std::string_view GetColumnName(const ui32 index) const; TString GetColumnNameString(const ui32 index) const { @@ -208,6 +225,16 @@ public: static std::shared_ptr<arrow::Schema> GetStatsSchema() { static arrow::FieldVector fields = { std::make_shared<arrow::Field>("name", arrow::binary()), std::make_shared<arrow::Field>("count", arrow::uint32()), std::make_shared<arrow::Field>("size", arrow::uint32()), + std::make_shared<arrow::Field>("accessor_type", arrow::uint8()), std::make_shared<arrow::Field>("value_type", arrow::uint8()) }; + static std::shared_ptr<arrow::Schema> result = std::make_shared<arrow::Schema>(fields); + return result; + } + + // Legacy schema without value_type; used only as the fallback when deserializing blobs written + // before native scalar columns existed. + static std::shared_ptr<arrow::Schema> GetStatsSchemaLegacy() { + static arrow::FieldVector fields = { std::make_shared<arrow::Field>("name", arrow::binary()), + std::make_shared<arrow::Field>("count", arrow::uint32()), std::make_shared<arrow::Field>("size", arrow::uint32()), std::make_shared<arrow::Field>("accessor_type", arrow::uint8()) }; static std::shared_ptr<arrow::Schema> result = std::make_shared<arrow::Schema>(fields); return result; diff --git a/ydb/core/formats/arrow/accessor/sub_columns/types.cpp b/ydb/core/formats/arrow/accessor/sub_columns/types.cpp new file mode 100644 index 00000000000..81c988c8e99 --- /dev/null +++ b/ydb/core/formats/arrow/accessor/sub_columns/types.cpp @@ -0,0 +1,134 @@ +#include "types.h" + +#include <contrib/libs/apache/arrow/cpp/src/arrow/array/array_binary.h> +#include <contrib/libs/apache/arrow/cpp/src/arrow/array/array_primitive.h> + +#include <library/cpp/json/json_reader.h> +#include <library/cpp/json/json_writer.h> + +#include <ydb/library/actors/core/log.h> + +#include <yql/essentials/types/binary_json/read.h> +#include <yql/essentials/types/binary_json/write.h> + +namespace NKikimr::NArrow::NAccessor::NSubColumns { + +namespace { + +NBinaryJson::TBinaryJson ToBinaryJson(const NJson::TJsonValue& json) { + auto result = NBinaryJson::SerializeToBinaryJson(NJson::WriteJson(&json, false)); + AFL_VERIFY(std::holds_alternative<NBinaryJson::TBinaryJson>(result)); + return std::get<NBinaryJson::TBinaryJson>(std::move(result)); +} + +EValueType ValueTypeForItem(const NBinaryJson::TBinaryJson& blob) { + auto reader = NBinaryJson::TBinaryJsonReader::Make(blob); + auto rootCursor = reader->GetRootCursor(); + if (rootCursor.GetType() != NBinaryJson::EContainerType::TopLevelScalar) { + return EValueType::BinaryJson; + } + switch (rootCursor.GetElement(0).GetType()) { + case NBinaryJson::EEntryType::String: + return EValueType::String; + case NBinaryJson::EEntryType::Number: + return EValueType::Double; + case NBinaryJson::EEntryType::BoolFalse: + case NBinaryJson::EEntryType::BoolTrue: + return EValueType::Bool; + case NBinaryJson::EEntryType::Container: + case NBinaryJson::EEntryType::Null: + return EValueType::BinaryJson; + } +} + +} // namespace + +std::shared_ptr<arrow::DataType> GetArrowTypeForValueType(const EValueType valueType) { + switch (valueType) { + case EValueType::BinaryJson: + case EValueType::String: + return arrow::binary(); + case EValueType::Double: + return arrow::float64(); + case EValueType::Bool: + return arrow::boolean(); + } +} + +bool DictionaryApplicableForValueType(const EValueType valueType) { + switch (valueType) { + case EValueType::BinaryJson: + case EValueType::String: + return true; + default: + return false; + } +} + +EValueType MergeValueTypes(const std::optional<EValueType>& acc, const EValueType next) { + if (!acc) { + return next; + } + return (*acc == next) ? *acc : EValueType::BinaryJson; +} + +EValueType DetectValueTypeForArray(const std::deque<NBinaryJson::TBinaryJson>& values) { + std::optional<EValueType> common; + for (const auto& v : values) { + common = MergeValueTypes(common, ValueTypeForItem(v)); + if (*common == EValueType::BinaryJson) { + break; + } + } + return common.value_or(EValueType::BinaryJson); +} + +TStringBuf ExtractStringScalar(const NBinaryJson::TBinaryJson& blob) { + auto reader = NBinaryJson::TBinaryJsonReader::Make(blob); + return reader->GetRootCursor().GetElement(0).GetString(); +} + +double ExtractDoubleScalar(const NBinaryJson::TBinaryJson& blob) { + auto reader = NBinaryJson::TBinaryJsonReader::Make(blob); + return reader->GetRootCursor().GetElement(0).GetNumber(); +} + +bool ExtractBoolScalar(const NBinaryJson::TBinaryJson& blob) { + auto reader = NBinaryJson::TBinaryJsonReader::Make(blob); + return reader->GetRootCursor().GetElement(0).GetType() == NBinaryJson::EEntryType::BoolTrue; +} + +NJson::TJsonValue ArrayElementToJsonValue(const arrow::Array& array, const i64 index, const EValueType valueType) { + switch (valueType) { + case EValueType::String: { + const auto view = static_cast<const arrow::BinaryArray&>(array).GetView(index); + return NJson::TJsonValue(TStringBuf(view.data(), view.size())); + } + case EValueType::Double: + return NJson::TJsonValue(static_cast<const arrow::DoubleArray&>(array).Value(index)); + case EValueType::Bool: + return NJson::TJsonValue(static_cast<const arrow::BooleanArray&>(array).Value(index)); + case EValueType::BinaryJson: { + const auto view = static_cast<const arrow::BinaryArray&>(array).GetView(index); + const auto text = NBinaryJson::SerializeToJson(TStringBuf(view.data(), view.size())); + NJson::TJsonValue result; + AFL_VERIFY(NJson::ReadJsonTree(text, &result)); + return result; + } + } +} + +NBinaryJson::TBinaryJson ArrayElementToBinaryJson(const arrow::Array& array, const i64 index, const EValueType valueType) { + switch (valueType) { + case EValueType::BinaryJson: { + const auto view = static_cast<const arrow::BinaryArray&>(array).GetView(index); + return NBinaryJson::TBinaryJson(view.data(), view.size()); + } + case EValueType::String: + case EValueType::Double: + case EValueType::Bool: + return ToBinaryJson(ArrayElementToJsonValue(array, index, valueType)); + } +} + +} // namespace NKikimr::NArrow::NAccessor::NSubColumns diff --git a/ydb/core/formats/arrow/accessor/sub_columns/types.h b/ydb/core/formats/arrow/accessor/sub_columns/types.h new file mode 100644 index 00000000000..29ae5b25f49 --- /dev/null +++ b/ydb/core/formats/arrow/accessor/sub_columns/types.h @@ -0,0 +1,47 @@ +#pragma once + +#include <deque> + +#include <contrib/libs/apache/arrow/cpp/src/arrow/array/array_base.h> + +#include <library/cpp/json/writer/json_value.h> + +#include <yql/essentials/types/binary_json/format.h> + +// Type conversions between BinaryJson, dedicated scalar types and arrow storage types. +namespace NKikimr::NArrow::NAccessor::NSubColumns { + +// Logical (as seen by external consumers) value type of data stored in a subcolumn. +// PERSISTED: these numeric codes are written to disk as the `value_type` column of TDictStats. +enum class EValueType : ui8 { + BinaryJson = 0, + Double = 1, + Bool = 2, + String = 3, +}; + +std::shared_ptr<arrow::DataType> GetArrowTypeForValueType(const EValueType valueType); + + +// Dictionary encoding only enabled for the binary-backed types (BinaryJson blobs and raw strings). +// Integral types would trade fixed-size position in array for fixed-size dictionary ref. +// May still be good for compression, but requires further experiments. +bool DictionaryApplicableForValueType(const EValueType valueType); + +// Element type to represent result of merging arrays with arg types +EValueType MergeValueTypes(const std::optional<EValueType>& acc, const EValueType next); + +EValueType DetectValueTypeForArray(const std::deque<NBinaryJson::TBinaryJson>& values); + +// Convert json blob to its contained scalar. +// Fail on verify if blob is not of specified type. +TStringBuf ExtractStringScalar(const NBinaryJson::TBinaryJson& blob); +double ExtractDoubleScalar(const NBinaryJson::TBinaryJson& blob); +bool ExtractBoolScalar(const NBinaryJson::TBinaryJson& blob); + +// Read element `index` of a materialized native column array (interpreted per valueType) as a JSON +// value (document reconstruction) or a BinaryJson blob. The array's physical arrow type must match valueType. +NJson::TJsonValue ArrayElementToJsonValue(const arrow::Array& array, const i64 index, const EValueType valueType); +NBinaryJson::TBinaryJson ArrayElementToBinaryJson(const arrow::Array& array, const i64 index, const EValueType valueType); + +} // namespace NKikimr::NArrow::NAccessor::NSubColumns diff --git a/ydb/core/formats/arrow/accessor/sub_columns/ut/ut_helpers.h b/ydb/core/formats/arrow/accessor/sub_columns/ut/ut_helpers.h new file mode 100644 index 00000000000..cedb7b07cd8 --- /dev/null +++ b/ydb/core/formats/arrow/accessor/sub_columns/ut/ut_helpers.h @@ -0,0 +1,38 @@ +#pragma once + +#include <contrib/libs/apache/arrow/cpp/src/arrow/array/array_binary.h> +#include <contrib/libs/apache/arrow/cpp/src/arrow/chunked_array.h> + +#include <ydb/library/actors/core/log.h> + +#include <yql/essentials/types/binary_json/read.h> + +#include <util/generic/string.h> + +namespace NKikimr::NArrow::NAccessor::NSubColumns::NTesting { + +// Reconstruct the JSON documents of a sub-columns array as text, for round-trip assertions. +inline TString PrintBinaryJsons(const std::shared_ptr<arrow::ChunkedArray>& array) { + TStringBuilder sb; + sb << "["; + for (auto&& i : array->chunks()) { + sb << "["; + AFL_VERIFY(i->type()->id() == arrow::binary()->id()); + auto views = std::static_pointer_cast<arrow::BinaryArray>(i); + for (ui32 r = 0; r < views->length(); ++r) { + if (views->IsNull(r)) { + sb << "null"; + } else { + sb << NKikimr::NBinaryJson::SerializeToJson(TStringBuf(views->GetView(r).data(), views->GetView(r).size())); + } + if (r + 1 != views->length()) { + sb << ","; + } + } + sb << "]"; + } + sb << "]"; + return sb; +} + +} // namespace NKikimr::NArrow::NAccessor::NSubColumns::NTesting diff --git a/ydb/core/formats/arrow/accessor/sub_columns/ut/ut_native_scalars.cpp b/ydb/core/formats/arrow/accessor/sub_columns/ut/ut_native_scalars.cpp new file mode 100644 index 00000000000..87bd9441d97 --- /dev/null +++ b/ydb/core/formats/arrow/accessor/sub_columns/ut/ut_native_scalars.cpp @@ -0,0 +1,210 @@ +#include <ydb/core/formats/arrow/accessor/common/chunk_data.h> +#include <ydb/core/formats/arrow/accessor/plain/accessor.h> +#include <ydb/core/formats/arrow/accessor/sub_columns/accessor.h> +#include <ydb/core/formats/arrow/accessor/sub_columns/constructor.h> +#include <ydb/core/formats/arrow/serializer/abstract.h> + +#include "ut_helpers.h" + +#include <library/cpp/testing/unittest/registar.h> +#include <yql/essentials/types/binary_json/read.h> +#include <yql/essentials/types/binary_json/write.h> + +#include <algorithm> + +using NKikimr::NArrow::NAccessor::NSubColumns::NTesting::PrintBinaryJsons; + +Y_UNIT_TEST_SUITE(SubColumnsNativeScalars) { + using namespace NKikimr; + using namespace NKikimr::NArrow; + using namespace NKikimr::NArrow::NAccessor; + using namespace NKikimr::NArrow::NAccessor::NSubColumns; + + NSubColumns::TSettings NativeSettings(const double dictFraction) { + NSubColumns::TSettings s(4, 1024, 0, 0, TDataAdapterContainer::GetDefault(), dictFraction); + s.SetEnableNativeColumns(true); + return s; + } + + NSubColumns::TSettings OffSettings() { + return NSubColumns::TSettings(4, 1024, 0, 0, TDataAdapterContainer::GetDefault()); + } + + std::shared_ptr<TSubColumnsArray> BuildSubColumns(const std::vector<TString>& jsons, const NSubColumns::TSettings& settings) { + TTrivialArray::TPlainBuilder<arrow::BinaryType> b; + ui32 idx = 0; + for (auto&& j : jsons) { + if (j != "null") { + auto v = NBinaryJson::SerializeToBinaryJson(j); + auto* bj = std::get_if<NBinaryJson::TBinaryJson>(&v); + UNIT_ASSERT(bj); + b.AddRecord(idx, std::string_view(bj->data(), bj->size())); + } + ++idx; + } + auto arr = b.Finish(jsons.size()); + return TSubColumnsArray::Make(arr, settings, arr->GetDataType()).DetachResult(); + } + + // Number of separated columns stored with the given value type. + ui32 CountValueType(const std::shared_ptr<TSubColumnsArray>& arr, const EValueType vt) { + const auto& stats = arr->GetColumnsData().GetStats(); + ui32 n = 0; + for (ui32 i = 0; i < stats.GetColumnsCount(); ++i) { + n += (stats.GetValueType(i) == vt); + } + return n; + } + + std::shared_ptr<TSubColumnsArray> SerializeRoundTrip(const std::shared_ptr<TSubColumnsArray>& arr, const NSubColumns::TSettings& settings) { + auto serializer = NSerialization::TSerializerContainer::GetDefaultSerializer(); + TChunkConstructionData cData(arr->GetRecordsCount(), nullptr, arrow::binary(), serializer); + NSubColumns::TConstructor constructor(settings); + return std::static_pointer_cast<TSubColumnsArray>(constructor.DeserializeFromString(arr->SerializeToString(cData), cData).DetachResult()); + } + + void AssertColumnType(const std::vector<TString>& docs, const EValueType vt, IChunkedArray::EType type, std::string description) { + auto arr = BuildSubColumns(docs, NativeSettings(1)); + const auto& stats = arr->GetColumnsData().GetStats(); + UNIT_ASSERT_VALUES_EQUAL(stats.GetColumnsCount(), 1); + UNIT_ASSERT_C(stats.GetAccessorType(0) == type && stats.GetValueType(0) == vt, + TStringBuilder() << "expected a " << description << " column with value_type " << (ui32)vt << ": " << arr->DebugJson().GetStringRobust()); + auto restored = SerializeRoundTrip(arr, NativeSettings(1)); + UNIT_ASSERT_VALUES_EQUAL(PrintBinaryJsons(restored->GetChunkedArray()), PrintBinaryJsons(arr->GetChunkedArray())); + } + + void AssertDictColumn(const std::vector<TString>& docs, const EValueType vt) { + AssertColumnType(docs, vt, IChunkedArray::EType::Dictionary, "Dictionary"); + } + + void AssertPlainNativeColumn(const std::vector<TString>& docs, const EValueType vt) { + AssertColumnType(docs, vt, IChunkedArray::EType::Array, "plain native"); + } + + Y_UNIT_TEST(Detection) { + const std::vector<TString> docs = { + R"({"s":"x","n":1,"b":true})", + R"({"s":"yy","n":2,"b":false})", + R"({"s":"zzz","n":3,"b":true})", + }; + auto arr = BuildSubColumns(docs, NativeSettings(0)); + UNIT_ASSERT_VALUES_EQUAL_C(CountValueType(arr, EValueType::String), 1, arr->DebugJson().GetStringRobust()); + UNIT_ASSERT_VALUES_EQUAL_C(CountValueType(arr, EValueType::Double), 1, arr->DebugJson().GetStringRobust()); + UNIT_ASSERT_VALUES_EQUAL_C(CountValueType(arr, EValueType::Bool), 1, arr->DebugJson().GetStringRobust()); + UNIT_ASSERT_VALUES_EQUAL(PrintBinaryJsons(arr->GetChunkedArray()), PrintBinaryJsons(BuildSubColumns(docs, OffSettings())->GetChunkedArray())); + } + + Y_UNIT_TEST(DoubleRoundTrip) { + const std::vector<TString> docs = { + R"({"n":3.5})", R"({"n":-2})", R"({"n":0})", R"({"n":1000000})", R"({"n":2.718281828})", + }; + auto arr = BuildSubColumns(docs, NativeSettings(0)); + UNIT_ASSERT_VALUES_EQUAL_C(CountValueType(arr, EValueType::Double), 1, arr->DebugJson().GetStringRobust()); + UNIT_ASSERT_VALUES_EQUAL(PrintBinaryJsons(arr->GetChunkedArray()), PrintBinaryJsons(BuildSubColumns(docs, OffSettings())->GetChunkedArray())); + } + + // Double and Bool are never dictionary-encoded, even at low cardinality: the compact native arrays + // make dictionary encoding excessive, so they stay plain native columns. + Y_UNIT_TEST(DoubleStaysPlainAtLowCardinality) { + std::vector<TString> docs; + for (ui32 i = 0; i < 40; ++i) { + docs.push_back(TStringBuilder() << R"({"n":)" << (i % 2 ? "1.5" : "2.5") << "}"); + } + AssertPlainNativeColumn(docs, EValueType::Double); + } + + Y_UNIT_TEST(BoolStaysPlainAtLowCardinality) { + std::vector<TString> docs; + for (ui32 i = 0; i < 40; ++i) { + docs.push_back(TStringBuilder() << R"({"b":)" << (i % 2 ? "true" : "false") << "}"); + } + AssertPlainNativeColumn(docs, EValueType::Bool); + } + + Y_UNIT_TEST(BoolRoundTrip) { + const std::vector<TString> docs = { R"({"b":true})", R"({"b":false})", R"({"b":true})" }; + auto arr = BuildSubColumns(docs, NativeSettings(0)); + UNIT_ASSERT_VALUES_EQUAL_C(CountValueType(arr, EValueType::Bool), 1, arr->DebugJson().GetStringRobust()); + auto restored = SerializeRoundTrip(arr, NativeSettings(0)); + UNIT_ASSERT_VALUES_EQUAL(PrintBinaryJsons(restored->GetChunkedArray()), PrintBinaryJsons(BuildSubColumns(docs, OffSettings())->GetChunkedArray())); + } + + // A container value (array/object) is not a scalar, so the key stays BinaryJson. + Y_UNIT_TEST(ContainerStaysBinaryJson) { + const std::vector<TString> docs = { R"({"a":[1,2]})", R"({"a":[3]})" }; + auto arr = BuildSubColumns(docs, NativeSettings(0)); + UNIT_ASSERT_VALUES_EQUAL(CountValueType(arr, EValueType::Double), 0); + UNIT_ASSERT_VALUES_EQUAL(CountValueType(arr, EValueType::String), 0); + UNIT_ASSERT_VALUES_EQUAL(CountValueType(arr, EValueType::BinaryJson), 1); + UNIT_ASSERT_VALUES_EQUAL(PrintBinaryJsons(arr->GetChunkedArray()), PrintBinaryJsons(BuildSubColumns(docs, OffSettings())->GetChunkedArray())); + } + + // With the setting off (the default), the same all-string key stays BinaryJson. + Y_UNIT_TEST(OffByDefault) { + const std::vector<TString> docs = { R"({"a":"x"})", R"({"a":"yy"})" }; + auto arr = BuildSubColumns(docs, OffSettings()); + UNIT_ASSERT_VALUES_EQUAL(CountValueType(arr, EValueType::String), 0); + UNIT_ASSERT_VALUES_EQUAL(PrintBinaryJsons(arr->GetChunkedArray()), R"([[{"a":"x"},{"a":"yy"}]])"); + } + + // A full serialize -> deserialize round-trip (exercising the 5-column stats blob) preserves the + // native String column and reconstructs identical documents. + Y_UNIT_TEST(SerializeRoundTripPreservesNative) { + const std::vector<TString> docs = { R"({"a":"x","n":1})", R"({"a":"yy","n":2})" }; + auto arr = BuildSubColumns(docs, NativeSettings(0)); + auto restored = SerializeRoundTrip(arr, NativeSettings(0)); + UNIT_ASSERT_VALUES_EQUAL(PrintBinaryJsons(restored->GetChunkedArray()), PrintBinaryJsons(arr->GetChunkedArray())); + UNIT_ASSERT_VALUES_EQUAL_C(CountValueType(restored, EValueType::String), 1, restored->DebugJson().GetStringRobust()); + } + + // A low-cardinality String key composes with dictionary encoding. + Y_UNIT_TEST(StringComposesWithDictionary) { + std::vector<TString> docs; + for (ui32 i = 0; i < 40; ++i) { + docs.push_back(TStringBuilder() << R"({"a":")" << (i % 2 ? "xxxx" : "yyyy") << R"("})"); + } + AssertDictColumn(docs, EValueType::String); + } + + Y_UNIT_TEST(MixedTypeBinaryJsonFallback) { + auto arr = BuildSubColumns({ R"({"a":"x"})", R"({"a":7})" }, NativeSettings(0)); + UNIT_ASSERT_VALUES_EQUAL(CountValueType(arr, EValueType::BinaryJson), 1); + UNIT_ASSERT_VALUES_EQUAL(PrintBinaryJsons(arr->GetChunkedArray()), R"([[{"a":"x"},{"a":7}]])"); + } + + Y_UNIT_TEST(NullBinaryJsonFallback) { + auto arr = BuildSubColumns({ R"({"a":"x"})", R"({"a":null})" }, NativeSettings(0)); + UNIT_ASSERT_VALUES_EQUAL(CountValueType(arr, EValueType::BinaryJson), 1); + UNIT_ASSERT_VALUES_EQUAL(PrintBinaryJsons(arr->GetChunkedArray()), R"([[{"a":"x"},{"a":null}]])"); + } + + Y_UNIT_TEST(SpecialCharsRoundTrip) { + const std::vector<TString> docs = { + R"({"a":"he said \"hi\""})", + R"({"a":"tab\tend"})", + R"({"a":"юникод"})", + }; + auto arr = BuildSubColumns(docs, NativeSettings(0)); + UNIT_ASSERT_VALUES_EQUAL_C(CountValueType(arr, EValueType::String), 1, arr->DebugJson().GetStringRobust()); + UNIT_ASSERT_VALUES_EQUAL(PrintBinaryJsons(arr->GetChunkedArray()), PrintBinaryJsons(BuildSubColumns(docs, OffSettings())->GetChunkedArray())); + } + + // Current compaction always restores BinaryJson, may be dropped once that is fixed + Y_UNIT_TEST(OrderedIteratorNormalizesNativeToBinaryJson) { + auto arr = BuildSubColumns({ R"({"s":"x","n":3.5,"b":true})" }, NativeSettings(0)); + auto it = arr->BuildOrderedIterator(); + std::vector<TString> jsons; + it->ReadRecord( + 0, [](ui32) {}, + [&](ui32, const NBinaryJson::TBinaryJson& json, bool) { + UNIT_ASSERT_C(NBinaryJson::IsValidBinaryJson(TStringBuf(json.data(), json.size())), "ordered value is not BinaryJson"); + jsons.push_back(TString(NBinaryJson::SerializeToJson(json))); + }, + []() {}); + std::sort(jsons.begin(), jsons.end()); + UNIT_ASSERT_VALUES_EQUAL(jsons.size(), 3u); + UNIT_ASSERT_VALUES_EQUAL(jsons[0], "\"x\""); + UNIT_ASSERT_VALUES_EQUAL(jsons[1], "3.5"); + UNIT_ASSERT_VALUES_EQUAL(jsons[2], "true"); + } +}; 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 40f53f35cc7..ae2fafda803 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 @@ -4,14 +4,22 @@ #include <ydb/core/formats/arrow/accessor/sub_columns/constructor.h> #include <ydb/core/formats/arrow/accessor/sub_columns/data_extractor.h> #include <ydb/core/formats/arrow/accessor/sub_columns/json_value_path.h> +#include <ydb/core/formats/arrow/arrow_helpers.h> #include <ydb/core/formats/arrow/serializer/abstract.h> +#include "ut_helpers.h" + +#include <contrib/libs/apache/arrow/cpp/src/arrow/array/builder_binary.h> +#include <contrib/libs/apache/arrow/cpp/src/arrow/array/builder_primitive.h> + #include <library/cpp/testing/unittest/registar.h> #include <yql/essentials/types/binary_json/read.h> #include <yql/essentials/types/binary_json/write.h> #include <regex> +using NKikimr::NArrow::NAccessor::NSubColumns::NTesting::PrintBinaryJsons; + Y_UNIT_TEST_SUITE(SubColumnsArrayAccessor) { using namespace NKikimr::NArrow::NAccessor; using namespace NKikimr::NArrow; @@ -21,29 +29,6 @@ Y_UNIT_TEST_SUITE(SubColumnsArrayAccessor) { return std::regex_replace(str, std::regex(" |\\n"), ""); } - TString PrintBinaryJsons(const std::shared_ptr<arrow::ChunkedArray>& array) { - TStringBuilder sb; - sb << "["; - for (auto&& i : array->chunks()) { - sb << "["; - AFL_VERIFY(i->type()->id() == arrow::binary()->id()); - auto views = std::static_pointer_cast<arrow::BinaryArray>(i); - for (ui32 r = 0; r < views->length(); ++r) { - if (views->IsNull(r)) { - sb << "null"; - } else { - sb << NBinaryJson::SerializeToJson(TStringBuf(views->GetView(r).data(), views->GetView(r).size())); - } - if (r + 1 != views->length()) { - sb << ","; - } - } - sb << "]"; - } - sb << "]"; - return sb; - } - Y_UNIT_TEST(EmptyOthers){ auto arrEmpty = NSubColumns::TOthersData::BuildEmpty(); auto arrSliceEmpty = arrEmpty.Slice(0, 1000, NSubColumns::TSettings()); @@ -97,31 +82,31 @@ Y_UNIT_TEST_SUITE(SubColumnsArrayAccessor) { { auto arrSlice = arrData->ISlice(0, 0); UNIT_ASSERT_VALUES_EQUAL(PrintBinaryJsons(arrSlice->GetChunkedArray()), R"([])"); - UNIT_ASSERT_VALUES_EQUAL(arrSlice->DebugJson()["internal"]["columns_data"]["stats"].GetStringRobust(), R"({"accessor":[],"size":[],"key_names":[],"records":[]})"); - UNIT_ASSERT_VALUES_EQUAL(arrSlice->DebugJson()["internal"]["others_data"]["stats"].GetStringRobust(), R"({"accessor":[],"size":[],"key_names":[],"records":[]})"); + UNIT_ASSERT_VALUES_EQUAL(arrSlice->DebugJson()["internal"]["columns_data"]["stats"].GetStringRobust(), R"({"accessor":[],"value_type":[],"size":[],"key_names":[],"records":[]})"); + UNIT_ASSERT_VALUES_EQUAL(arrSlice->DebugJson()["internal"]["others_data"]["stats"].GetStringRobust(), R"({"accessor":[],"value_type":[],"size":[],"key_names":[],"records":[]})"); } { auto arrSlice = arrData->ISlice(0, 2); UNIT_ASSERT_VALUES_EQUAL(PrintBinaryJsons(arrSlice->GetChunkedArray()), R"([[{"a":1,"b":1,"c":"1111"},null]])"); if (colsCount == 1) { - UNIT_ASSERT_VALUES_EQUAL(arrSlice->DebugJson()["internal"]["columns_data"]["stats"].GetStringRobust(), R"({"accessor":[1],"size":[34],"key_names":["\"c\""],"records":[1]})"); - UNIT_ASSERT_VALUES_EQUAL(arrSlice->DebugJson()["internal"]["others_data"]["stats"].GetStringRobust(), R"({"accessor":[1,1],"size":[24,24],"key_names":["\"a\"","\"b\""],"records":[1,1]})"); + UNIT_ASSERT_VALUES_EQUAL(arrSlice->DebugJson()["internal"]["columns_data"]["stats"].GetStringRobust(), R"({"accessor":[1],"value_type":[0],"size":[34],"key_names":["\"c\""],"records":[1]})"); + UNIT_ASSERT_VALUES_EQUAL(arrSlice->DebugJson()["internal"]["others_data"]["stats"].GetStringRobust(), R"({"accessor":[1,1],"value_type":[0,0],"size":[24,24],"key_names":["\"a\"","\"b\""],"records":[1,1]})"); } } { auto arrSlice = arrData->ISlice(0, 3); UNIT_ASSERT_VALUES_EQUAL(PrintBinaryJsons(arrSlice->GetChunkedArray()), R"([[{"a":1,"b":1,"c":"1111"},null,{"a1":2,"b":2,"c":"2222"}]])"); if (colsCount == 1) { - UNIT_ASSERT_VALUES_EQUAL(arrSlice->DebugJson()["internal"]["columns_data"]["stats"].GetStringRobust(), R"({"accessor":[1],"size":[63],"key_names":["\"c\""],"records":[2]})"); - UNIT_ASSERT_VALUES_EQUAL(arrSlice->DebugJson()["internal"]["others_data"]["stats"].GetStringRobust(), R"({"accessor":[1,1,1],"size":[24,24,48],"key_names":["\"a\"","\"a1\"","\"b\""],"records":[1,1,2]})"); + UNIT_ASSERT_VALUES_EQUAL(arrSlice->DebugJson()["internal"]["columns_data"]["stats"].GetStringRobust(), R"({"accessor":[1],"value_type":[0],"size":[63],"key_names":["\"c\""],"records":[2]})"); + UNIT_ASSERT_VALUES_EQUAL(arrSlice->DebugJson()["internal"]["others_data"]["stats"].GetStringRobust(), R"({"accessor":[1,1,1],"value_type":[0,0,0],"size":[24,24,48],"key_names":["\"a\"","\"a1\"","\"b\""],"records":[1,1,2]})"); } } { auto arrSlice = arrData->ISlice(3, 3); UNIT_ASSERT_VALUES_EQUAL(PrintBinaryJsons(arrSlice->GetChunkedArray()), R"([[{"a":3,"b":3,"c":"3333"},null,{"a":5,"b1":5}]])"); if (colsCount == 1) { - UNIT_ASSERT_VALUES_EQUAL(arrSlice->DebugJson()["internal"]["columns_data"]["stats"].GetStringRobust(), R"({"accessor":[1],"size":[38],"key_names":["\"c\""],"records":[1]})"); - UNIT_ASSERT_VALUES_EQUAL(arrSlice->DebugJson()["internal"]["others_data"]["stats"].GetStringRobust(), R"({"accessor":[1,1,1],"size":[48,24,24],"key_names":["\"a\"","\"b\"","\"b1\""],"records":[2,1,1]})"); + UNIT_ASSERT_VALUES_EQUAL(arrSlice->DebugJson()["internal"]["columns_data"]["stats"].GetStringRobust(), R"({"accessor":[1],"value_type":[0],"size":[38],"key_names":["\"c\""],"records":[1]})"); + UNIT_ASSERT_VALUES_EQUAL(arrSlice->DebugJson()["internal"]["others_data"]["stats"].GetStringRobust(), R"({"accessor":[1,1,1],"value_type":[0,0,0],"size":[48,24,24],"key_names":["\"a\"","\"b\"","\"b1\""],"records":[2,1,1]})"); } } } @@ -557,3 +542,84 @@ Y_UNIT_TEST_SUITE(SubColumnsArrayAccessor) { } } }; + +Y_UNIT_TEST_SUITE(SubColumnsDictStats) { + using namespace NKikimr; + using namespace NKikimr::NArrow; + using namespace NKikimr::NArrow::NAccessor; + using namespace NKikimr::NArrow::NAccessor::NSubColumns; + + struct TStatsRow { + TString Name; + ui32 Records; + ui32 Size; + IChunkedArray::EType Accessor; + EValueType ValueType; + }; + + void AssertStatsMatch(const TDictStats& stats, const std::vector<TStatsRow>& expected) { + UNIT_ASSERT_VALUES_EQUAL(stats.GetColumnsCount(), expected.size()); + for (ui32 i = 0; i < expected.size(); ++i) { + UNIT_ASSERT_VALUES_EQUAL(stats.GetColumnNameString(i), expected[i].Name); + UNIT_ASSERT_VALUES_EQUAL(stats.GetColumnRecordsCount(i), expected[i].Records); + UNIT_ASSERT_VALUES_EQUAL(stats.GetColumnSize(i), expected[i].Size); + UNIT_ASSERT_VALUES_EQUAL((ui32)stats.GetAccessorType(i), (ui32)expected[i].Accessor); + UNIT_ASSERT_VALUES_EQUAL((ui32)stats.GetValueType(i), (ui32)expected[i].ValueType); + } + } + + Y_UNIT_TEST(ValueTypeCodesArePersistent) { + UNIT_ASSERT_VALUES_EQUAL((ui32)EValueType::BinaryJson, 0u); + UNIT_ASSERT_VALUES_EQUAL((ui32)EValueType::Double, 1u); + UNIT_ASSERT_VALUES_EQUAL((ui32)EValueType::Bool, 2u); + UNIT_ASSERT_VALUES_EQUAL((ui32)EValueType::String, 3u); + } + + // The current format (5-column, with value_type) round-trips through serialization, preserving every field. + Y_UNIT_TEST(NewFormatRoundTrip) { + const std::vector<TStatsRow> rows = { + { "a", 3, 30, IChunkedArray::EType::Array, EValueType::String }, + { "b", 5, 40, IChunkedArray::EType::Dictionary, EValueType::BinaryJson }, + { "c", 2, 16, IChunkedArray::EType::SparsedArray, EValueType::Double }, + { "d", 1, 8, IChunkedArray::EType::Array, EValueType::Bool }, + }; + auto builder = TDictStats::MakeBuilder(); + for (const auto& r : rows) { + builder.Add(r.Name, r.Records, r.Size, r.Accessor, r.ValueType); + } + auto stats = builder.Finish(); + auto restored = TDictStats::DeserializeFromBlob(stats.SerializeAsString(nullptr)); + AssertStatsMatch(restored, rows); + } + + Y_UNIT_TEST(LegacyFourColumnFormatDeserializes) { + const std::vector<TStatsRow> rows = { + { "a", 3, 30, IChunkedArray::EType::Array, EValueType::BinaryJson }, + { "b", 5, 40, IChunkedArray::EType::Dictionary, EValueType::BinaryJson }, + { "c", 2, 16, IChunkedArray::EType::SparsedArray, EValueType::BinaryJson }, + }; + auto legacySchema = std::make_shared<arrow::Schema>(arrow::FieldVector{ + std::make_shared<arrow::Field>("name", arrow::binary()), std::make_shared<arrow::Field>("count", arrow::uint32()), + std::make_shared<arrow::Field>("size", arrow::uint32()), std::make_shared<arrow::Field>("accessor_type", arrow::uint8()) }); + + arrow::BinaryBuilder names; + arrow::UInt32Builder count; + arrow::UInt32Builder size; + arrow::UInt8Builder acc; + for (const auto& r : rows) { + UNIT_ASSERT(names.Append(r.Name.data(), r.Name.size()).ok()); + UNIT_ASSERT(count.Append(r.Records).ok()); + UNIT_ASSERT(size.Append(r.Size).ok()); + UNIT_ASSERT(acc.Append((ui8)r.Accessor).ok()); + } + std::shared_ptr<arrow::Array> namesArr, countArr, sizeArr, accArr; + UNIT_ASSERT(names.Finish(&namesArr).ok()); + UNIT_ASSERT(count.Finish(&countArr).ok()); + UNIT_ASSERT(size.Finish(&sizeArr).ok()); + UNIT_ASSERT(acc.Finish(&accArr).ok()); + auto legacyBatch = arrow::RecordBatch::Make(legacySchema, rows.size(), { namesArr, countArr, sizeArr, accArr }); + + auto restored = TDictStats::DeserializeFromBlob(NArrow::SerializeBatchNoCompression(legacyBatch)); + AssertStatsMatch(restored, rows); + } +} 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 d3124eb9463..45aa769c406 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/ut/ya.make +++ b/ydb/core/formats/arrow/accessor/sub_columns/ut/ya.make @@ -11,6 +11,7 @@ PEERDIR( SRCS( ut_sub_columns.cpp + ut_native_scalars.cpp ut_dictionary.cpp ) diff --git a/ydb/core/formats/arrow/accessor/sub_columns/ya.make b/ydb/core/formats/arrow/accessor/sub_columns/ya.make index 11f68ba9128..afc29f20367 100644 --- a/ydb/core/formats/arrow/accessor/sub_columns/ya.make +++ b/ydb/core/formats/arrow/accessor/sub_columns/ya.make @@ -27,6 +27,7 @@ SRCS( json_value_path.cpp accessor.cpp direct_builder.cpp + types.cpp settings.cpp stats.cpp others_storage.cpp diff --git a/ydb/core/kqp/ut/olap/types/json_ut.cpp b/ydb/core/kqp/ut/olap/types/json_ut.cpp index 6b7abdf8977..1efdabfc900 100644 --- a/ydb/core/kqp/ut/olap/types/json_ut.cpp +++ b/ydb/core/kqp/ut/olap/types/json_ut.cpp @@ -8,6 +8,7 @@ #include <ydb/core/kqp/ut/olap/helpers/writer.h> #include <ydb/core/base/tablet_pipecache.h> +#include <ydb/core/formats/arrow/accessor/sub_columns/types.h> #include <ydb/core/formats/arrow/serializer/native.h> #include <ydb/core/kqp/ut/common/columnshard.h> #include <ydb/core/tx/columnshard/engines/reader/common_reader/iterator/source.h> @@ -1573,4 +1574,109 @@ Y_UNIT_TEST_SUITE(KqpOlapJson) { } +namespace { + +using NArrow::NAccessor::NSubColumns::EValueType; + +constexpr const char* CreateColumnTableDdl = R"(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);)"; + +// STOP_COMPACTION + create the JsonDocument table + pin the SIMPLE scan reader, then turn Col2 into a +// SUB_COLUMNS column that separates every key (OTHERS_ALLOWED_FRACTION=0) with native scalar storage on. +TString NativeTableSetup() { + TStringBuilder script; + script << R"( + STOP_COMPACTION + ------ + SCHEMA: + )" << CreateColumnTableDdl << R"( + ------ + 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`, + `OTHERS_ALLOWED_FRACTION`=`0`, `ENABLE_NATIVE_COLUMNS`=`true`) + ------)"; + return script; +} + +// primary_index_stats assertion that every separated Col2 sub-column in every chunk was stored with the +// given native `expectedValueType` (mirrors dictionary_ut's AccessorTypeCheck, but on $.columns.value_type; +// the value_type numeric code is what ChunkDetails serializes). 0 (BinaryJson) is the non-native default. +TString NativeValueTypeCheck(const EValueType expectedValueType) { + return Sprintf(R"SQL(READ: $All = SELECT COUNT(*) AS cnt FROM `/Root/ColumnTable/.sys/primary_index_stats` + WHERE Activity == 1 AND EntityName = 'Col2'; + $Ok = SELECT SUM(CASE + WHEN JSON_EXISTS(CAST(ChunkDetails AS JsonDocument), "$.columns.value_type[0]") AND + NOT JSON_EXISTS(CAST(ChunkDetails AS JsonDocument), "$.columns.value_type[*] ? (@ != %u)") + THEN 1 ELSE 0 END) AS ok + FROM `/Root/ColumnTable/.sys/primary_index_stats` + WHERE Activity == 1 AND EntityName = 'Col2'; + SELECT ($All > 0u) AND ($All == $Ok); + EXPECTED: [[[%%true]]])SQL", + (ui32)expectedValueType); +} + +} // namespace + +Y_UNIT_TEST_SUITE(KqpOlapJsonNativeScalars) { + + Y_UNIT_TEST(Double) { + const TString script = TStringBuilder() << NativeTableSetup() << R"( + DATA: + REPLACE INTO `/Root/ColumnTable` (Col1, Col2) VALUES (1u, JsonDocument('{"a" : 1.5}')), (2u, JsonDocument('{"a" : 2.5}')), + (3u, JsonDocument('{"a" : 3.5}')) + ------ + READ: SELECT * FROM `/Root/ColumnTable` ORDER BY Col1; + EXPECTED: [[1u;["{\"a\":1.5}"]];[2u;["{\"a\":2.5}"]];[3u;["{\"a\":3.5}"]]] + ------ + )" << NativeValueTypeCheck(EValueType::Double); + Variator::ToExecutor(Variator::SingleScript(script)).Execute(); + } + + Y_UNIT_TEST(Bool) { + const TString script = TStringBuilder() << NativeTableSetup() << R"( + DATA: + REPLACE INTO `/Root/ColumnTable` (Col1, Col2) VALUES (1u, JsonDocument('{"a" : true}')), (2u, JsonDocument('{"a" : false}')), + (3u, JsonDocument('{"a" : true}')) + ------ + READ: SELECT * FROM `/Root/ColumnTable` ORDER BY Col1; + EXPECTED: [[1u;["{\"a\":true}"]];[2u;["{\"a\":false}"]];[3u;["{\"a\":true}"]]] + ------ + )" << NativeValueTypeCheck(EValueType::Bool); + Variator::ToExecutor(Variator::SingleScript(script)).Execute(); + } + + Y_UNIT_TEST(String) { + const TString script = TStringBuilder() << NativeTableSetup() << R"( + DATA: + REPLACE INTO `/Root/ColumnTable` (Col1, Col2) VALUES (1u, JsonDocument('{"a" : "x"}')), (2u, JsonDocument('{"a" : "yy"}')), + (3u, JsonDocument('{"a" : "zzz"}')) + ------ + READ: SELECT * FROM `/Root/ColumnTable` ORDER BY Col1; + EXPECTED: [[1u;["{\"a\":\"x\"}"]];[2u;["{\"a\":\"yy\"}"]];[3u;["{\"a\":\"zzz\"}"]]] + ------ + )" << NativeValueTypeCheck(EValueType::String); + Variator::ToExecutor(Variator::SingleScript(script)).Execute(); + } + + Y_UNIT_TEST(MixedDocument) { + const TString script = TStringBuilder() << NativeTableSetup() << R"( + DATA: + REPLACE INTO `/Root/ColumnTable` (Col1, Col2) VALUES (1u, JsonDocument('{"b" : true, "n" : 1.5, "s" : "x"}')), + (2u, JsonDocument('{"b" : false, "n" : 2.5, "s" : "yy"}')) + ------ + READ: SELECT * FROM `/Root/ColumnTable` ORDER BY Col1; + EXPECTED: [[1u;["{\"b\":true,\"n\":1.5,\"s\":\"x\"}"]];[2u;["{\"b\":false,\"n\":2.5,\"s\":\"yy\"}"]]] + )"; + Variator::ToExecutor(Variator::SingleScript(script)).Execute(); + } +} + } // namespace NKikimr::NKqp diff --git a/ydb/core/sys_view/show_create/create_table_formatter.cpp b/ydb/core/sys_view/show_create/create_table_formatter.cpp index 0b13b887a17..37bcdacf971 100644 --- a/ydb/core/sys_view/show_create/create_table_formatter.cpp +++ b/ydb/core/sys_view/show_create/create_table_formatter.cpp @@ -1839,13 +1839,20 @@ void TCreateTableFormatter::FormatAlterColumn(const TString& fullPath, const NKi EscapeValue(settings.GetOthersAllowedFraction(), paramsStr); del = ", "; } - if (settings.HasDictionaryUniqueFraction() && settings.GetDictionaryUniqueFraction()) { + if (settings.HasDictionaryUniqueFraction()) { paramsStr << del; EscapeName("DICTIONARY_UNIQUE_FRACTION", paramsStr); paramsStr << "="; EscapeValue(settings.GetDictionaryUniqueFraction(), paramsStr); del = ", "; } + if (settings.HasEnableNativeColumns()) { + paramsStr << del; + EscapeName("ENABLE_NATIVE_COLUMNS", paramsStr); + paramsStr << "="; + EscapeValue(settings.GetEnableNativeColumns(), paramsStr); + del = ", "; + } if (settings.HasDataExtractor()) { const auto& dataExtractor = settings.GetDataExtractor(); if (dataExtractor.HasClassName() && !dataExtractor.GetClassName().empty()) { diff --git a/ydb/core/sys_view/ut_show_create.cpp b/ydb/core/sys_view/ut_show_create.cpp index b4b04514101..9fe670e0697 100644 --- a/ydb/core/sys_view/ut_show_create.cpp +++ b/ydb/core/sys_view/ut_show_create.cpp @@ -7,6 +7,7 @@ #include <ydb/public/lib/ydb_cli/dump/util/query_utils.h> +#include <util/string/cast.h> #include <util/string/subst.h> namespace NKikimr { @@ -2046,7 +2047,7 @@ Y_UNIT_TEST(TablePartitionPolicyIndexTable) { ); } -Y_UNIT_TEST(TableColumnAlterColumn) { +void CheckAlterColumnShowCreate(double dictionaryUniqueFraction, bool enableNativeColumns) { TTestEnv env(1, 4, {.StoragePools = 3, .ShowCreateTable = true, .AlterObjectEnabled = true, .EnableSparsedColumns = true, .EnableOlapCompression = true, .EnableCsDictionaryEncoding = true}); env.GetServer().GetRuntime()->SetLogPriority(NKikimrServices::KQP_EXECUTER, NActors::NLog::PRI_DEBUG); @@ -2056,8 +2057,10 @@ Y_UNIT_TEST(TableColumnAlterColumn) { TShowCreateChecker checker(env); - checker.CheckShowCreateTable( - R"( + const TString fractionStr = ToString(dictionaryUniqueFraction); + const TString nativeColumnsStr = enableNativeColumns ? "true" : "false"; + + const TString query = Sprintf(R"( CREATE TABLE `/Root/test_show_create` ( Col1 Uint64 NOT NULL, Col2 JsonDocument, @@ -2068,15 +2071,16 @@ Y_UNIT_TEST(TableColumnAlterColumn) { ) PARTITION BY HASH(Col1) WITH (STORE = COLUMN, AUTO_PARTITIONING_MIN_PARTITIONS_COUNT = 2); - ALTER OBJECT `/Root/test_show_create` (TYPE TABLE) SET (ACTION=ALTER_COLUMN, NAME=Col2, `FORCE_SIMD_PARSING`=`true`, `DATA_ACCESSOR_CONSTRUCTOR.CLASS_NAME`=`SUB_COLUMNS`, `OTHERS_ALLOWED_FRACTION`=`0.5`, `DICTIONARY_UNIQUE_FRACTION`=`0.5`); + ALTER OBJECT `/Root/test_show_create` (TYPE TABLE) SET (ACTION=ALTER_COLUMN, NAME=Col2, `FORCE_SIMD_PARSING`=`true`, `DATA_ACCESSOR_CONSTRUCTOR.CLASS_NAME`=`SUB_COLUMNS`, `OTHERS_ALLOWED_FRACTION`=`0.5`, `DICTIONARY_UNIQUE_FRACTION`=`%s`, `ENABLE_NATIVE_COLUMNS`=`%s`); ALTER OBJECT `/Root/test_show_create` (TYPE TABLE) SET (ACTION=ALTER_COLUMN, NAME=Col3, `DEFAULT_VALUE`=`5`); ALTER TABLE `/Root/test_show_create` ALTER COLUMN Col2 SET COMPRESSION (algorithm=zstd, level=4); ALTER TABLE `/Root/test_show_create` ALTER COLUMN Col3 SET ENCODING (DICT); ALTER TABLE `/Root/test_show_create` ALTER COLUMN Col4 SET ENCODING (); ALTER TABLE `/Root/test_show_create` ALTER COLUMN Col4 SET COMPRESSION (); ALTER TABLE `/Root/test_show_create` ALTER COLUMN Col5 SET ENCODING (OFF); - )", "test_show_create", - R"( + )", fractionStr.c_str(), nativeColumnsStr.c_str()); + + const TString expected = Sprintf(R"( CREATE TABLE `test_show_create` ( `Col1` Uint64 NOT NULL, `Col2` JsonDocument COMPRESSION (algorithm = zstd, level = 4), @@ -2091,11 +2095,20 @@ Y_UNIT_TEST(TableColumnAlterColumn) { AUTO_PARTITIONING_MIN_PARTITIONS_COUNT = 2 ); - ALTER OBJECT `/Root/test_show_create` (TYPE TABLE) SET (ACTION = ALTER_COLUMN, NAME = Col2, `DATA_ACCESSOR_CONSTRUCTOR.CLASS_NAME` = `SUB_COLUMNS`, `SPARSED_DETECTOR_KFF` = `20`, `COLUMNS_LIMIT` = `1024`, `MEM_LIMIT_CHUNK` = `52428800`, `OTHERS_ALLOWED_FRACTION` = `0.5`, `DICTIONARY_UNIQUE_FRACTION` = `0.5`, `DATA_EXTRACTOR_CLASS_NAME` = `JSON_SCANNER`, `SCAN_FIRST_LEVEL_ONLY` = `false`, `FORCE_SIMD_PARSING` = `true`); + ALTER OBJECT `/Root/test_show_create` (TYPE TABLE) SET (ACTION = ALTER_COLUMN, NAME = Col2, `DATA_ACCESSOR_CONSTRUCTOR.CLASS_NAME` = `SUB_COLUMNS`, `SPARSED_DETECTOR_KFF` = `20`, `COLUMNS_LIMIT` = `1024`, `MEM_LIMIT_CHUNK` = `52428800`, `OTHERS_ALLOWED_FRACTION` = `0.5`, `DICTIONARY_UNIQUE_FRACTION` = `%s`, `ENABLE_NATIVE_COLUMNS` = `%s`, `DATA_EXTRACTOR_CLASS_NAME` = `JSON_SCANNER`, `SCAN_FIRST_LEVEL_ONLY` = `false`, `FORCE_SIMD_PARSING` = `true`); ALTER OBJECT `/Root/test_show_create` (TYPE TABLE) SET (ACTION = ALTER_COLUMN, NAME = Col3, `DEFAULT_VALUE` = `5`); - )" - ); + )", fractionStr.c_str(), nativeColumnsStr.c_str()); + + checker.CheckShowCreateTable(query.c_str(), "test_show_create", expected); +} + +Y_UNIT_TEST(TableColumnAlterColumn) { + CheckAlterColumnShowCreate(0.5, true); +} + +Y_UNIT_TEST(TableColumnAlterColumnExplicitlyUnsetValues) { + CheckAlterColumnShowCreate(0, false); } Y_UNIT_TEST(TableColumnUpsertOptions) { diff --git a/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/builder.cpp b/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/builder.cpp index 13208100b93..ca0f399e936 100644 --- a/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/builder.cpp +++ b/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/builder.cpp @@ -62,8 +62,9 @@ void TMergedBuilder::FlushData() { if (ColumnBuilders[idx].GetFilledRecordsCount()) { auto accessor = ColumnBuilders[idx].Finish(RecordIndex); accessor = MaybeDictionaryEncode(accessor, ColumnBuilders[idx].GetFilledRecordsCount()); + // For now compaction reencodes all columns into BinaryJson statsBuilder.Add(ResultColumnStats.GetColumnName(idx), ColumnBuilders[idx].GetFilledRecordsCount(), - ColumnBuilders[idx].GetFilledRecordsSize(), accessor->GetType()); + ColumnBuilders[idx].GetFilledRecordsSize(), accessor->GetType(), NArrow::NAccessor::NSubColumns::EValueType::BinaryJson); arrays.emplace_back(std::move(accessor)); } } diff --git a/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/logic.cpp b/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/logic.cpp index 821128a5a6e..9cbda31bac9 100644 --- a/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/logic.cpp +++ b/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/logic.cpp @@ -50,13 +50,14 @@ TColumnPortionResult TSubColumnsMerger::DoExecute(const TChunkMergeContext& cont const auto startRecord = [&](const ui32 /*sourceRecordIndex*/) { builder.StartRecord(); }; - const auto addKV = [&](const ui32 sourceKeyIndex, const std::string_view value, const bool isColumnKey) { + const auto addKV = [&](const ui32 sourceKeyIndex, const NBinaryJson::TBinaryJson& value, const bool isColumnKey) { auto commonKeyInfo = RemapKeyIndex.RemapIndex(sourceIdx, sourceKeyIndex, isColumnKey); + TStringBuf valueBuf(value.data(), value.size()); if (commonKeyInfo.GetIsColumnKey()) { - builder.AddColumnKV(commonKeyInfo.GetCommonKeyIndex(), value); + builder.AddColumnKV(commonKeyInfo.GetCommonKeyIndex(), valueBuf); columnStats.Add(value.size()); } else { - builder.AddOtherKV(commonKeyInfo.GetCommonKeyIndex(), value); + builder.AddOtherKV(commonKeyInfo.GetCommonKeyIndex(), valueBuf); otherStats.Add(value.size()); } }; diff --git a/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/remap.cpp b/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/remap.cpp index 2d06564f061..c1206d05bf0 100644 --- a/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/remap.cpp +++ b/ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/remap.cpp @@ -1,5 +1,7 @@ #include "remap.h" +#include <ydb/core/formats/arrow/accessor/sub_columns/types.h> + namespace NKikimr::NOlap::NCompaction::NSubColumns { TRemapColumns::TOthersData::TFinishContext TRemapColumns::BuildRemapInfo( @@ -17,7 +19,9 @@ TRemapColumns::TOthersData::TFinishContext TRemapColumns::BuildRemapInfo( } builder.Add(i.first, statsByKeyIndex[i.second].GetRecordsCount(), statsByKeyIndex[i.second].GetDataSize(), settings.IsSparsed(statsByKeyIndex[i.second].GetRecordsCount(), recordsCount) ? NArrow::NAccessor::IChunkedArray::EType::SparsedArray - : NArrow::NAccessor::IChunkedArray::EType::Array); + : NArrow::NAccessor::IChunkedArray::EType::Array, + // For now others always encode in BinaryJson + NArrow::NAccessor::NSubColumns::EValueType::BinaryJson); remap[i.second] = idx++; } return TOthersData::TFinishContext(builder.Finish(), remap); diff --git a/ydb/library/formats/arrow/protos/accessor.proto b/ydb/library/formats/arrow/protos/accessor.proto index e27bac75137..229c2f36174 100644 --- a/ydb/library/formats/arrow/protos/accessor.proto +++ b/ydb/library/formats/arrow/protos/accessor.proto @@ -37,10 +37,13 @@ message TRequestedConstructor { optional uint32 ChunkMemoryLimit = 3 [default = 50000000]; optional double OthersAllowedFraction = 4 [default = 0.05]; optional TDataExtractor DataExtractor = 5; - // Apply dictionary encoding to subcolumn with string storage (BinaryJson ones or with deduced string types in future) + // Apply dictionary encoding to subcolumn with string storage (BinaryJson ones or with deduced string type) // if the fraction of unique values (<unique count> / <total present count>) <= this value. // Value of `1` means always dictionary-encode; `0` (default) - never dictionary-encode. optional double DictionaryUniqueFraction = 6 [default = 0]; + // Store a separated sub-column whose present values are all scalars of one type in the corresponding + // native Arrow array instead of per-value BinaryJson blobs. `false` (default) - always BinaryJson. + optional bool EnableNativeColumns = 7 [default = false]; } optional TSettings Settings = 1; } @@ -73,6 +76,7 @@ message TConstructor { optional double OthersAllowedFraction = 4 [default = 0.05]; optional TDataExtractor DataExtractor = 5; optional double DictionaryUniqueFraction = 6 [default = 0]; + optional bool EnableNativeColumns = 7 [default = false]; } optional TSettings Settings = 1; } |
