diff options
| author | Yuriy Kaminskiy <[email protected]> | 2026-07-13 13:01:58 +0300 |
|---|---|---|
| committer | GitHub <[email protected]> | 2026-07-13 13:01:58 +0300 |
| commit | 58eef0cf74d327029af94f4b613ea38117e464bf (patch) | |
| tree | 3850c3bddaa4a5972eb8db81f8ff41ce7d8d1d25 | |
| parent | 702bf720887d55613d1cb112ae1fd6e4e3b77e78 (diff) | |
shared reading: support structued nested json parsing (#40462)
Co-authored-by: Pisarenko Grigoriy <[email protected]>
16 files changed, 924 insertions, 67 deletions
diff --git a/ydb/core/fq/libs/row_dispatcher/common/row_dispatcher_settings.cpp b/ydb/core/fq/libs/row_dispatcher/common/row_dispatcher_settings.cpp index dd2fcdde46f..f00a8aa1f20 100644 --- a/ydb/core/fq/libs/row_dispatcher/common/row_dispatcher_settings.cpp +++ b/ydb/core/fq/libs/row_dispatcher/common/row_dispatcher_settings.cpp @@ -46,9 +46,11 @@ TRowDispatcherSettings::TRowDispatcherSettings(const NConfig::TRowDispatcherConf , ConsumerMode(config.GetWithoutConsumer() ? EConsumerMode::Without : EConsumerMode::Required) {} -TRowDispatcherSettings::TRowDispatcherSettings(const NKikimrConfig::TStreamingQueriesConfig::TExternalStorageConfig& config) +TRowDispatcherSettings::TRowDispatcherSettings(const NKikimrConfig::TStreamingQueriesConfig::TExternalStorageConfig& config, bool enableStructuredJsonParsing) : Coordinator(config) { + JsonParser.SetStructuredParsing(enableStructuredJsonParsing); + std::optional<size_t> userPoolSize; if (NKikimr::HasAppData() && NActors::TlsActivationContext && NActors::TlsActivationContext->ActorSystem()) { userPoolSize = NActors::TlsActivationContext->ActorSystem()->GetPoolThreadsCount(NKikimr::AppData()->UserPoolId); diff --git a/ydb/core/fq/libs/row_dispatcher/common/row_dispatcher_settings.h b/ydb/core/fq/libs/row_dispatcher/common/row_dispatcher_settings.h index ad15ad0bc55..7dfe8b4642f 100644 --- a/ydb/core/fq/libs/row_dispatcher/common/row_dispatcher_settings.h +++ b/ydb/core/fq/libs/row_dispatcher/common/row_dispatcher_settings.h @@ -33,6 +33,7 @@ public: YDB_ACCESSOR(ui64, BatchSizeBytes, 1_MB); YDB_ACCESSOR(TDuration, BatchCreationTimeout, TDuration::Seconds(1)); YDB_ACCESSOR(ui64, BufferCellCount, 1000'000); + YDB_ACCESSOR(bool, StructuredParsing, true); }; class TCompileServiceSettings { @@ -65,7 +66,7 @@ public: TRowDispatcherSettings() = default; TRowDispatcherSettings(const NConfig::TRowDispatcherConfig& config); - TRowDispatcherSettings(const NKikimrConfig::TStreamingQueriesConfig_TExternalStorageConfig& config); + TRowDispatcherSettings(const NKikimrConfig::TStreamingQueriesConfig_TExternalStorageConfig& config, bool enableStructuredJsonParsing = false); private: YDB_ACCESSOR_MUTABLE(TCoordinatorSettings, Coordinator, {}); diff --git a/ydb/core/fq/libs/row_dispatcher/format_handler/parsers/json_parser.cpp b/ydb/core/fq/libs/row_dispatcher/format_handler/parsers/json_parser.cpp index 477a9e7dc8c..888bd4da50b 100644 --- a/ydb/core/fq/libs/row_dispatcher/format_handler/parsers/json_parser.cpp +++ b/ydb/core/fq/libs/row_dispatcher/format_handler/parsers/json_parser.cpp @@ -13,11 +13,29 @@ #include <yql/essentials/minikql/mkql_node_cast.h> #include <yql/essentials/minikql/mkql_string_util.h> #include <yql/essentials/minikql/mkql_type_ops.h> +#include <yql/essentials/minikql/computation/mkql_computation_node_holders.h> + +#include <sstream> namespace NFq::NRowDispatcher { namespace { +TStringBuf GetTypeName(const NKikimr::NMiniKQL::TType* type) { + if (type->GetKind() == NKikimr::NMiniKQL::TTypeBase::EKind::Data) { + if (auto maybeDataSlot = AS_TYPE(NKikimr::NMiniKQL::TDataType, type)->GetDataSlot()) { + return NYql::NUdf::GetDataTypeInfo(*maybeDataSlot).Name; + } + } + return type->GetKindAsStr(); +} + +static std::string JsonTypeToString(const simdjson::builtin::ondemand::json_type& json_type) { + std::stringstream ss; + ss << json_type; + return std::move(ss).str(); +} + #define CHECK_JSON_ERROR(value) \ const simdjson::error_code error = value; \ if (Y_UNLIKELY(error)) \ @@ -85,7 +103,7 @@ public: TString TypeYson; public: - TStatus InitParser(const TString& name, const TString& typeYson, std::span<ui16> parsedRows, const NKikimr::NMiniKQL::TType* typeMkql, bool skipErrors) { + TStatus InitParser(const TString& name, const TString& typeYson, std::span<ui16> parsedRows, const NKikimr::NMiniKQL::TType* typeMkql, bool skipErrors, NKikimr::NMiniKQL::THolderFactory* holderFactory) { Name = name; TypeYson = typeYson; IsOptional = false; @@ -93,6 +111,8 @@ public: Status = TStatus::Success(); ParsedRowsCount = 0; ParsedRows = parsedRows; + TypeMkql = nullptr; + HolderFactory = holderFactory; return Status = ExtractDataSlot(typeMkql); } @@ -112,6 +132,211 @@ public: return Status; } + bool ParseNestedValue(simdjson::builtin::ondemand::value jsonValue, NYql::NUdf::TUnboxedValue& resultValue, TStatus& status, const NKikimr::NMiniKQL::TType* type, bool isOptional) const { + Y_ENSURE(HolderFactory); // should be already verified by ParseNestedType + simdjson::builtin::ondemand::json_type cellType; + CHECK_JSON_ERROR(jsonValue.type().get(cellType)) { + SetParsingError(error, jsonValue, "determine json value type", status); + return false; + } + + if (cellType == simdjson::builtin::ondemand::json_type::null) { + if (isOptional || type->GetKind() == NKikimr::NMiniKQL::TTypeBase::EKind::Optional) { + resultValue = NYql::NUdf::TUnboxedValuePod(); + return true; + } + status = TStatus::Fail(EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Found unexpected null value, expected non optional type " << GetTypeName(type)); + return false; + } + + switch (type->GetKind()) { + case NKikimr::NMiniKQL::TTypeBase::EKind::Data: { + auto maybeDataSlot = AS_TYPE(NKikimr::NMiniKQL::TDataType, type)->GetDataSlot(); + Y_ENSURE(maybeDataSlot); + auto dataSlot = *maybeDataSlot; + if (dataSlot != NYql::NUdf::EDataSlot::Json) { + return ParseDataType(std::move(jsonValue), resultValue, status, dataSlot, isOptional, NYql::NUdf::GetDataTypeInfo(dataSlot).Name); + } else { + return ParseJsonType(std::move(jsonValue), resultValue, status); + } + } + + case NKikimr::NMiniKQL::TTypeBase::EKind::Optional: { + Y_ENSURE(!isOptional); + return ParseNestedValue(std::move(jsonValue), resultValue, status, AS_TYPE(NKikimr::NMiniKQL::TOptionalType, type)->GetItemType(), true); + } + + case NKikimr::NMiniKQL::TTypeBase::EKind::List: { + if (cellType != simdjson::builtin::ondemand::json_type::array) { + status = TStatus::Fail(EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Failed to parse nested json value (List), expected array, but got " << JsonTypeToString(cellType)); + return false; + } + auto listType = AS_TYPE(NKikimr::NMiniKQL::TListType, type); + auto itemType = listType->GetItemType(); + auto listBuilder = HolderFactory->NewList(); + for (auto elt : jsonValue.get_array()) { + simdjson::builtin::ondemand::value eltValue; + CHECK_JSON_ERROR(elt.get(eltValue)) { + SetParsingError(error, jsonValue, "parse as array", status); + return false; + } + NYql::NUdf::TUnboxedValue value; + if (!ParseNestedValue(std::move(eltValue), value, status, itemType, false)) { + return false; + } + listBuilder->Add(std::move(value)); + } + resultValue = listBuilder->Build(); + break; + } + + case NKikimr::NMiniKQL::TTypeBase::EKind::Tuple: { + if (cellType != simdjson::builtin::ondemand::json_type::array) { + status = TStatus::Fail(EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Failed to parse nested json value (Tuple), expected array, but got " << JsonTypeToString(cellType)); + return false; + } + auto tupleType = AS_TYPE(NKikimr::NMiniKQL::TTupleType, type); + auto elementsCount = tupleType->GetElementsCount(); + NYql::NUdf::TUnboxedValue *resultValues; + resultValue = HolderFactory->CreateDirectArrayHolder(elementsCount, resultValues); + size_t idx = 0; + for (auto elt : jsonValue.get_array()) { + if (idx == elementsCount) { + break; + } + simdjson::builtin::ondemand::value eltValue; + CHECK_JSON_ERROR(elt.get(eltValue)) { + SetParsingError(error, jsonValue, "parse as array (tuple)", status); + return false; + } + if (!ParseNestedValue(std::move(eltValue), resultValues[idx], status, tupleType->GetElementType(idx), false)) { + return false; + } + ++idx; + } + for (; idx != elementsCount; ++idx) { + // tail must be optional + if (tupleType->GetElementType(idx)->GetKind() != NKikimr::NMiniKQL::TTypeBase::EKind::Optional) { + status = TStatus::Fail(EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Failed to parse nested json value (Tuple), short json array at index " << idx << ", expected non optional type " << GetTypeName(tupleType->GetElementType(idx))); + return false; + } + } + break; + } + + case NKikimr::NMiniKQL::TTypeBase::EKind::Struct: { + if (cellType != simdjson::builtin::ondemand::json_type::object) { + status = TStatus::Fail(EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Failed to parse nested json value (Struct), expected object, but got " << JsonTypeToString(cellType)); + return false; + } + auto structType = AS_TYPE(NKikimr::NMiniKQL::TStructType, type); + auto membersCount = structType->GetMembersCount(); + + auto it = StructMembers.find(structType); + Y_ENSURE(it != StructMembers.end()); + const auto& memberNames = it->second; + Y_ENSURE(memberNames.size() == membersCount); + + NYql::NUdf::TUnboxedValue *resultValues; + resultValue = HolderFactory->CreateDirectArrayHolder(membersCount, resultValues); + for (auto elt : jsonValue.get_object()) { + std::string_view name; + { + CHECK_JSON_ERROR(elt.escaped_key().get(name)) { + SetParsingError(error, jsonValue, "parse as object", status); + return false; + } + } + auto itName = memberNames.find(name); + if (itName == memberNames.end()) { + continue; + } + auto idx = itName->second; + + simdjson::builtin::ondemand::value eltValue; + CHECK_JSON_ERROR(elt.value().get(eltValue)) { + SetParsingError(error, jsonValue, "parse as object", status); + return false; + } + + if (!ParseNestedValue(std::move(eltValue), resultValues[idx], status, structType->GetMemberType(idx), false)) { + return false; + } + } + + for (ui32 idx = 0; idx != membersCount; ++idx) { + if (!resultValues[idx] && structType->GetMemberType(idx)->GetKind() != NKikimr::NMiniKQL::TTypeBase::EKind::Optional) { + status = TStatus::Fail(EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Failed to parse nested json value (Struct), expected non-optional field " << structType->GetMemberName(idx)); + return false; + } + } + break; + } + + case NKikimr::NMiniKQL::TTypeBase::EKind::Dict: { + if (cellType != simdjson::builtin::ondemand::json_type::object) { + status = TStatus::Fail(EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Failed to parse nested json value (Dict), expected object, but got " << JsonTypeToString(cellType)); + return false; + } + auto jsonObject = jsonValue.get_object(); + { + bool isEmpty; + CHECK_JSON_ERROR(jsonObject.is_empty().get(isEmpty)) { + SetParsingError(error, jsonValue, "parse as object", status); + return false; + } + if (isEmpty) { + resultValue = HolderFactory->GetEmptyContainerLazy(); + return true; + } + } + auto dictType = AS_TYPE(NKikimr::NMiniKQL::TDictType, type); + auto keyType = dictType->GetKeyType(); + Y_ENSURE(keyType->GetKind() == NKikimr::NMiniKQL::TTypeBase::EKind::Data); // should be already verified, not user-data-error + auto keyDataSlot = AS_TYPE(NKikimr::NMiniKQL::TDataType, keyType)->GetDataSlot(); + Y_ENSURE(keyDataSlot); // should be already verified, not user-data-error + auto payloadType = dictType->GetPayloadType(); + NKikimr::NMiniKQL::TKeyTypes keyTypes; + bool isTuple; + bool encoded; + bool useIHash; + GetDictionaryKeyTypes(keyType, keyTypes, isTuple, encoded, useIHash); + Y_ENSURE(!(isTuple || encoded || useIHash)); + status = TStatus::Success(); + resultValue = HolderFactory->CreateDirectHashedDictHolder( + [&, this](auto& map) { + for (auto elt : jsonObject) { + std::string_view name; + { + CHECK_JSON_ERROR(elt.escaped_key().get(name)) { + SetParsingError(error, jsonValue, "parse as object", status); + return; + } + } + simdjson::builtin::ondemand::value eltValue; + CHECK_JSON_ERROR(elt.value().get(eltValue)) { + SetParsingError(error, jsonValue, "parse as object", status); + return; + } + NYql::NUdf::TUnboxedValue payload; + if (!ParseNestedValue(std::move(eltValue), payload, status, payloadType, false)) { + return; + } + map.emplace(NKikimr::NMiniKQL::ValueFromString(*keyDataSlot, name), std::move(payload)); + } + }, + keyTypes, false, true, nullptr, nullptr, nullptr); + return status.IsSuccess(); + } + + default: + // should've been handled in ParseNestedType + status = TStatus::Fail(EStatusId::UNSUPPORTED, TStringBuilder() << "Unsupported type kind: " << type->GetKindAsStr()); + return false; + } + return true; + } + bool ParseJsonValue(ui64 offset, ui16 rowId, simdjson::builtin::ondemand::value jsonValue, NYql::NUdf::TUnboxedValue& resultValue) { if (Y_UNLIKELY(!SkipErrors && Status.IsFail())) { return false; @@ -122,15 +347,15 @@ public: } ParsedRows[ParsedRowsCount++] = rowId; bool success = false; - if (DataSlot != NYql::NUdf::EDataSlot::Json) { + if (TypeMkql) { + success = ParseNestedValue(std::move(jsonValue), resultValue, Status, TypeMkql, IsOptional); + } else if (DataSlot != NYql::NUdf::EDataSlot::Json) { success = ParseDataType(std::move(jsonValue), resultValue, Status); } else { success = ParseJsonType(std::move(jsonValue), resultValue, Status); } - - if (IsOptional && resultValue) { - resultValue = resultValue.MakeOptional(); - } + Y_DEBUG_ABORT_UNLESS(ValidateObjectRefs(resultValue)); + resultValue = LockObject(std::move(resultValue)); if (Y_UNLIKELY(!SkipErrors && Status.IsFail())) { Status.AddParentIssue(TStringBuilder() << "Failed to parse json string at offset " << offset << ", got parsing error for column '" << Name << "' with type " << TypeYson); @@ -161,7 +386,214 @@ public: return false; } + bool ValidateObjectRefs(const NYql::NUdf::TUnboxedValue& value) { + return TypeMkql ? ValidateObjectRefs(TypeMkql, value) : true; + } private: + + bool ValidateObjectRefs(const NKikimr::NMiniKQL::TType* type, const NYql::NUdf::TUnboxedValue& value, i32 expectedRefs = 1) { + if (value.RefCount() == -1) { // NULL/POD/Embedded String + return true; + } + if (type->IsOptional()) { + type = AS_TYPE(NKikimr::NMiniKQL::TOptionalType, type)->GetItemType(); + } + switch (type->GetKind()) { + case NKikimr::NMiniKQL::TTypeBase::EKind::Data: + if (value.RefCount() != expectedRefs) { + Y_DEBUG_ABORT("%s", (TStringBuilder() << value.RefCount() << "!=" << expectedRefs).c_str()); + return false; + } + return true; + + case NKikimr::NMiniKQL::TTypeBase::EKind::Struct: { + if (!value.IsBoxed()) { + Y_DEBUG_ABORT(); + return false; + } + + auto structType = AS_TYPE(NKikimr::NMiniKQL::TStructType, type); + auto membersCount = structType->GetMembersCount(); + if (value.RefCount() != expectedRefs) { + if (membersCount == 0) { + return true; + } + Y_DEBUG_ABORT("%s", (TStringBuilder() << value.RefCount() << "!=" << expectedRefs).c_str()); + return false; + } + auto elements = value.GetElements(); + for (ui32 i = 0; i < membersCount; ++i) { + if (!ValidateObjectRefs(structType->GetMemberType(i), elements[i])) { + Y_DEBUG_ABORT(); + return false; + } + } + break; + } + case NKikimr::NMiniKQL::TTypeBase::EKind::Tuple: { + if (!value.IsBoxed()) { + Y_DEBUG_ABORT(); + return false; + } + auto tupleType = AS_TYPE(NKikimr::NMiniKQL::TTupleType, type); + auto elementsCount = tupleType->GetElementsCount(); + if (value.RefCount() != expectedRefs) { + if (elementsCount == 0) { // special case: shared empty list + return true; + } + Y_DEBUG_ABORT("%s", (TStringBuilder() << value.RefCount() << "!=" << expectedRefs).c_str()); + return false; + } + auto elements = value.GetElements(); + for (ui32 i = 0; i < elementsCount; ++i) { + if (!ValidateObjectRefs(tupleType->GetElementType(i), elements[i])) { + Y_DEBUG_ABORT("%d", i); + return false; + } + } + break; + } + case NKikimr::NMiniKQL::TTypeBase::EKind::List: { + if (!value.IsBoxed()) { + Y_DEBUG_ABORT(); + return false; + } + if (value.RefCount() != expectedRefs) { + if (value.GetListLength() == 0) { // special case: shared empty list + return true; + } + Y_DEBUG_ABORT("%s", (TStringBuilder() << value.RefCount() << "!=" << expectedRefs).c_str()); + return false; + } + auto listType = AS_TYPE(NKikimr::NMiniKQL::TListType, type); + auto itemType = listType->GetItemType(); + auto listIterator = value.GetListIterator(); + + for(NYql::NUdf::TUnboxedValue itemValue; listIterator.Next(itemValue); ) { + if (!ValidateObjectRefs(itemType, itemValue, 2)) { // two refs expected: one from object, one from itemValue + Y_DEBUG_ABORT(); + return false; + } + } + break; + } + + case NKikimr::NMiniKQL::TTypeBase::EKind::Dict: { + if (!value.IsBoxed()) { + Y_DEBUG_ABORT(); + return false; + } + if (value.RefCount() != expectedRefs) { + if (value.GetDictLength() == 0) { // special case: shared empty dict + return true; + } + Y_DEBUG_ABORT("%s", (TStringBuilder() << value.RefCount() << "!=" << expectedRefs).c_str()); + return false; + } + auto dictType = AS_TYPE(NKikimr::NMiniKQL::TDictType, type); + auto keyType = dictType->GetKeyType(); + auto valueType = dictType->GetPayloadType(); + auto dictIterator = value.GetDictIterator(); + + for(NYql::NUdf::TUnboxedValue keyValue, payload; dictIterator.NextPair(keyValue, payload); ) { + if (!ValidateObjectRefs(keyType, keyValue, 2)) { // two refs expected: one from object, one from keyValue + Y_DEBUG_ABORT(); + return false; + } + if (!ValidateObjectRefs(valueType, payload, 2)) { // two refs expected: one from object, one from payload + Y_DEBUG_ABORT(); + return false; + } + } + break; + } + default: + Y_ABORT(); + } + return true; + } + + TStatus ParseNestedType(const NKikimr::NMiniKQL::TType* type, bool isOptional = false) { + if (!HolderFactory) { + return TStatus::Fail(EStatusId::UNSUPPORTED, "Structured Json Parsing is disabled. Please contact your system administrator to enable it"); + } + switch (type->GetKind()) { + case NKikimr::NMiniKQL::TTypeBase::EKind::Data: { + auto dataSlot = AS_TYPE(NKikimr::NMiniKQL::TDataType, type)->GetDataSlot(); + if (!dataSlot) { + return TStatus::Fail(EStatusId::UNSUPPORTED, TStringBuilder() << "DataType has no DataSlot"); + } + return TStatus::Success(); + } + + case NKikimr::NMiniKQL::TTypeBase::EKind::Optional: { + if (isOptional) { + return TStatus::Fail(EStatusId::UNSUPPORTED, TStringBuilder() << "Nested optionals is not supported as input type"); + } + return ParseNestedType(AS_TYPE(NKikimr::NMiniKQL::TOptionalType, type)->GetItemType(), true); + } + + case NKikimr::NMiniKQL::TTypeBase::EKind::Struct: { + auto structType = AS_TYPE(NKikimr::NMiniKQL::TStructType, type); + auto membersCount = structType->GetMembersCount(); + + auto [it, inserted] = StructMembers.try_emplace(structType); + auto& memberNames = it->second; + if (!inserted) { + Y_ENSURE(membersCount == memberNames.size()); + break; + } + + for (ui32 idx = 0; idx != membersCount; ++idx) { + auto success = ParseNestedType(structType->GetMemberType(idx)); + if (!success) { + return success; + } + + auto [_, nameInserted] = memberNames.emplace(structType->GetMemberName(idx), idx); + Y_ENSURE(nameInserted); + } + break; + } + + case NKikimr::NMiniKQL::TTypeBase::EKind::Dict: { + auto dictType = AS_TYPE(NKikimr::NMiniKQL::TDictType, type); + auto keyType = dictType->GetKeyType(); + if (keyType->GetKind() != NKikimr::NMiniKQL::TTypeBase::EKind::Data) { + return TStatus::Fail(EStatusId::UNSUPPORTED, TStringBuilder() << "Dict key type: expected either String, or Utf8, but got " << keyType->GetKindAsStr()); + } + auto keyTypeSlot = AS_TYPE(NKikimr::NMiniKQL::TDataType, keyType)->GetDataSlot(); + if (!IsIn({NYql::NUdf::EDataSlot::String, NYql::NUdf::EDataSlot::Utf8}, keyTypeSlot)) { + return TStatus::Fail(EStatusId::UNSUPPORTED, TStringBuilder() << "Dict key type: expected either String, or Utf8, but got " << (keyTypeSlot ? NYql::NUdf::GetDataTypeInfo(*keyTypeSlot).Name : "<empty>")); + } + return ParseNestedType(dictType->GetPayloadType()); + } + + case NKikimr::NMiniKQL::TTypeBase::EKind::List: { + auto listType = AS_TYPE(NKikimr::NMiniKQL::TListType, type); + return ParseNestedType(listType->GetItemType()); + } + + case NKikimr::NMiniKQL::TTypeBase::EKind::Tuple: { + auto tupleType = AS_TYPE(NKikimr::NMiniKQL::TTupleType, type); + auto elementsCount = tupleType->GetElementsCount(); + + for (ui32 idx = 0; idx != elementsCount; ++idx) { + auto success = ParseNestedType(tupleType->GetElementType(idx)); + if (!success) { + return success; + } + } + break; + } + + default: { + return TStatus::Fail(EStatusId::UNSUPPORTED, TStringBuilder() << "Unsupported type kind: " << type->GetKindAsStr()); + } + } + return TStatus::Success(); + } + TStatus ExtractDataSlot(const NKikimr::NMiniKQL::TType* type) { switch (type->GetKind()) { case NKikimr::NMiniKQL::TTypeBase::EKind::Data: { @@ -181,6 +613,13 @@ private: IsOptional = true; return ExtractDataSlot(AS_TYPE(NKikimr::NMiniKQL::TOptionalType, type)->GetItemType()); } + case NKikimr::NMiniKQL::TTypeBase::EKind::Dict: + case NKikimr::NMiniKQL::TTypeBase::EKind::Tuple: + case NKikimr::NMiniKQL::TTypeBase::EKind::Struct: + case NKikimr::NMiniKQL::TTypeBase::EKind::List: { + TypeMkql = type; + return ParseNestedType(type); + } default: { return TStatus::Fail(EStatusId::UNSUPPORTED, TStringBuilder() << "Unsupported type kind: " << type->GetKindAsStr()); @@ -189,6 +628,10 @@ private: } Y_FORCE_INLINE bool ParseDataType(simdjson::builtin::ondemand::value jsonValue, NYql::NUdf::TUnboxedValue& resultValue, TStatus& status) const { + return ParseDataType(jsonValue, resultValue, status, DataSlot, IsOptional, DataTypeName); + } + + Y_FORCE_INLINE bool ParseDataType(simdjson::builtin::ondemand::value jsonValue, NYql::NUdf::TUnboxedValue& resultValue, TStatus& status, NYql::NUdf::EDataSlot dataSlot, bool isOptional, const TStringBuf dataTypeName) const { simdjson::builtin::ondemand::json_type cellType; CHECK_JSON_ERROR(jsonValue.type().get(cellType)) { SetParsingError(error, jsonValue, "determine json value type", status); @@ -197,7 +640,7 @@ private: switch (cellType) { case simdjson::builtin::ondemand::json_type::number: { - switch (DataSlot) { + switch (dataSlot) { case NYql::NUdf::EDataSlot::Int8: ParseJsonNumber<i8>(jsonValue.get_int64(), resultValue, status); break; @@ -232,11 +675,11 @@ private: break; default: - status = TStatus::Fail(EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Number value is not expected for data type " << DataTypeName); + status = TStatus::Fail(EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Number value is not expected for data type " << dataTypeName); break; } if (Y_UNLIKELY(status.IsFail())) { - status.AddParentIssue(TStringBuilder() << "Failed to parse data type " << DataTypeName << " from json number (raw: '" << TruncateString(jsonValue.raw_json_token()) << "')"); + status.AddParentIssue(TStringBuilder() << "Failed to parse data type " << dataTypeName << " from json number (raw: '" << TruncateString(jsonValue.raw_json_token()) << "')"); } return resultValue.HasValue(); } @@ -248,9 +691,9 @@ private: return false; } - resultValue = LockObject(NKikimr::NMiniKQL::ValueFromString(DataSlot, rawString)); + resultValue = NKikimr::NMiniKQL::ValueFromString(dataSlot, rawString); if (Y_UNLIKELY(!resultValue)) { - status = TStatus::Fail(EStatusId::BAD_REQUEST, TStringBuilder() << "Failed to parse data type " << DataTypeName << " from json string: '" << TruncateString(rawString) << "'"); + status = TStatus::Fail(EStatusId::BAD_REQUEST, TStringBuilder() << "Failed to parse data type " << dataTypeName << " from json string: '" << TruncateString(rawString) << "'"); } return resultValue.HasValue(); } @@ -262,13 +705,13 @@ private: SetParsingError(error, jsonValue, "extract json value", status); return false; } - status = TStatus::Fail(EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Found unexpected nested value (raw: '" << TruncateString(rawJson) << "'), expected data type " <<DataTypeName << ", please use Json type for nested values"); + status = TStatus::Fail(EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Found unexpected nested value (raw: '" << TruncateString(rawJson) << "'), expected data type " << dataTypeName << ", please use Json type for nested values"); return false; } case simdjson::builtin::ondemand::json_type::boolean: { - if (Y_UNLIKELY(DataSlot != NYql::NUdf::EDataSlot::Bool)) { - status = TStatus::Fail(EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Found unexpected bool value, expected data type " << DataTypeName); + if (Y_UNLIKELY(dataSlot != NYql::NUdf::EDataSlot::Bool)) { + status = TStatus::Fail(EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Found unexpected bool value, expected data type " << dataTypeName); return false; } @@ -283,8 +726,8 @@ private: } case simdjson::builtin::ondemand::json_type::null: { - if (Y_UNLIKELY(!IsOptional)) { - status = TStatus::Fail(EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Found unexpected null value, expected non optional data type " << DataTypeName); + if (Y_UNLIKELY(!isOptional)) { + status = TStatus::Fail(EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Found unexpected null value, expected non optional data type " << dataTypeName); return false; } @@ -297,7 +740,7 @@ private: SetParsingError(error, jsonValue, "extract json string", status); return false; } - status = TStatus::Fail(EStatusId::BAD_REQUEST, TStringBuilder() << "Failed to parse data type " << DataTypeName << " from json string: '" << TruncateString(rawString) << "'"); + status = TStatus::Fail(EStatusId::BAD_REQUEST, TStringBuilder() << "Failed to parse data type " << dataTypeName << " from json string: '" << TruncateString(rawString) << "'"); return false; } } @@ -315,7 +758,7 @@ private: return false; } - resultValue = LockObject(NKikimr::NMiniKQL::MakeString(rawJson)); + resultValue = NKikimr::NMiniKQL::MakeString(rawJson); return true; } @@ -350,6 +793,9 @@ private: } private: + NKikimr::NMiniKQL::THolderFactory* HolderFactory; + THashMap<const NKikimr::NMiniKQL::TStructType*, THashMap<std::string_view, size_t>> StructMembers; + const NKikimr::NMiniKQL::TType* TypeMkql = nullptr; // for complex types NYql::NUdf::EDataSlot DataSlot; TString DataTypeName; bool IsOptional = false; @@ -405,7 +851,7 @@ public: return TStatus(typeStatus).AddParentIssue(TStringBuilder() << "Failed to parse column '" << name << "' type " << typeYson); } - if (auto status = Columns[i].InitParser(name, typeYson, parsedRowsIdxSpan.subspan(i * MaxNumberRows, MaxNumberRows), typeStatus.DetachResult(), Config.SkipErrors); status.IsFail()) { + if (auto status = Columns[i].InitParser(name, typeYson, parsedRowsIdxSpan.subspan(i * MaxNumberRows, MaxNumberRows), typeStatus.DetachResult(), Config.SkipErrors, Config.StructuredParsing ? HolderFactory.get() : nullptr); status.IsFail()) { return status.AddParentIssue(TStringBuilder() << "Failed to create parser for column '" << name << "' with type " << typeYson); } if (!Columns[i].GetIsOptional()) { @@ -533,6 +979,7 @@ protected: const auto parsedRows = column.GetParsedRows(); const auto parsedRowsCount = column.GetParsedRowsCount(); for (ui16 rowId = 0; rowId < parsedRowsCount; ++rowId) { + Y_DEBUG_ABORT_UNLESS(column.ValidateObjectRefs(parsedColumn[parsedRows[rowId]])); ClearObject(parsedColumn[parsedRows[rowId]]); } @@ -715,6 +1162,7 @@ private: void ClearRowBuffer(ui16 outputRowId) { for (size_t columnId = 0; columnId < Columns.size(); ++columnId) { if (Columns[columnId].ClearParsedRow(outputRowId)) { + Y_DEBUG_ABORT_UNLESS(Columns[columnId].ValidateObjectRefs(ParsedValues[columnId][outputRowId])); ClearObject(ParsedValues[columnId][outputRowId]); } } @@ -783,7 +1231,8 @@ TJsonParserConfig CreateJsonParserConfig(const TRowDispatcherSettings::TJsonPars .BatchSize = parserConfig.GetBatchSizeBytes(), .LatencyLimit = parserConfig.GetBatchCreationTimeout(), .BufferCellCount = parserConfig.GetBufferCellCount(), - .SkipErrors = skipErrors + .SkipErrors = skipErrors, + .StructuredParsing = parserConfig.GetStructuredParsing(), }; } diff --git a/ydb/core/fq/libs/row_dispatcher/format_handler/parsers/json_parser.h b/ydb/core/fq/libs/row_dispatcher/format_handler/parsers/json_parser.h index a86b3466cc4..358c7b30fcf 100644 --- a/ydb/core/fq/libs/row_dispatcher/format_handler/parsers/json_parser.h +++ b/ydb/core/fq/libs/row_dispatcher/format_handler/parsers/json_parser.h @@ -18,6 +18,7 @@ struct TJsonParserConfig { TDuration LatencyLimit; ui64 BufferCellCount = 1000000; // (number rows) * (number columns) limit bool SkipErrors = false; + bool StructuredParsing = true; }; TValueStatus<ITopicParser::TPtr> CreateJsonParser(IParsedDataConsumer::TPtr consumer, const TJsonParserConfig& config, const TCountersDesc& counters); diff --git a/ydb/core/fq/libs/row_dispatcher/format_handler/parsers/parser_base.cpp b/ydb/core/fq/libs/row_dispatcher/format_handler/parsers/parser_base.cpp index 036a4989867..fe36169181a 100644 --- a/ydb/core/fq/libs/row_dispatcher/format_handler/parsers/parser_base.cpp +++ b/ydb/core/fq/libs/row_dispatcher/format_handler/parsers/parser_base.cpp @@ -14,10 +14,13 @@ TTypeParser::TTypeParser(const TSourceLocation& location, const NKikimr::NMiniKQ , FunctionRegistry(functionRegistry) , TypeEnv(std::make_unique<NKikimr::NMiniKQL::TTypeEnvironment>(Alloc)) , ProgramBuilder(std::make_unique<NKikimr::NMiniKQL::TProgramBuilder>(*TypeEnv, *FunctionRegistry)) + , MemInfo("SharedReadingParser") + , HolderFactory(std::make_unique<NKikimr::NMiniKQL::THolderFactory>(Alloc.Ref(), MemInfo, functionRegistry)) {} TTypeParser::~TTypeParser() { with_lock (Alloc) { + HolderFactory.reset(); TypeEnv.reset(); ProgramBuilder.reset(); } @@ -111,20 +114,23 @@ void TTopicParserBase::ParseBuffer() { //// Functions NYql::NUdf::TUnboxedValue LockObject(NYql::NUdf::TUnboxedValue&& value) { - // All UnboxedValue's with type Boxed or String should be locked - // because after parsing they will be used under another MKQL allocator in purecalc filters - - const i32 numberRefs = value.LockRef(); - - // -1 - value is embbeded or empty, otherwise value should have exactly one ref - Y_ENSURE(numberRefs == -1 || numberRefs == 1); - - return value; + // Object must be one of: + // 1) null (refs = -1) + // 2) embedded string/POD (refs = -1) + // 3) large string (refs = 1) + // 4a) dict as boxed: (refs = 1) + // 4b) (struct | tuple | list) as boxed -> direct array holder (refs = 1, except for special case: zero-length direct array holder points to special shared zero-length container) + Y_ABORT_UNLESS(value.RefCount() == -1 || value.RefCount() == 1 || (value.IsBoxed() && value.GetListLength() == 0)); + // Note that GetListLength will trigger ABORT if called on non-List (i.e. Dict), so there are no much point in turning it into YQL_ENSURE/Y_VALIDATE. + // It is debatable if this should be Y_DEBUG_ABORT_UNLESS (if anything else keeps reference on our object, it must be outside of our code and will result in freeing-with-wrong-allocator, which is UB anyway, and better die early than later). + // It is debatable if it is NOT sufficient check for nested structures (but that check would be expensive and thus belong to DEBUG checks) + return std::move(value); } void ClearObject(NYql::NUdf::TUnboxedValue& value) { - // Value should be unlocked with same number of refs - value.UnlockRef(1); + // Every other reference to value must be cleared (except for special case - zero-length list) + Y_ABORT_UNLESS(value.RefCount() == -1 || value.RefCount() == 1 || (value.IsBoxed() && value.GetListLength() == 0)); + // (again, using YQL_ENSURE/Y_VALIDATE pointless here) value.Clear(); } diff --git a/ydb/core/fq/libs/row_dispatcher/format_handler/parsers/parser_base.h b/ydb/core/fq/libs/row_dispatcher/format_handler/parsers/parser_base.h index b7de126b0e4..4c43b2e1e35 100644 --- a/ydb/core/fq/libs/row_dispatcher/format_handler/parsers/parser_base.h +++ b/ydb/core/fq/libs/row_dispatcher/format_handler/parsers/parser_base.h @@ -3,6 +3,7 @@ #include "parser_abstract.h" #include <yql/essentials/minikql/mkql_program_builder.h> +#include <yql/essentials/minikql/computation/mkql_computation_node_holders.h> namespace NFq::NRowDispatcher { @@ -18,6 +19,8 @@ protected: const NKikimr::NMiniKQL::IFunctionRegistry* FunctionRegistry; std::unique_ptr<NKikimr::NMiniKQL::TTypeEnvironment> TypeEnv; std::unique_ptr<NKikimr::NMiniKQL::TProgramBuilder> ProgramBuilder; + NKikimr::NMiniKQL::TMemoryUsageInfo MemInfo; + std::unique_ptr<NKikimr::NMiniKQL::THolderFactory> HolderFactory; }; class TTopicParserBase : public ITopicParser, public TTypeParser { diff --git a/ydb/core/fq/libs/row_dispatcher/format_handler/ut/format_handler_ut.cpp b/ydb/core/fq/libs/row_dispatcher/format_handler/ut/format_handler_ut.cpp index ed2a0bca86e..0e7042711c5 100644 --- a/ydb/core/fq/libs/row_dispatcher/format_handler/ut/format_handler_ut.cpp +++ b/ydb/core/fq/libs/row_dispatcher/format_handler/ut/format_handler_ut.cpp @@ -480,7 +480,7 @@ Y_UNIT_TEST_SUITE(TestFormatHandler) { CheckError( MakeClient( - {{"data_2", "[ListType; [DataType; String]]"}}, + {{"data_2", "[EmptyListType; [DataType; String]]"}}, "", "", EmptyCheck(), diff --git a/ydb/core/fq/libs/row_dispatcher/format_handler/ut/topic_parser_ut.cpp b/ydb/core/fq/libs/row_dispatcher/format_handler/ut/topic_parser_ut.cpp index 6a796988265..0d277f7fb14 100644 --- a/ydb/core/fq/libs/row_dispatcher/format_handler/ut/topic_parser_ut.cpp +++ b/ydb/core/fq/libs/row_dispatcher/format_handler/ut/topic_parser_ut.cpp @@ -327,6 +327,203 @@ Y_UNIT_TEST_SUITE(TestJsonParser) { }); } + Y_UNIT_TEST_F(NestedListTypes, TJsonParserFixture) { + ExpectedBatches = 1; + + CheckSuccess(MakeParser({{"nested", "[OptionalType; [ListType; [DataType; String]]]"}, {"a1", "[DataType; String]"}}, [&](ui64 numberRows, TVector<std::span<NYql::NUdf::TUnboxedValue>> result) { + UNIT_ASSERT_VALUES_EQUAL(4, numberRows); + + UNIT_ASSERT_VALUES_EQUAL(2, result.size()); + UNIT_ASSERT(result[0][0]); + UNIT_ASSERT_VALUES_EQUAL(result[0][0].GetListLength(), 2); + { + auto val = result[0][0].Lookup(NYql::NUdf::TUnboxedValuePod(0)); + UNIT_ASSERT_VALUES_EQUAL("key1", TString(val.AsStringRef())); + } + { + auto val = result[0][0].Lookup(NYql::NUdf::TUnboxedValuePod(1)); + UNIT_ASSERT_VALUES_EQUAL("key2", TString(val.AsStringRef())); + } + UNIT_ASSERT_VALUES_EQUAL("hello1", TString(result[1][0].AsStringRef())); + + UNIT_ASSERT(result[0][1]); + UNIT_ASSERT_VALUES_EQUAL(result[0][1].GetListLength(), 0); + UNIT_ASSERT_VALUES_EQUAL("hello2", TString(result[1][1].AsStringRef())); + + UNIT_ASSERT(!result[0][2]); + UNIT_ASSERT_VALUES_EQUAL("hello3", TString(result[1][2].AsStringRef())); + UNIT_ASSERT(!result[0][3]); + UNIT_ASSERT_VALUES_EQUAL("hello4", TString(result[1][3].AsStringRef())); + })); + + Parser->ParseMessages({ + GetMessage(FIRST_OFFSET, R"({"a1": "hello1", "nested": ["key1", "key2"]})"), + GetMessage(FIRST_OFFSET + 1, R"({"a1": "hello2", "nested": []})"), + GetMessage(FIRST_OFFSET + 2, R"({"a1": "hello3"})"), + GetMessage(FIRST_OFFSET + 3, R"({"a1": "hello4", "nested": null})"), + }); + } + + Y_UNIT_TEST_F(NestedTupleTypes, TJsonParserFixture) { + ExpectedBatches = 1; + + CheckSuccess(MakeParser({{"nested", "[OptionalType; [TupleType; [[DataType; String]; [OptionalType; [DataType; Int64]]]]]"}, {"a1", "[DataType; String]"}}, [&](ui64 numberRows, TVector<std::span<NYql::NUdf::TUnboxedValue>> result) { + UNIT_ASSERT_VALUES_EQUAL(4, numberRows); + + UNIT_ASSERT_VALUES_EQUAL(2, result.size()); + UNIT_ASSERT(result[0][0]); + UNIT_ASSERT_VALUES_EQUAL(result[0][0].GetListLength(), 2); + { + auto val = result[0][0].GetElement(0); + UNIT_ASSERT_VALUES_EQUAL("key1", TString(val.AsStringRef())); + } + UNIT_ASSERT_VALUES_EQUAL(12, result[0][0].GetElement(1).Get<i64>()); + UNIT_ASSERT_VALUES_EQUAL("hello1", TString(result[1][0].AsStringRef())); + + UNIT_ASSERT(result[0][1]); + UNIT_ASSERT_VALUES_EQUAL(result[0][1].GetListLength(), 2); + { + auto val = result[0][1].GetElement(0); + UNIT_ASSERT_VALUES_EQUAL("key2", TString(val.AsStringRef())); + } + UNIT_ASSERT(!result[0][1].GetElement(1)); + UNIT_ASSERT_VALUES_EQUAL("hello2", TString(result[1][1].AsStringRef())); + + UNIT_ASSERT(!result[0][2]); + UNIT_ASSERT_VALUES_EQUAL("hello3", TString(result[1][2].AsStringRef())); + UNIT_ASSERT(!result[0][3]); + UNIT_ASSERT_VALUES_EQUAL("hello4", TString(result[1][3].AsStringRef())); + })); + + Parser->ParseMessages({ + GetMessage(FIRST_OFFSET, R"({"a1": "hello1", "nested": ["key1", 12, true]})"), + GetMessage(FIRST_OFFSET + 1, R"({"a1": "hello2", "nested": ["key2"]})"), + GetMessage(FIRST_OFFSET + 2, R"({"a1": "hello3"})"), + GetMessage(FIRST_OFFSET + 3, R"({"a1": "hello4", "nested": null})"), + }); + } + + Y_UNIT_TEST_F(NestedEmptyTupleTypes, TJsonParserFixture) { + ExpectedBatches = 1; + + CheckSuccess(MakeParser({{"nested", "[OptionalType; [TupleType; []]]"}, {"a1", "[DataType; String]"}}, [&](ui64 numberRows, TVector<std::span<NYql::NUdf::TUnboxedValue>> result) { + UNIT_ASSERT_VALUES_EQUAL(4, numberRows); + + UNIT_ASSERT_VALUES_EQUAL(2, result.size()); + UNIT_ASSERT(result[0][0]); + UNIT_ASSERT_VALUES_EQUAL(result[0][0].GetListLength(), 0); + UNIT_ASSERT_VALUES_EQUAL("hello1", TString(result[1][0].AsStringRef())); + + UNIT_ASSERT(result[0][1]); + UNIT_ASSERT_VALUES_EQUAL(result[0][1].GetListLength(), 0); + UNIT_ASSERT_VALUES_EQUAL("hello2", TString(result[1][1].AsStringRef())); + + UNIT_ASSERT(!result[0][2]); + UNIT_ASSERT_VALUES_EQUAL("hello3", TString(result[1][2].AsStringRef())); + UNIT_ASSERT(!result[0][3]); + UNIT_ASSERT_VALUES_EQUAL("hello4", TString(result[1][3].AsStringRef())); + })); + + Parser->ParseMessages({ + GetMessage(FIRST_OFFSET, R"({"a1": "hello1", "nested": []})"), + GetMessage(FIRST_OFFSET + 1, R"({"a1": "hello2", "nested": ["foobar", 123, true, null]})"), + GetMessage(FIRST_OFFSET + 2, R"({"a1": "hello3"})"), + GetMessage(FIRST_OFFSET + 3, R"({"a1": "hello4", "nested": null})"), + }); + } + + Y_UNIT_TEST_F(NestedStructTypes, TJsonParserFixture) { + ExpectedBatches = 1; + + CheckSuccess(MakeParser({{"nested", "[OptionalType; [StructType; [[a; [DataType; String]]; [b; [OptionalType; [DataType; Int64]]]]]]"}, {"a1", "[DataType; String]"}}, [&](ui64 numberRows, TVector<std::span<NYql::NUdf::TUnboxedValue>> result) { + UNIT_ASSERT_VALUES_EQUAL(4, numberRows); + + UNIT_ASSERT_VALUES_EQUAL(2, result.size()); + UNIT_ASSERT(result[0][0]); + UNIT_ASSERT_VALUES_EQUAL(result[0][0].GetListLength(), 2); + { + auto val = result[0][0].GetElement(0); + UNIT_ASSERT_VALUES_EQUAL("key1", TString(val.AsStringRef())); + } + UNIT_ASSERT_VALUES_EQUAL(12, result[0][0].GetElement(1).Get<i64>()); + UNIT_ASSERT_VALUES_EQUAL("hello1", TString(result[1][0].AsStringRef())); + + UNIT_ASSERT(result[0][1]); + UNIT_ASSERT_VALUES_EQUAL(result[0][1].GetListLength(), 2); + { + auto val = result[0][1].GetElement(0); + UNIT_ASSERT_VALUES_EQUAL("key2", TString(val.AsStringRef())); + } + UNIT_ASSERT(!result[0][1].GetElement(1)); + UNIT_ASSERT_VALUES_EQUAL("hello2", TString(result[1][1].AsStringRef())); + + UNIT_ASSERT(!result[0][2]); + UNIT_ASSERT_VALUES_EQUAL("hello3", TString(result[1][2].AsStringRef())); + UNIT_ASSERT(!result[0][3]); + UNIT_ASSERT_VALUES_EQUAL("hello4", TString(result[1][3].AsStringRef())); + })); + + Parser->ParseMessages({ + GetMessage(FIRST_OFFSET, R"({"a1": "hello1", "nested": {"a":"key1", "b": 12}})"), + GetMessage(FIRST_OFFSET + 1, R"({"a1": "hello2", "nested": {"a":"key2"}})"), + GetMessage(FIRST_OFFSET + 2, R"({"a1": "hello3"})"), + GetMessage(FIRST_OFFSET + 3, R"({"a1": "hello4", "nested": null})"), + }); + } + + Y_UNIT_TEST_F(NestedDictTypes, TJsonParserFixture) { + ExpectedBatches = 1; + + CheckSuccess(MakeParser({{"nested", "[OptionalType; [TupleType; [[DictType; [DataType; String]; [OptionalType; [ListType; [DataType; Int64]]]]]]]"}, {"a1", "[DataType; String]"}}, [&](ui64 numberRows, TVector<std::span<NYql::NUdf::TUnboxedValue>> result) { + UNIT_ASSERT_VALUES_EQUAL(5, numberRows); + UNIT_ASSERT_VALUES_EQUAL(2, result.size()); + + UNIT_ASSERT(result[0][0]); + UNIT_ASSERT_VALUES_EQUAL(result[0][0].GetListLength(), 1); + { + const auto& dict = result[0][0].GetElement(0); + UNIT_ASSERT(dict); + UNIT_ASSERT_VALUES_EQUAL(dict.GetDictLength(), 2); + UNIT_ASSERT(dict.Contains(result[1][0])); + UNIT_ASSERT(dict.Contains(result[1][1])); + } + UNIT_ASSERT_VALUES_EQUAL("hello1", TString(result[1][0].AsStringRef())); + + UNIT_ASSERT(result[0][1]); + UNIT_ASSERT_VALUES_EQUAL(result[0][1].GetListLength(), 1); + { + const auto& dict = result[0][1].GetElement(0); + UNIT_ASSERT(dict); + UNIT_ASSERT_VALUES_EQUAL(dict.GetDictLength(), 1); + UNIT_ASSERT(!dict.Contains(result[1][0])); + UNIT_ASSERT(dict.Contains(result[1][1])); + } + UNIT_ASSERT_VALUES_EQUAL("hello2", TString(result[1][1].AsStringRef())); + + UNIT_ASSERT(!result[0][2]); + UNIT_ASSERT_VALUES_EQUAL("hello3", TString(result[1][2].AsStringRef())); + UNIT_ASSERT(!result[0][3]); + UNIT_ASSERT_VALUES_EQUAL("hello4", TString(result[1][3].AsStringRef())); + + UNIT_ASSERT(result[0][4]); + UNIT_ASSERT_VALUES_EQUAL(result[0][4].GetListLength(), 1); + { + const auto& dict = result[0][4].GetElement(0); + UNIT_ASSERT(dict); + UNIT_ASSERT_VALUES_EQUAL(dict.GetDictLength(), 0); + } + UNIT_ASSERT_VALUES_EQUAL("hello5", TString(result[1][4].AsStringRef())); + })); + + Parser->ParseMessages({ + GetMessage(FIRST_OFFSET, R"({"a1": "hello1", "nested": [{"hello1": [10], "hello2": null},"foo"]})"), + GetMessage(FIRST_OFFSET + 1, R"({"a1": "hello2", "nested": [{"hello2": []},42]})"), + GetMessage(FIRST_OFFSET + 2, R"({"a1": "hello3"})"), + GetMessage(FIRST_OFFSET + 3, R"({"a1": "hello4", "nested": null})"), + GetMessage(FIRST_OFFSET + 4, R"({"a1": "hello5", "nested": [{}, {}]})"), + }); + } + Y_UNIT_TEST_F(SimpleBooleans, TJsonParserFixture) { ExpectedBatches = 1; @@ -466,6 +663,33 @@ Y_UNIT_TEST_SUITE(TestJsonParser) { }); } + Y_UNIT_TEST_F(MissingRepeatedNestedFieldsValidation, TJsonParserFixture) { + CheckSuccess(MakeParser({"a1","a2"}, "[OptionalType; [StructType; [[a1; [DataType; Uint64]]; [a2; [DataType; Uint64]]]]]")); + ParserHandler->ExpectColumnError(1, EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Failed to parse nested json value (Struct), expected non-optional field a1"); + ParserHandler->ExpectColumnError(0, EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Failed to parse nested json value (Struct), expected non-optional field a2"); + ExpectedBatches++; + Parser->ParseMessages({ + GetMessage(FIRST_OFFSET, R"({"a1":{"a1": 101, "a1": 102}})"), + GetMessage(FIRST_OFFSET + 1, R"({"a2": {"a2":103, "a2": 104}})") + }); + } + + Y_UNIT_TEST_F(MissingNestedStructFieldsValidation, TJsonParserFixture) { + CheckSuccess(MakeParser({{"a1", "[StructType; [[a; [DataType; String]]]]"}, {"a2", "[StructType; [[a; [DataType; Uint64]]]]"}})); + CheckColumnError(R"({"a2": {"a":105}, "event": "event1"})", 0, EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Failed to parse json messages, found 1 missing values in non optional column 'a1' with type [StructType; [[a; [DataType; String]]]], buffered offsets: " << FIRST_OFFSET); + CheckColumnError(R"({"a1": {"a":"hello1"}, "a2": null, "event": "event1"})", 1, EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Failed to parse json string at offset " << FIRST_OFFSET + 1 << ", got parsing error for column 'a2' with type [StructType; [[a; [DataType; Uint64]]]] subissue: { <main>: Error: Found unexpected null value, expected non optional type Struct }"); + CheckColumnError(R"({"a2": {"a":105}, "a1":{}, "event": "event1"})", 0, EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Error: Failed to parse nested json value (Struct), expected non-optional field a"); + CheckColumnError(R"({"a1": {"a":"hello1"}, "a2": {"a":null}, "event": "event1"})", 1, EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Failed to parse json string at offset " << FIRST_OFFSET + 3 << ", got parsing error for column 'a2' with type [StructType; [[a; [DataType; Uint64]]]] subissue: { <main>: Error: Found unexpected null value, expected non optional type Uint64 }"); + } + + Y_UNIT_TEST_F(MissingNestedTupleFieldsValidation, TJsonParserFixture) { + CheckSuccess(MakeParser({{"a1", "[TupleType; [[DataType; String]]]"}, {"a2", "[TupleType; [[DataType; Uint64]]]"}})); + CheckColumnError(R"({"a2": [105], "event": "event1"})", 0, EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Failed to parse json messages, found 1 missing values in non optional column 'a1' with type [TupleType; [[DataType; String]]], buffered offsets: " << FIRST_OFFSET); + CheckColumnError(R"({"a1": ["hello1"], "a2": null, "event": "event1"})", 1, EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Failed to parse json string at offset " << FIRST_OFFSET + 1 << ", got parsing error for column 'a2' with type [TupleType; [[DataType; Uint64]]] subissue: { <main>: Error: Found unexpected null value, expected non optional type Tuple }"); + CheckColumnError(R"({"a2": [105], "a1":[], "event": "event1"})", 0, EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Error: Failed to parse nested json value (Tuple), short json array at index 0, expected non optional type String"); + CheckColumnError(R"({"a1": ["hello1"], "a2": [null], "event": "event1"})", 1, EStatusId::PRECONDITION_FAILED, TStringBuilder() << "Error: Found unexpected null value, expected non optional type Uint64"); + } + Y_UNIT_TEST_F(TypeKindsValidation, TJsonParserFixture) { CheckError( MakeParser({{"a1", "[[BAD TYPE]]"}}), @@ -473,9 +697,34 @@ Y_UNIT_TEST_SUITE(TestJsonParser) { "Failed to parse column 'a1' type [[BAD TYPE]] subissue: { <main>: Error: Failed to parse type from yson: Failed to parse scheme from YSON:" ); CheckError( - MakeParser({{"a2", "[OptionalType; [DataType; String]]"}, {"a1", "[ListType; [DataType; String]]"}}), + MakeParser({{"a2", "[OptionalType; [DataType; String]]"}, {"a1", "[EmptyListType; [DataType; Uint8]]"}}), + EStatusId::UNSUPPORTED, + "Failed to create parser for column 'a1' with type [EmptyListType; [DataType; Uint8]] subissue: { <main>: Error: Unsupported type kind: EmptyList }" + ); + CheckError( + MakeParser({{"a2", "[OptionalType; [OptionalType; [DataType; String]]]"}}), EStatusId::UNSUPPORTED, - "Failed to create parser for column 'a1' with type [ListType; [DataType; String]] subissue: { <main>: Error: Unsupported type kind: List }" + "Failed to create parser for column 'a2' with type [OptionalType; [OptionalType; [DataType; String]]] subissue: { <main>: Error: Nested optionals is not supported as input type }" + ); + CheckError( + MakeParser({{"a2", "[ListType; [OptionalType; [OptionalType; [DataType; String]]]]"}}), + EStatusId::UNSUPPORTED, + "Failed to create parser for column 'a2' with type [ListType; [OptionalType; [OptionalType; [DataType; String]]]] subissue: { <main>: Error: Nested optionals is not supported as input type }" + ); + CheckError( + MakeParser({{"a2", "[DictType; [DataType; Int64]; [DataType; String]]"}}), + EStatusId::UNSUPPORTED, + "Failed to create parser for column 'a2' with type [DictType; [DataType; Int64]; [DataType; String]] subissue: { <main>: Error: Dict key type: expected either String, or Utf8, but got Int64 }" + ); + CheckError( + MakeParser({{"a2", "[DictType; [ListType; [DataType; String]]; [DataType; String]]"}}), + EStatusId::UNSUPPORTED, + "Failed to create parser for column 'a2' with type [DictType; [ListType; [DataType; String]]; [DataType; String]] subissue: { <main>: Error: Dict key type: expected either String, or Utf8, but got List }" + ); + CheckError( + MakeParser({{"a2", "[OptionalType; [DataType; String]]"}, {"a1", "[ListType; [EmptyListType; [DataType; Uint8]]]"}}), + EStatusId::UNSUPPORTED, + "Failed to create parser for column 'a1' with type [ListType; [EmptyListType; [DataType; Uint8]]] subissue: { <main>: Error: Unsupported type kind: EmptyList }" ); } @@ -521,11 +770,13 @@ Y_UNIT_TEST_SUITE(TestJsonParser) { {"a4", "[DataType; Int8]"}, {"a5", "[DataType; Bool]"}, {"a6", "[DataType; Double]"}, - {"a7", "[DataType; Utf8]"}}; + {"a7", "[DataType; Utf8]"}, + {"a8", "[ListType; [DataType; Utf8]]"}, + }; CheckSuccess(MakeParser(columns, [](ui64 numberRows, TVector<std::span<NYql::NUdf::TUnboxedValue>> result) { UNIT_ASSERT_VALUES_EQUAL(1, numberRows); - UNIT_ASSERT_VALUES_EQUAL(7, result.size()); + UNIT_ASSERT_VALUES_EQUAL(8, result.size()); UNIT_ASSERT_VALUES_EQUAL("hello1", TString(result[0][0].AsStringRef())); UNIT_ASSERT_VALUES_EQUAL(101, result[1][0].Get<ui64>()); UNIT_ASSERT_VALUES_EQUAL(102, result[2][0].Get<i64>()); @@ -533,8 +784,9 @@ Y_UNIT_TEST_SUITE(TestJsonParser) { UNIT_ASSERT_VALUES_EQUAL(true, result[4][0].Get<bool>()); UNIT_ASSERT_VALUES_EQUAL("146.4", ToString(result[5][0].Get<double>())); UNIT_ASSERT_VALUES_EQUAL("hi", TString(result[6][0].AsStringRef())); + UNIT_ASSERT_VALUES_EQUAL(0, result[7][0].GetListLength()); })); - PushToParser(FIRST_OFFSET, R"({"a1": "hello1", "a2": 101, "a3": 102, "a4": -2, "a5": true, "a6": 146.4, "a7": "hi", "event": "event1"})"); + PushToParser(FIRST_OFFSET, R"({"a1": "hello1", "a2": 101, "a3": 102, "a4": -2, "a5": true, "a6": 146.4, "a7": "hi", "a8":[], "event": "event1"})"); } Y_UNIT_TEST_F(SkipMissingRepeatedFieldsValidation, TJsonParserFixtureSkipErrors) { @@ -562,11 +814,13 @@ Y_UNIT_TEST_SUITE(TestJsonParser) { {"a4", "[OptionalType; [DataType; Int8]]"}, {"a5", "[OptionalType; [DataType; Bool]]"}, {"a6", "[OptionalType; [DataType; Double]]"}, - {"a7", "[OptionalType; [DataType; Utf8]]"}}; + {"a7", "[OptionalType; [DataType; Utf8]]"}, + {"a8", "[OptionalType; [ListType; [DataType; Utf8]]]"}, + }; CheckSuccess(MakeParser(columns, [](ui64 numberRows, TVector<std::span<NYql::NUdf::TUnboxedValue>> result) { UNIT_ASSERT_VALUES_EQUAL(1, numberRows); - UNIT_ASSERT_VALUES_EQUAL(7, result.size()); + UNIT_ASSERT_VALUES_EQUAL(8, result.size()); UNIT_ASSERT_VALUES_EQUAL("hello1", TString(result[0][0].AsStringRef())); UNIT_ASSERT_VALUES_EQUAL(101, result[1][0].Get<ui64>()); UNIT_ASSERT_VALUES_EQUAL(102, result[2][0].Get<i64>()); @@ -574,8 +828,9 @@ Y_UNIT_TEST_SUITE(TestJsonParser) { UNIT_ASSERT_VALUES_EQUAL(true, result[4][0].Get<bool>()); UNIT_ASSERT_VALUES_EQUAL("146.4", ToString(result[5][0].Get<double>())); UNIT_ASSERT_VALUES_EQUAL("hi", TString(result[6][0].AsStringRef())); + UNIT_ASSERT_VALUES_EQUAL(1, result[7][0].GetListLength()); })); - PushToParser(FIRST_OFFSET, R"({"a1": "hello1", "a2": 101, "a3": 102, "a4": -2, "a5": true, "a6": 146.4, "a7": "hi", "event": "event1"})"); + PushToParser(FIRST_OFFSET, R"({"a1": "hello1", "a2": 101, "a3": 102, "a4": -2, "a5": true, "a6": 146.4, "a7": "hi", "a8": ["a"], "event": "event1"})"); } Y_UNIT_TEST_F(SkipErrorsOptionalNull, TJsonParserFixtureSkipErrors) { @@ -586,11 +841,13 @@ Y_UNIT_TEST_SUITE(TestJsonParser) { {"a4", "[OptionalType; [DataType; Int8]]"}, {"a5", "[OptionalType; [DataType; Bool]]"}, {"a6", "[OptionalType; [DataType; Double]]"}, - {"a7", "[OptionalType; [DataType; Utf8]]"}}; + {"a7", "[OptionalType; [DataType; Utf8]]"}, + {"a8", "[OptionalType; [ListType; [DataType; Utf8]]]"}, + }; CheckSuccess(MakeParser(columns, [](ui64 numberRows, TVector<std::span<NYql::NUdf::TUnboxedValue>> result) { UNIT_ASSERT_VALUES_EQUAL(1, numberRows); - UNIT_ASSERT_VALUES_EQUAL(7, result.size()); + UNIT_ASSERT_VALUES_EQUAL(8, result.size()); UNIT_ASSERT(!result[0][0]); UNIT_ASSERT(!result[1][0]); UNIT_ASSERT(!result[2][0]); @@ -598,8 +855,9 @@ Y_UNIT_TEST_SUITE(TestJsonParser) { UNIT_ASSERT(!result[4][0]); UNIT_ASSERT(!result[5][0]); UNIT_ASSERT(!result[6][0]); + UNIT_ASSERT(!result[7][0]); })); - PushToParser(FIRST_OFFSET, R"({"a1": null, "a2": null, "a3": null, "a4": null, "a5": null, "a6": null, "a7": null, "event": "event1"})"); + PushToParser(FIRST_OFFSET, R"({"a1": null, "a2": null, "a3": null, "a4": null, "a5": null, "a6": null, "a7": null, "a8": null, "event": "event1"})"); } Y_UNIT_TEST_F(SkipErrors_StringValidation, TJsonParserFixtureSkipErrors) { @@ -652,22 +910,34 @@ Y_UNIT_TEST_SUITE(TestJsonParser) { GetMessage(FIRST_OFFSET + 5, "\x80"), }); } - + Y_UNIT_TEST_F(SkipErrors_Optional, TJsonParserFixtureSkipErrors) { ExpectedBatches = 1; - CheckSuccess(MakeParser({{"a1", "[OptionalType; [DataType; String]]"}, {"a2", "[OptionalType; [DataType; String]]"}}, [&](ui64 numberRows, TVector<std::span<NYql::NUdf::TUnboxedValue>> result) { + CheckSuccess(MakeParser( + { + {"a1", "[OptionalType; [DataType; String]]"}, + {"a2", "[OptionalType; [DataType; String]]"}, + {"a3", "[OptionalType; [ListType; [DataType; String]]]"}, + }, + [&](ui64 numberRows, TVector<std::span<NYql::NUdf::TUnboxedValue>> result) { UNIT_ASSERT_VALUES_EQUAL(2, numberRows); - UNIT_ASSERT_VALUES_EQUAL(2, result.size()); + UNIT_ASSERT_VALUES_EQUAL(3, result.size()); UNIT_ASSERT_VALUES_EQUAL("hello0", TString(result[0][0].AsStringRef())); + UNIT_ASSERT(result[2][0]); + UNIT_ASSERT_VALUES_EQUAL(result[2][0].GetListLength(), 2); UNIT_ASSERT_VALUES_EQUAL("100", TString(result[1][0].AsStringRef())); UNIT_ASSERT(!result[0][1]); UNIT_ASSERT_VALUES_EQUAL("102", TString(result[1][1].AsStringRef())); + UNIT_ASSERT(!result[2][1]); }, false)); Parser->ParseMessages({ - GetMessage(FIRST_OFFSET, R"({"a1": "hello0", "a2": "100"})"), + GetMessage(FIRST_OFFSET, R"({"a1": "hello0", "a2": "100", "a3": ["a", "b"]})"), GetMessage(FIRST_OFFSET + 1, R"({"a1": "hello1", "a2": 101})"), - GetMessage(FIRST_OFFSET + 2, R"({"a2": "102"})") + GetMessage(FIRST_OFFSET + 2, R"({"a1": "hello2", "a2": "100", "a3": 123})"), + GetMessage(FIRST_OFFSET + 3, R"({"a1": "hello2", "a2": "100", "a3": [123]})"), + GetMessage(FIRST_OFFSET + 4, R"({"a2": "102"})"), + GetMessage(FIRST_OFFSET + 5, R"({"a1": "hello2", "a2": "100", "a3": {}})"), }); } diff --git a/ydb/library/yql/providers/common/pushdown/collection.cpp b/ydb/library/yql/providers/common/pushdown/collection.cpp index cf3377b71f5..36ecd97eee2 100644 --- a/ydb/library/yql/providers/common/pushdown/collection.cpp +++ b/ydb/library/yql/providers/common/pushdown/collection.cpp @@ -219,13 +219,28 @@ private: return false; } - bool IsMemberColumn(const TCoMember& member) const { - // We allow member access only for top level predicate argument - return member.Struct().Raw() == LambdaArg.Raw() - && Settings.IsMemberEnabled(TString(member.Name().Value())); + bool IsMemberColumn(const TCoMember& member) { + // Allow member access for top level predicate argument + if (member.Struct().Raw() == LambdaArg.Raw()) { + return Settings.IsMemberEnabled(TString(member.Name().Value())); + } + if (Settings.IsEnabled(EFlag::AnyExpressionExceptMember)) { + return true; + } + if (Settings.IsEnabled(EFlag::StructOperators)) { + return CheckExpressionNodeForPushdown(member.Struct()); + } + return false; } - bool IsMemberColumn(const TExprBase& node) const { + bool IsSupportedNth(const TCoNth& nth) { + if (Settings.IsEnabled(EFlag::StructOperators)) { + return CheckExpressionNodeForPushdown(nth.Tuple()); + } + return false; + } + + bool IsMemberColumn(const TExprBase& node) { if (const auto member = node.Maybe<TCoMember>()) { return IsMemberColumn(member.Cast()); } @@ -329,6 +344,9 @@ public: if (Settings.IsEnabled(EFlag::AnyExpressionExceptMember)) { return true; } + if (auto maybeNth = node.Maybe<TCoNth>()) { + return IsSupportedNth(maybeNth.Cast()); + } if (auto maybeSafeCast = node.Maybe<TCoSafeCast>()) { return IsSupportedSafeCast(maybeSafeCast.Cast()); } @@ -623,7 +641,7 @@ private: return IsComparableArguments(left, right, true); } - bool JsonExistsCanBePushed(const TCoJsonExists& jsonExists) const { + bool JsonExistsCanBePushed(const TCoJsonExists& jsonExists) { if (!Settings.IsEnabled(EFlag::JsonExistsOperator)) { return false; } @@ -646,7 +664,7 @@ private: return predicateTree.CanBePushed; } - bool ExistsCanBePushed(const TCoExists& exists) const { + bool ExistsCanBePushed(const TCoExists& exists) { return IsMemberColumn(exists.Optional()); } diff --git a/ydb/library/yql/providers/common/pushdown/settings.h b/ydb/library/yql/providers/common/pushdown/settings.h index 4194da8611d..269f210129b 100644 --- a/ydb/library/yql/providers/common/pushdown/settings.h +++ b/ydb/library/yql/providers/common/pushdown/settings.h @@ -49,6 +49,7 @@ struct TSettings { DateCtor = 1 << 30, PredicateAsExpression = ui64{1} << 31, // Predicates can be used in expressions (e.g. (a = b) = (c = d)) AnyExpressionExceptMember = ui64{1} << 32, + StructOperators = ui64{1} << 33, // Struct operators can be used }; explicit TSettings(NLog::EComponent logComponent) diff --git a/ydb/library/yql/providers/generic/connector/api/service/protos/connector.proto b/ydb/library/yql/providers/generic/connector/api/service/protos/connector.proto index 08907f3d0c3..ca6a23f844d 100644 --- a/ydb/library/yql/providers/generic/connector/api/service/protos/connector.proto +++ b/ydb/library/yql/providers/generic/connector/api/service/protos/connector.proto @@ -387,6 +387,16 @@ message TExpression { repeated TExpression operands = 1; } + message TStructMember { + TExpression operand = 1; + string field = 2; + } + + message TTupleNth { + TExpression operand = 1; + uint64 field = 2; + } + oneof payload { // A scalar value Ydb.TypedValue typed_value = 1; @@ -412,6 +422,10 @@ message TExpression { TCurrentUtcTimestamp current_utc_timestamp = 11; TPredicate predicate = 12; + + TStructMember struct_member = 13; + + TTupleNth tuple_nth = 14; } } diff --git a/ydb/library/yql/providers/generic/provider/yql_generic_predicate_pushdown.cpp b/ydb/library/yql/providers/generic/provider/yql_generic_predicate_pushdown.cpp index bed00078eaf..d582354998c 100644 --- a/ydb/library/yql/providers/generic/provider/yql_generic_predicate_pushdown.cpp +++ b/ydb/library/yql/providers/generic/provider/yql_generic_predicate_pushdown.cpp @@ -28,6 +28,8 @@ namespace NYql { TString FormatCoalesce(const TExpression::TCoalesce& coalesce); TString FormatIfExpression(const TExpression::TIf& sqlIf); TString FormatUnwrap(const TExpression::TUnwrap& unwrap); + TString FormatStructMember(const TExpression::TStructMember& structMember); + TString FormatTupleNth(const TExpression::TTupleNth& tupleNth); TString FormatMinOf(const TExpression::TMinOf& minOf); TString FormatMaxOf(const TExpression::TMaxOf& maxOf); TString FormatCurrentUtcTimestamp(const TExpression::TCurrentUtcTimestamp& currentUtcTimestamp); @@ -55,13 +57,27 @@ namespace NYql { TExprContext& Ctx; }; - bool SerializeMember(const TCoMember& member, TExpression* proto, TSerializationContext& ctx) { - if (member.Struct().Raw() != ctx.Arg.Raw()) { // member callable called not for lambda argument - ctx.Err << "member callable called not for lambda argument"; + bool SerializeExpression(const TExprBase& expression, TExpression* proto, TSerializationContext& ctx, ui64 depth); + + bool SerializeMember(const TCoMember& member, TExpression* proto, TSerializationContext& ctx, ui64 depth) { + if (member.Struct().Raw() == ctx.Arg.Raw()) { // member callable called for lambda argument + proto->set_column(member.Name().StringValue()); + return true; + } + auto structMember = proto->mutable_struct_member(); + structMember->set_field(member.Name().StringValue()); + return SerializeExpression(member.Struct(), structMember->mutable_operand(), ctx, depth + 1); + } + + bool SerializeNth(const TCoNth& nth, TExpression* proto, TSerializationContext& ctx, ui64 depth) { + auto tupleNth = proto->mutable_tuple_nth(); + auto index = TryFromString<ui64>(nth.Index().StringValue()); + if (!index) { + ctx.Err << "Nth: expected ui64, got " << nth.Index().StringValue(); return false; } - proto->set_column(member.Name().StringValue()); - return true; + tupleNth->set_field(*index); + return SerializeExpression(nth.Tuple(), tupleNth->mutable_operand(), ctx, depth + 1); } bool SerializeLambdaArgument(const TExprBase& node, TExpression* proto, TSerializationContext& ctx) { @@ -84,7 +100,6 @@ namespace NYql { return TString(from); } - bool SerializeExpression(const TExprBase& expression, TExpression* proto, TSerializationContext& ctx, ui64 depth); bool SerializeCompare(const TCoCompare& compare, TPredicate* predicateProto, TSerializationContext& ctx, ui64 depth); bool SerializeApply(const TCoApply& apply, TPredicate* proto, TSerializationContext& ctx, ui64 depth); bool SerializeExists(const TCoExists& exists, TPredicate* proto, TSerializationContext& ctx, bool withNot, ui64 depth); @@ -301,7 +316,10 @@ namespace NYql { bool SerializeExpression(const TExprBase& expression, TExpression* proto, TSerializationContext& ctx, ui64 depth) { if (auto member = expression.Maybe<TCoMember>()) { - return SerializeMember(member.Cast(), proto, ctx); + return SerializeMember(member.Cast(), proto, ctx, depth); + } + if (auto nth = expression.Maybe<TCoNth>()) { + return SerializeNth(nth.Cast(), proto, ctx, depth); } if (auto coalesce = expression.Maybe<TCoCoalesce>()) { return SerializeCoalesceExpression(coalesce.Cast(), proto, ctx, depth); @@ -558,8 +576,12 @@ namespace NYql { return SerializePredicate(notExpr.Value(), dstProto->mutable_operand(), ctx, depth + 1); } - bool SerializeMember(const TCoMember& member, TPredicate* proto, TSerializationContext& ctx) { - return SerializeMember(member, proto->mutable_bool_expression()->mutable_value(), ctx); + bool SerializeMember(const TCoMember& member, TPredicate* proto, TSerializationContext& ctx, ui64 depth) { + return SerializeMember(member, proto->mutable_bool_expression()->mutable_value(), ctx, depth); + } + + bool SerializeNth(const TCoNth& nth, TPredicate* proto, TSerializationContext& ctx, ui64 depth) { + return SerializeNth(nth, proto->mutable_bool_expression()->mutable_value(), ctx, depth); } bool SerializeRegexp(const TCoUdf& regexp, const TExprNode::TListType& children, TPredicate* proto, TSerializationContext& ctx, ui64 depth) { @@ -618,7 +640,10 @@ namespace NYql { return SerializeNot(notExpr.Cast(), proto, ctx, depth); } if (auto member = predicate.Maybe<TCoMember>()) { - return SerializeMember(member.Cast(), proto, ctx); + return SerializeMember(member.Cast(), proto, ctx, depth); + } + if (auto nth = predicate.Maybe<TCoNth>()) { + return SerializeNth(nth.Cast(), proto, ctx, depth); } if (auto exists = predicate.Maybe<TCoExists>()) { return SerializeExists(exists.Cast(), proto, ctx, false, depth); @@ -652,6 +677,14 @@ namespace NYql { return NFq::EncloseAndEscapeString(value, '`'); } + TString FormatStructMember(const TExpression::TStructMember& structMember) { + return TStringBuilder() << '(' << FormatExpression(structMember.operand()) << ')' << "." << NFq::EncloseAndEscapeString(structMember.field(), '`'); + } + + TString FormatTupleNth(const TExpression::TTupleNth& tupleNth) { + return TStringBuilder() << '(' << FormatExpression(tupleNth.operand()) << ')' << ".`" << tupleNth.field() << "`"; + } + TString FormatValue(const Ydb::Value& value) { switch (value.value_case()) { case Ydb::Value::kBoolValue: @@ -851,6 +884,10 @@ namespace NYql { return FormatCast(expression.cast()); case TExpression::kUnwrap: return FormatUnwrap(expression.unwrap()); + case TExpression::kStructMember: + return FormatStructMember(expression.struct_member()); + case TExpression::kTupleNth: + return FormatTupleNth(expression.tuple_nth()); case TExpression::kMinOf: return FormatMinOf(expression.min_of()); case TExpression::kMaxOf: diff --git a/ydb/library/yql/providers/generic/pushdown/yql_generic_match_predicate.cpp b/ydb/library/yql/providers/generic/pushdown/yql_generic_match_predicate.cpp index ff86c0b7c64..b9dc5769fbc 100644 --- a/ydb/library/yql/providers/generic/pushdown/yql_generic_match_predicate.cpp +++ b/ydb/library/yql/providers/generic/pushdown/yql_generic_match_predicate.cpp @@ -61,6 +61,8 @@ namespace NYql::NGenericPushDown { case NYql::NConnector::NApi::TExpression::kMaxOf: case NYql::NConnector::NApi::TExpression::kCurrentUtcTimestamp: case NYql::NConnector::NApi::TExpression::kPredicate: + case NYql::NConnector::NApi::TExpression::kStructMember: + case NYql::NConnector::NApi::TExpression::kTupleNth: case NYql::NConnector::NApi::TExpression::PAYLOAD_NOT_SET: return false; } @@ -82,6 +84,8 @@ namespace NYql::NGenericPushDown { case NYql::NConnector::NApi::TExpression::kMaxOf: case NYql::NConnector::NApi::TExpression::kCurrentUtcTimestamp: case NYql::NConnector::NApi::TExpression::kPredicate: + case NYql::NConnector::NApi::TExpression::kStructMember: + case NYql::NConnector::NApi::TExpression::kTupleNth: case NYql::NConnector::NApi::TExpression::PAYLOAD_NOT_SET: return false; } @@ -305,6 +309,8 @@ namespace NYql::NGenericPushDown { case NYql::NConnector::NApi::TExpression::kMaxOf: case NYql::NConnector::NApi::TExpression::kCurrentUtcTimestamp: case NYql::NConnector::NApi::TExpression::kPredicate: + case NYql::NConnector::NApi::TExpression::kStructMember: + case NYql::NConnector::NApi::TExpression::kTupleNth: case NYql::NConnector::NApi::TExpression::PAYLOAD_NOT_SET: return Triple::Unknown; } diff --git a/ydb/library/yql/providers/pq/provider/yql_pq_datasource_type_ann.cpp b/ydb/library/yql/providers/pq/provider/yql_pq_datasource_type_ann.cpp index b87d6354518..d2e2410a8a2 100644 --- a/ydb/library/yql/providers/pq/provider/yql_pq_datasource_type_ann.cpp +++ b/ydb/library/yql/providers/pq/provider/yql_pq_datasource_type_ann.cpp @@ -66,6 +66,7 @@ struct TWatermarkPushdownSettings: public NPushdown::TSettings { EFlag::ToBytesFromStringExpressions | EFlag::ToStringFromStringExpressions | EFlag::FlatMapOverOptionals | + EFlag::StructOperators | EFlag::NonDeterministic ); } diff --git a/ydb/library/yql/providers/pq/provider/yql_pq_logical_opt.cpp b/ydb/library/yql/providers/pq/provider/yql_pq_logical_opt.cpp index 38465f00462..bb03df1a701 100644 --- a/ydb/library/yql/providers/pq/provider/yql_pq_logical_opt.cpp +++ b/ydb/library/yql/providers/pq/provider/yql_pq_logical_opt.cpp @@ -40,6 +40,7 @@ struct TPushdownSettings: public NPushdown::TSettings { EFlag::StringTypes | EFlag::LikeOperator | EFlag::DoNotCheckCompareArgumentsTypes | EFlag::InOperator | EFlag::IsDistinctOperator | EFlag::JustPassthroughOperators | EFlag::DivisionExpressions | EFlag::CastExpression | EFlag::ToBytesFromStringExpressions | EFlag::FlatMapOverOptionals | EFlag::PredicateAsExpression | + EFlag::StructOperators | // Split features EFlag::SplitOrOperator diff --git a/ydb/tests/fq/yds/test_row_dispatcher.py b/ydb/tests/fq/yds/test_row_dispatcher.py index 757777c585a..1570455bbc3 100644 --- a/ydb/tests/fq/yds/test_row_dispatcher.py +++ b/ydb/tests/fq/yds/test_row_dispatcher.py @@ -319,6 +319,53 @@ class TestPqRowDispatcher(TestYdsBase): assert "Row dispatcher will use the predicate:" in issues, "Incorrect Issues: " + issues @yq_v1 + def test_nested_structured_types(self, kikimr, client): + self.init(client, "test_nested_structured_types") + + large_string = "abcdefghjkl1234567890+abcdefghjkl1234567890" + sql = Rf''' + INSERT INTO {YDS_CONNECTION}.`{self.output_topic}` + SELECT ToBytes(Unwrap(Json::SerializeJson(Yson::From(data)))) + FROM {YDS_CONNECTION}.`{self.input_topic}` + WITH ( + format = json_each_row, + SCHEMA = ( + time UInt64 NOT NULL, + data Struct<key: String, second_key: Tuple<Int64, Bool?>> NOT NULL, + event String NOT NULL + ) + ) + WHERE (event = "event1" or event = "event2" or event = "event4" or event = "event5") + AND (data.key = "value" OR data.key = "ev2" OR data.key = "{large_string}") + AND data.second_key.0 > 0 AND data.second_key.1 IS DISTINCT FROM true; + ''' + + query_id = start_yds_query(kikimr, client, sql) + wait_actor_count(kikimr, "FQ_ROW_DISPATCHER_SESSION", 1) + + data = [ + '{"time": 101, "data": {"key": "value", "second_key":[1, false]}, "event": "event1"}', + '{"time": 102, "data": {"second_key":[2], "key":"' + large_string + '"}, "event": "event2"}', + '{"time": 103, "data": {"key":"ev2", "second_key":[3]}, "event": "event3"}', + '{"time": 104, "data": {"key":"ev2", "second_key":[42, true]}, "event": "event4"}', + '{"time": 105, "data": {"key":"ev2", "second_key":[42, false]}, "event": "event5"}', + ] + + self.write_stream(data) + expected = [ + '{"key":"value","second_key":[1,false]}', + '{"key":"' + large_string + '","second_key":[2,null]}', + '{"key":"ev2","second_key":[42,false]}', + ] + assert self.read_stream(len(expected), topic_path=self.output_topic) == expected + + wait_actor_count(kikimr, "DQ_PQ_READ_ACTOR", 1) + stop_yds_query(client, query_id) + + issues = str(client.describe_query(query_id).result.query.transient_issue) + assert "Row dispatcher will use the predicate:" in issues and "second_key" in issues, "Incorrect Issues: " + issues + + @yq_v1 def test_nested_types_without_predicate(self, kikimr, client): self.init(client, "test_nested_types_without_predicate") |
