summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorrisenberg <[email protected]>2026-07-11 10:22:48 +0300
committerGitHub <[email protected]>2026-07-11 10:22:48 +0300
commitebff6a5cc9608f384a445589efbaa55099729cf5 (patch)
treede1c7d575141156c1e23632c79aa244994a3f40b
parent13a3a42e8d6f4e7c7926ed27e620425b467f4483 (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.
-rw-r--r--ydb/core/formats/arrow/accessor/plain/accessor.h7
-rw-r--r--ydb/core/formats/arrow/accessor/sub_columns/accessor.cpp3
-rw-r--r--ydb/core/formats/arrow/accessor/sub_columns/columns_storage.cpp20
-rw-r--r--ydb/core/formats/arrow/accessor/sub_columns/columns_storage.h22
-rw-r--r--ydb/core/formats/arrow/accessor/sub_columns/direct_builder.cpp75
-rw-r--r--ydb/core/formats/arrow/accessor/sub_columns/direct_builder.h20
-rw-r--r--ydb/core/formats/arrow/accessor/sub_columns/header.cpp8
-rw-r--r--ydb/core/formats/arrow/accessor/sub_columns/iterators.cpp10
-rw-r--r--ydb/core/formats/arrow/accessor/sub_columns/iterators.h29
-rw-r--r--ydb/core/formats/arrow/accessor/sub_columns/others_storage.cpp4
-rw-r--r--ydb/core/formats/arrow/accessor/sub_columns/others_storage.h11
-rw-r--r--ydb/core/formats/arrow/accessor/sub_columns/request.cpp5
-rw-r--r--ydb/core/formats/arrow/accessor/sub_columns/settings.h41
-rw-r--r--ydb/core/formats/arrow/accessor/sub_columns/stats.cpp53
-rw-r--r--ydb/core/formats/arrow/accessor/sub_columns/stats.h51
-rw-r--r--ydb/core/formats/arrow/accessor/sub_columns/types.cpp134
-rw-r--r--ydb/core/formats/arrow/accessor/sub_columns/types.h47
-rw-r--r--ydb/core/formats/arrow/accessor/sub_columns/ut/ut_helpers.h38
-rw-r--r--ydb/core/formats/arrow/accessor/sub_columns/ut/ut_native_scalars.cpp210
-rw-r--r--ydb/core/formats/arrow/accessor/sub_columns/ut/ut_sub_columns.cpp128
-rw-r--r--ydb/core/formats/arrow/accessor/sub_columns/ut/ya.make1
-rw-r--r--ydb/core/formats/arrow/accessor/sub_columns/ya.make1
-rw-r--r--ydb/core/kqp/ut/olap/types/json_ut.cpp106
-rw-r--r--ydb/core/sys_view/show_create/create_table_formatter.cpp9
-rw-r--r--ydb/core/sys_view/ut_show_create.cpp31
-rw-r--r--ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/builder.cpp3
-rw-r--r--ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/logic.cpp7
-rw-r--r--ydb/core/tx/columnshard/engines/changes/compaction/sub_columns/remap.cpp6
-rw-r--r--ydb/library/formats/arrow/protos/accessor.proto6
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;
}