diff --git a/ydb/apps/ydbd/ya.make b/ydb/apps/ydbd/ya.make index 2e92c7c532798..8802ff056dbb6 100644 --- a/ydb/apps/ydbd/ya.make +++ b/ydb/apps/ydbd/ya.make @@ -55,6 +55,7 @@ PEERDIR( ydb/library/yql/udfs/common/hybrid_search ydb/library/yql/udfs/common/knn ydb/library/yql/udfs/common/roaring + ydb/library/yql/udfs/common/rowid ydb/library/yql/udfs/statistics_internal yql/essentials/parser/pg_wrapper yql/essentials/sql/pg diff --git a/ydb/core/engine/minikql/minikql_engine_host.cpp b/ydb/core/engine/minikql/minikql_engine_host.cpp index 6b01c7651213b..20b1b2494554a 100644 --- a/ydb/core/engine/minikql/minikql_engine_host.cpp +++ b/ydb/core/engine/minikql/minikql_engine_host.cpp @@ -1086,6 +1086,7 @@ NUdf::TUnboxedValue GetCellValue(const TCell& cell, NScheme::TTypeInfo type) { case NYql::NProto::TypeIds::Yson: case NYql::NProto::TypeIds::Json: case NYql::NProto::TypeIds::Uuid: + case NYql::NProto::TypeIds::Rowid: case NYql::NProto::TypeIds::JsonDocument: case NYql::NProto::TypeIds::DyNumber: return MakeString(NUdf::TStringRef(cell.Data(), cell.Size())); diff --git a/ydb/core/engine/mkql_proto.cpp b/ydb/core/engine/mkql_proto.cpp index 05442014ad39e..a070854756d52 100644 --- a/ydb/core/engine/mkql_proto.cpp +++ b/ydb/core/engine/mkql_proto.cpp @@ -7,6 +7,7 @@ #include #include #include +#include #include @@ -223,6 +224,15 @@ bool CellsFromTuple(const NKikimrMiniKQL::TType* tupleType, } break; } + case NScheme::NTypeIds::Rowid: + { + if (!v.HasBytes()) { + CHECK_OR_RETURN_ERROR(false, Sprintf("Cannot parse value of type Rowid in tuple at position %" PRIu32, i)); + } + CHECK_OR_RETURN_ERROR(v.GetBytes().size() == NRowid::ROWID_LEN, Sprintf("Invalid Rowid size in tuple at position %" PRIu32, i)); + c = TCell(v.GetBytes().data(), v.GetBytes().size()); + break; + } case NScheme::NTypeIds::Decimal: { if (v.HasLow128() && v.HasHi128()) { @@ -366,6 +376,11 @@ bool CellToValue(NScheme::TTypeInfo type, const TCell& c, NKikimrMiniKQL::TValue break; } + case NScheme::NTypeIds::Rowid: + Y_ABORT_UNLESS(c.Size() == NRowid::ROWID_LEN); + val.MutableOptional()->SetBytes(c.Data(), c.Size()); + break; + default: errStr = "Unknown type: " + ToString(typeId); return false; diff --git a/ydb/core/kqp/common/kqp_resolve.cpp b/ydb/core/kqp/common/kqp_resolve.cpp index 0af663929a5bf..43b6777dd0e90 100644 --- a/ydb/core/kqp/common/kqp_resolve.cpp +++ b/ydb/core/kqp/common/kqp_resolve.cpp @@ -68,6 +68,10 @@ NUdf::TUnboxedValue MakeDefaultValueByType(NKikimr::NMiniKQL::TType* type) { buf.half[1] = 0; return NKikimr::NMiniKQL::MakeString(NUdf::TStringRef(buf.bytes, 16)); } + case NUdf::TDataType::Id: { + char bytes[NUdf::ROWID_SIZE] = {}; + return NKikimr::NMiniKQL::MakeString(NUdf::TStringRef(bytes, sizeof(bytes))); + } default: return NKikimr::NMiniKQL::MakeString(""); } diff --git a/ydb/core/kqp/common/result_set_format/kqp_formats_arrow.cpp b/ydb/core/kqp/common/result_set_format/kqp_formats_arrow.cpp index 8af3de9eba729..46355007cdbca 100644 --- a/ydb/core/kqp/common/result_set_format/kqp_formats_arrow.cpp +++ b/ydb/core/kqp/common/result_set_format/kqp_formats_arrow.cpp @@ -23,7 +23,9 @@ std::shared_ptr BuildArrowType(NUdf::EDataSlot slot) { template <> std::shared_ptr BuildArrowType(NUdf::EDataSlot slot) { - Y_UNUSED(slot); + if (slot == NUdf::EDataSlot::Rowid) { + return arrow::fixed_size_binary(NUdf::ROWID_SIZE); + } return arrow::fixed_size_binary(NScheme::FSB_SIZE); } @@ -376,6 +378,12 @@ void AppendDataValue(arrow::ArrayBuilder* builder, N break; } + case NUdf::EDataSlot::Rowid: { + auto data = value.AsStringRef(); + status = typedBuilder->Append(data.Data()); + break; + } + case NUdf::EDataSlot::Decimal: { auto intVal = value.GetInt128(); status = typedBuilder->Append(reinterpret_cast(&intVal)); diff --git a/ydb/core/kqp/common/result_set_format/kqp_formats_arrow.h b/ydb/core/kqp/common/result_set_format/kqp_formats_arrow.h index b7cadd66ac683..96ebd8d1017ae 100644 --- a/ydb/core/kqp/common/result_set_format/kqp_formats_arrow.h +++ b/ydb/core/kqp/common/result_set_format/kqp_formats_arrow.h @@ -32,7 +32,7 @@ constexpr size_t MAX_VARIANT_DEPTH = 2; * - Temporal types: Date, Datetime, Timestamp, Interval (and their extended variants) * - String types: Utf8, Json, JsonDocument (serialized to string), DyNumber (serialized to string) -> arrow::StringType * - Binary types: String, Yson -> arrow::BinaryType - * - Fixed-size binary: Decimal, Uuid -> arrow::FixedSizeBinaryType + * - Fixed-size binary: Decimal, Uuid, Rowid -> arrow::FixedSizeBinaryType * - Timezone-aware: TzDate, TzDatetime, TzTimestamp -> arrow::StructType * * @tparam TFunc Callable type accepting a single template parameter (Arrow type) @@ -94,6 +94,7 @@ bool SwitchMiniKQLDataTypeToArrowType(NUdf::EDataSlot typeId, TFunc&& callback) case NUdf::EDataSlot::Decimal: case NUdf::EDataSlot::Uuid: + case NUdf::EDataSlot::Rowid: return callback.template operator()(); case NUdf::EDataSlot::TzDate: diff --git a/ydb/core/kqp/opt/kqp_statistics_transformer.cpp b/ydb/core/kqp/opt/kqp_statistics_transformer.cpp index 2c5d51f966efe..d7586c46bf968 100644 --- a/ydb/core/kqp/opt/kqp_statistics_transformer.cpp +++ b/ydb/core/kqp/opt/kqp_statistics_transformer.cpp @@ -781,6 +781,8 @@ double EstimateRowSize(const TStructExprType& rowType, const TString& format, co break; case EDataSlot::Uuid: break; + case EDataSlot::Rowid: + break; case EDataSlot::Date: result += decoded ? 2.0 : 1.51; break; diff --git a/ydb/core/kqp/provider/yql_kikimr_provider.cpp b/ydb/core/kqp/provider/yql_kikimr_provider.cpp index 3d962b94eb76c..9dd3e55f6b4bb 100644 --- a/ydb/core/kqp/provider/yql_kikimr_provider.cpp +++ b/ydb/core/kqp/provider/yql_kikimr_provider.cpp @@ -17,6 +17,7 @@ #include #include #include +#include #include namespace NYql { @@ -833,6 +834,11 @@ void FillLiteralProto(const NNodes::TCoDataCtor& literal, Ydb::TypedValue& proto protoValue.set_high_128(uuidData[1]); break; } + case EDataSlot::Rowid: { + YQL_ENSURE(value.size() == NKikimr::NRowid::ROWID_LEN, "Invalid Rowid literal size"); + protoValue.set_bytes_value(value.data(), value.size()); + break; + } default: YQL_ENSURE(false, "Unexpected type slot " << slot); diff --git a/ydb/core/kqp/runtime/kqp_scan_data.cpp b/ydb/core/kqp/runtime/kqp_scan_data.cpp index e81abf7976c5c..fa05283b3e072 100644 --- a/ydb/core/kqp/runtime/kqp_scan_data.cpp +++ b/ydb/core/kqp/runtime/kqp_scan_data.cpp @@ -72,6 +72,12 @@ TBytesStatistics GetUnboxedValueSize(const NUdf::TUnboxedValue& value, const NSc Y_VERIFY_DEBUG_S(size == sizeof(TGUID), "Wrong Uuid size: " << size); return { sizeof(NUdf::TUnboxedValue) + size, size }; } + case NTypeIds::Rowid: + { + const auto size = value.AsStringRef().Size(); + Y_VERIFY_DEBUG_S(size == NUdf::ROWID_SIZE, "Wrong Rowid size: " << size); + return { sizeof(NUdf::TUnboxedValue) + size, size }; + } case NTypeIds::String: case NTypeIds::Utf8: case NTypeIds::Json: diff --git a/ydb/core/scheme/scheme_tablecell.h b/ydb/core/scheme/scheme_tablecell.h index 0a9c2d9d5e002..98144734e4231 100644 --- a/ydb/core/scheme/scheme_tablecell.h +++ b/ydb/core/scheme/scheme_tablecell.h @@ -342,6 +342,13 @@ inline int CompareTypedCells(const TCell& a, const TCell& b, const NScheme::TTyp return CompareCellsAsByteString(a, b, type.IsDescending()); } + case NKikimr::NScheme::NTypeIds::Rowid: + { + Y_ASSERT(a.Size() == 14); + Y_ASSERT(b.Size() == 14); + return CompareCellsAsByteString(a, b, type.IsDescending()); + } + case NKikimr::NScheme::NTypeIds::Decimal: { Y_ASSERT(a.Size() == sizeof(std::pair)); @@ -505,6 +512,7 @@ inline ui64 GetValueHash(NScheme::TTypeInfo info, const TCell& cell) { case NYql::NProto::TypeIds::JsonDocument: case NYql::NProto::TypeIds::DyNumber: case NYql::NProto::TypeIds::Uuid: + case NYql::NProto::TypeIds::Rowid: return ComputeHash(TStringBuf{cell.Data(), cell.Size()}); default: diff --git a/ydb/core/scheme_types/scheme_type_registry.cpp b/ydb/core/scheme_types/scheme_type_registry.cpp index 8c02ef2b4a6c0..675eb0ad9ef64 100644 --- a/ydb/core/scheme_types/scheme_type_registry.cpp +++ b/ydb/core/scheme_types/scheme_type_registry.cpp @@ -42,6 +42,7 @@ TTypeRegistry::TTypeRegistry() RegisterType(); RegisterType(); RegisterType(); + RegisterType(); RegisterType(); RegisterType(); RegisterType(); diff --git a/ydb/core/scheme_types/scheme_types_defs.cpp b/ydb/core/scheme_types/scheme_types_defs.cpp index f3adc09d8866a..bb01e185f09d0 100644 --- a/ydb/core/scheme_types/scheme_types_defs.cpp +++ b/ydb/core/scheme_types/scheme_types_defs.cpp @@ -41,6 +41,7 @@ namespace NNames { DECLARE_TYPED_TYPE_NAME(DyNumber); DECLARE_TYPED_TYPE_NAME(Uuid); + DECLARE_TYPED_TYPE_NAME(Rowid); } void WriteEscapedValue(IOutputStream &out, const char *data, size_t size) { diff --git a/ydb/core/scheme_types/scheme_types_defs.h b/ydb/core/scheme_types/scheme_types_defs.h index 8ecb5008cc264..6daa545c9ffd0 100644 --- a/ydb/core/scheme_types/scheme_types_defs.h +++ b/ydb/core/scheme_types/scheme_types_defs.h @@ -229,6 +229,7 @@ using TLargeBoundedString = TBoundedString<0x200000, NTypeIds::String2m, NNames: namespace NNames { extern const char Decimal[8]; extern const char Uuid[5]; + extern const char Rowid[6]; } class TDecimal : public IIntegerPair {}; @@ -237,6 +238,10 @@ class TUuid : public TTypedType { public: }; +class TRowid : public TTypedType { +public: +}; + //////////////////////////////////////////////////////// /// Datetime types namespace NNames { @@ -289,6 +294,7 @@ class TInterval64 : public IIntegerTypeWithKeyString #include #include +#include #include #include @@ -345,6 +346,11 @@ Y_FORCE_INLINE void ConvertData(NUdf::TDataTypeId typeId, const NKikimrMiniKQL:: res.set_high_128(value.GetHi128()); break; } + case NUdf::TDataType::Id: { + const auto& stringRef = value.GetBytes(); + res.set_bytes_value(stringRef.data(), stringRef.size()); + break; + } case NUdf::TDataType::Id: { const auto json = NBinaryJson::SerializeToJson(value.GetBytes()); res.set_text_value(json); @@ -511,6 +517,15 @@ Y_FORCE_INLINE void ConvertData(NUdf::TDataTypeId typeId, const Ydb::Value& valu res.SetLow128(value.low_128()); res.SetHi128(value.high_128()); break; + case NUdf::TDataType::Id: { + CheckTypeId(value.value_case(), Ydb::Value::kBytesValue, "Rowid"); + const auto& stringRef = value.bytes_value(); + if (!NRowid::IsValidRowidBytes(TStringBuf(stringRef.data(), stringRef.size()))) { + throw yexception() << "Invalid Rowid value"; + } + res.SetBytes(stringRef.data(), stringRef.size()); + break; + } case NUdf::TDataType::Id: { CheckTypeId(value.value_case(), Ydb::Value::kTextValue, "JsonDocument"); const auto binaryJson = NBinaryJson::SerializeToBinaryJson(value.text_value()); @@ -1162,6 +1177,10 @@ bool CheckValueData(NScheme::TTypeInfo type, const TCell& cell, TString& err) { // Uuid value was verified at parsing time break; + case NScheme::NTypeIds::Rowid: + ok = NRowid::IsValidRowidBytes(cell.AsBuf()); + break; + case NScheme::NTypeIds::Pg: // no pg validation here break; @@ -1301,6 +1320,20 @@ bool CellFromProtoVal(const NScheme::TTypeInfo& type, i32 typmod, const Ydb::Val c = TCell((const char*)&valInPool, sizeof(valInPool)); break; } + case NScheme::NTypeIds::Rowid : { + if (!val.Hasbytes_value()) { + err = "Cannot parse value of type Rowid"; + return false; + } + const auto& bytes = val.Getbytes_value(); + if (!NRowid::IsValidRowidBytes(TStringBuf(bytes.data(), bytes.size()))) { + err = "Invalid Rowid value"; + return false; + } + const auto bytesInPool = valueDataPool.AppendString(TStringBuf(bytes.data(), bytes.size())); + c = TCell(bytesInPool.data(), bytesInPool.size()); + break; + } case NScheme::NTypeIds::Pg : { TString binary; bool isText = false; @@ -1451,6 +1484,9 @@ void ProtoValueFromCell(NYdb::TValueBuilder& vb, const NScheme::TTypeInfo& typeI vb.Uuid(TUuidValue(lo, hi)); break; } + case EPrimitiveType::Rowid: + vb.Rowid(TRowidValue(cell.AsBuf().data(), cell.AsBuf().size())); + break; case EPrimitiveType::JsonDocument: vb.JsonDocument(NBinaryJson::SerializeToJson(getString())); break; diff --git a/ydb/library/mkql_proto/mkql_proto.cpp b/ydb/library/mkql_proto/mkql_proto.cpp index 74e61c299084f..1c0404d27e6f3 100644 --- a/ydb/library/mkql_proto/mkql_proto.cpp +++ b/ydb/library/mkql_proto/mkql_proto.cpp @@ -13,6 +13,7 @@ #include #include #include +#include #include #include @@ -142,6 +143,11 @@ Y_FORCE_INLINE void HandleKindDataExport(const TType* type, const NUdf::TUnboxed UuidToYdbProto(stringRef.Data(), stringRef.Size(), res); break; } + case NUdf::TDataType::Id: { + const auto& stringRef = value.AsStringRef(); + res.set_bytes_value(stringRef.Data(), stringRef.Size()); + break; + } case NUdf::TDataType::Id: { NUdf::TUnboxedValue json = ValueToString(NUdf::EDataSlot::JsonDocument, value); const auto stringRef = json.AsStringRef(); @@ -510,6 +516,11 @@ Y_FORCE_INLINE void HandleKindDataExport(const TType* type, const NUdf::TUnboxed UuidToMkqlProto(stringRef.Data(), stringRef.Size(), res); break; } + case NUdf::TDataType::Id: { + auto stringRef = value.AsStringRef(); + res.SetBytes(stringRef.Data(), stringRef.Size()); + break; + } case NUdf::TDataType::Id: { auto stringRef = value.AsStringRef(); res.SetBytes(stringRef.Data(), stringRef.Size()); @@ -843,6 +854,10 @@ Y_FORCE_INLINE NUdf::TUnboxedValue HandleKindDataImport(const TType* type, const buf.half[1] = value.GetHi128(); return MakeString(NUdf::TStringRef(buf.bytes, 16)); } + case NUdf::TDataType::Id: + MKQL_ENSURE_S(oneOfCase == NKikimrMiniKQL::TValue::ValueValueCase::kBytes); + MKQL_ENSURE_S(value.GetBytes().size() == NRowid::ROWID_LEN); + return MakeString(value.GetBytes()); default: MKQL_ENSURE_S(oneOfCase == NKikimrMiniKQL::TValue::ValueValueCase::kBytes, "got: " << (int) oneOfCase << ", type: " << (int) dataType->GetSchemeType()); @@ -877,6 +892,7 @@ void ExportPrimitiveTypeToProto(ui32 schemeType, Ydb::Type& output) { case NYql::NProto::TypeIds::Yson: case NYql::NProto::TypeIds::Json: case NYql::NProto::TypeIds::Uuid: + case NYql::NProto::TypeIds::Rowid: case NYql::NProto::TypeIds::JsonDocument: case NYql::NProto::TypeIds::DyNumber: case NYql::NProto::TypeIds::Date32: @@ -1282,6 +1298,10 @@ TNode* TProtoImporter::ImportNodeFromProto(TType* type, const NKikimrMiniKQL::TV dataNode = TDataLiteral::Create(env.NewStringValue(NUdf::TStringRef(buf.bytes, 16)), dataType, env); break; } + case NUdf::TDataType::Id: + MKQL_ENSURE(value.GetBytes().size() == NRowid::ROWID_LEN, "Invalid Rowid size"); + dataNode = TDataLiteral::Create(env.NewStringValue(NUdf::TStringRef(value.GetBytes().data(), value.GetBytes().size())), dataType, env); + break; default: MKQL_ENSURE(false, TStringBuilder() << "Unknown data type: " << schemeType); } @@ -1643,6 +1663,14 @@ Y_FORCE_INLINE NUdf::TUnboxedValue KindDataImport(const TType* type, const Ydb:: buf.half[1] = value.high_128(); return MakeString(NUdf::TStringRef(buf.bytes, 16)); } + case NUdf::TDataType::Id: { + CheckTypeId(value.value_case(), Ydb::Value::kBytesValue, "Rowid"); + const auto& stringRef = value.bytes_value(); + if (stringRef.size() != NRowid::ROWID_LEN) { + throw yexception() << "Invalid Rowid value"; + } + return MakeString(TStringBuf(stringRef.data(), stringRef.size())); + } case NUdf::TDataType::Id: { CheckTypeId(value.value_case(), Ydb::Value::kTextValue, "JsonDocument"); const auto binaryJson = NBinaryJson::SerializeToBinaryJson(value.text_value()); diff --git a/ydb/library/yql/dq/runtime/dq_arrow_helpers.cpp b/ydb/library/yql/dq/runtime/dq_arrow_helpers.cpp index 962023393bcdf..85b5d4e220975 100644 --- a/ydb/library/yql/dq/runtime/dq_arrow_helpers.cpp +++ b/ydb/library/yql/dq/runtime/dq_arrow_helpers.cpp @@ -92,6 +92,7 @@ bool SwitchMiniKQLDataTypeToArrowType(NUdf::EDataSlot type, TFunc&& callback) { return callback(TTypeWrapper()); case NUdf::EDataSlot::Decimal: case NUdf::EDataSlot::Uuid: + case NUdf::EDataSlot::Rowid: return callback(TTypeWrapper()); case NUdf::EDataSlot::TzDate: case NUdf::EDataSlot::TzDatetime: @@ -220,7 +221,9 @@ std::shared_ptr CreateEmptyArrowImpl(NUdf::EDataSlot slot) { template <> std::shared_ptr CreateEmptyArrowImpl(NUdf::EDataSlot slot) { - Y_UNUSED(slot); + if (slot == NUdf::EDataSlot::Rowid) { + return arrow::fixed_size_binary(NUdf::ROWID_SIZE); + } return arrow::fixed_size_binary(NScheme::FSB_SIZE); } @@ -534,7 +537,7 @@ void AppendFixedSizeDataValue(arrow::ArrayBuilder* builder, NUdf::TUnboxedValue if (!value.HasValue()) { status = typedBuilder->AppendNull(); } else { - if (dataSlot == NUdf::EDataSlot::Uuid) { + if (dataSlot == NUdf::EDataSlot::Uuid || dataSlot == NUdf::EDataSlot::Rowid) { auto data = value.AsStringRef(); status = typedBuilder->Append(data.Data()); } else if (dataSlot == NUdf::EDataSlot::Decimal) { diff --git a/ydb/library/yql/dq/runtime/dq_output_consumer.cpp b/ydb/library/yql/dq/runtime/dq_output_consumer.cpp index 5b2f537780aac..5c8d06a144aac 100644 --- a/ydb/library/yql/dq/runtime/dq_output_consumer.cpp +++ b/ydb/library/yql/dq/runtime/dq_output_consumer.cpp @@ -247,6 +247,7 @@ struct TColumnShardHashV1 { break; } case NYql::NProto::Uuid: + case NYql::NProto::Rowid: case NYql::NProto::String: case NYql::NProto::Utf8: { auto value = uv.AsStringRef(); diff --git a/ydb/library/yql/dq/runtime/dq_transport.cpp b/ydb/library/yql/dq/runtime/dq_transport.cpp index 64a796233063b..654f0c7b16a3c 100644 --- a/ydb/library/yql/dq/runtime/dq_transport.cpp +++ b/ydb/library/yql/dq/runtime/dq_transport.cpp @@ -220,6 +220,8 @@ std::optional EstimateIntegralDataSize(const TDataType* dataType) { case NUdf::EDataSlot::Uuid: case NUdf::EDataSlot::Decimal: return 16; + case NUdf::EDataSlot::Rowid: + return NUdf::ROWID_SIZE; case NUdf::EDataSlot::String: case NUdf::EDataSlot::Utf8: case NUdf::EDataSlot::DyNumber: diff --git a/ydb/library/yql/dq/type_ann/dq_type_ann.cpp b/ydb/library/yql/dq/type_ann/dq_type_ann.cpp index fc26e3ec42c61..71e744466433d 100644 --- a/ydb/library/yql/dq/type_ann/dq_type_ann.cpp +++ b/ydb/library/yql/dq/type_ann/dq_type_ann.cpp @@ -1610,6 +1610,7 @@ bool IsTypeSupportedInMergeCn(EDataSlot type) { case EDataSlot::String: case EDataSlot::Utf8: case EDataSlot::Uuid: + case EDataSlot::Rowid: case EDataSlot::Date: case EDataSlot::Datetime: case EDataSlot::Timestamp: diff --git a/ydb/library/yql/udfs/common/rowid/rowid.cpp b/ydb/library/yql/udfs/common/rowid/rowid.cpp new file mode 100644 index 0000000000000..8271de6b84042 --- /dev/null +++ b/ydb/library/yql/udfs/common/rowid/rowid.cpp @@ -0,0 +1,457 @@ +#include "rowid_keygen.h" + +#include +#include + +#include + +#include + +// Rowid UDF: key-friendly Rowid generators. +// +// Returned values are raw 14-byte Rowid representation. +// Primary-key helpers (layouts from the pk_generation RFC): +// - newRowKey: shard spread via 12-bit prefix + time locality within a prefix; +// - newColumnKey: chronological clustering by creation time (seconds); +// - newRowGroup: batch of row keys sharing a common prefix (Uint64 or Rowid). +// +// Optional dependency arguments [T1, ...] work like RandomUuid(): they control +// when the function is evaluated per row, not the value contents. + +using namespace NYql; +using namespace NYql::NUdf; + +namespace { + +constexpr ui32 MaxDepArgs = 32; + +enum class EPrefixArgType { + None, + Uint64, + Rowid, +}; + +TString BuildDepArgKindsPredicate(TStringBuf argName) { + return TStringBuilder() << R"( +{cmd=or;value=[ + {cmd=kind;arg=)" << argName << R"(;value=Data}; + {cmd=kind;arg=)" << argName << R"(;value=Optional}; + {cmd=kind;arg=)" << argName << R"(;value=Tuple}; + {cmd=kind;arg=)" << argName << R"(;value=Struct}; + {cmd=kind;arg=)" << argName << R"(;value=List}; + {cmd=kind;arg=)" << argName << R"(;value=Dict}; + {cmd=kind;arg=)" << argName << R"(;value=Stream}; + {cmd=kind;arg=)" << argName << R"(;value=Null}; + {cmd=kind;arg=)" << argName << R"(;value=Void} +]} +)"; +} + +TString BuildAndDepArgKindsPredicate(ui32 depCount, ui32 firstArgIndex = 0) { + Y_ENSURE(depCount > 0); + TStringBuilder sb; + sb << "{cmd=and;value=["; + for (ui32 i = 0; i < depCount; ++i) { + if (i > 0) { + sb << ";"; + } + sb << BuildDepArgKindsPredicate(TStringBuilder() << "T" << (firstArgIndex + i)); + } + sb << "]}"; + return sb; +} + +TString BuildCallableTypeWithUniversalDeps(ui32 depCount, EPrefixArgType prefixArg) { + TStringBuilder sb; + sb << "[CallableType;[];[];["; + if (prefixArg != EPrefixArgType::None) { + const TStringBuf prefixTypeName = prefixArg == EPrefixArgType::Rowid ? "Rowid" : "Uint64"; + sb << "[[DataType;" << prefixTypeName << "]"; + for (ui32 i = 0; i < depCount; ++i) { + sb << ";[UniversalType]"; + } + sb << ";[[DataType;Rowid]]]"; + } else { + for (ui32 i = 0; i < depCount; ++i) { + sb << "[UniversalType]"; + if (i + 1 < depCount) { + sb << ";"; + } + } + if (depCount > 0) { + sb << ";"; + } + sb << "[[DataType;Rowid]]]"; + } + sb << "]]"; + return sb; +} + +TString BuildCallableTypeRowGroup(ui32 depCount, EPrefixArgType prefixArg) { + Y_ENSURE(prefixArg != EPrefixArgType::None); + const TStringBuf prefixTypeName = prefixArg == EPrefixArgType::Rowid ? "Rowid" : "Uint64"; + TStringBuilder sb; + sb << "[CallableType;[];[];[[[DataType;" << prefixTypeName << "];[DataType;Uint64]"; + for (ui32 i = 0; i < depCount; ++i) { + sb << ";[UniversalType]"; + } + sb << ";[[ListType;[DataType;Rowid]]]]]"; + return sb; +} + +void AppendNoPrefixPolyArgRule(TStringBuilder& sb, ui32 depCount) { + sb << "["; + if (depCount == 0) { + sb << "[]"; + } else { + sb << BuildAndDepArgKindsPredicate(depCount); + } + sb << "; {type=" << BuildCallableTypeWithUniversalDeps(depCount, EPrefixArgType::None) << "}]"; +} + +void AppendRowGroupPolyArgRule(TStringBuilder& sb, ui32 depCount, EPrefixArgType prefixArg) { + Y_ENSURE(prefixArg != EPrefixArgType::None); + const TStringBuf prefixTypeName = prefixArg == EPrefixArgType::Rowid ? "Rowid" : "Uint64"; + + sb << "["; + if (depCount == 0) { + sb << "{cmd=and;value=[" + << "{cmd=type;arg=T0;value=[DataType;" << prefixTypeName << "]};" + << "{cmd=type;arg=T1;value=[DataType;Uint64]}" + << "]}"; + } else { + sb << "{cmd=and;value=[" + << "{cmd=type;arg=T0;value=[DataType;" << prefixTypeName << "]};" + << "{cmd=type;arg=T1;value=[DataType;Uint64]}"; + for (ui32 i = 0; i < depCount; ++i) { + sb << ";" << BuildDepArgKindsPredicate(TStringBuilder() << "T" << (i + 2)); + } + sb << "]}"; + } + sb << "; {type=" << BuildCallableTypeRowGroup(depCount, prefixArg) << "}]"; +} + +TString BuildNoPrefixPolyArgs(TStringBuf errorMessage) { + TStringBuilder sb; + sb << "[["; + bool first = true; + for (ui32 depCount = MaxDepArgs; depCount > 0; --depCount) { + if (!first) { + sb << ";"; + } + first = false; + AppendNoPrefixPolyArgRule(sb, depCount); + } + if (!first) { + sb << ";"; + } + AppendNoPrefixPolyArgRule(sb, 0); + sb << "; [{cmd=error;message=\"" << errorMessage << "\"}; {}]]"; + return sb; +} + +TString BuildRowGroupPolyArgs(TStringBuf errorMessage) { + TStringBuilder sb; + sb << "[["; + bool first = true; + for (ui32 depCount = MaxDepArgs; depCount > 0; --depCount) { + if (!first) { + sb << ";"; + } + first = false; + AppendRowGroupPolyArgRule(sb, depCount, EPrefixArgType::Rowid); + sb << ";"; + AppendRowGroupPolyArgRule(sb, depCount, EPrefixArgType::Uint64); + } + if (!first) { + sb << ";"; + } + AppendRowGroupPolyArgRule(sb, 0, EPrefixArgType::Rowid); + sb << ";"; + AppendRowGroupPolyArgRule(sb, 0, EPrefixArgType::Uint64); + sb << "; [{cmd=error;message=\"" << errorMessage << "\"}; {}]]"; + return sb; +} + +ui64 ReadPrefixArg(const TUnboxedValuePod& arg, bool prefixFromRowid) { + if (prefixFromRowid) { + const auto ref = arg.AsStringRef(); + if (ref.Size() != NRowidKeyGen::RowidLen) { + throw std::runtime_error("Expected Rowid value of 14 bytes"); + } + return NRowidKeyGen::ExtractPrefixFromRowidBytes( + reinterpret_cast(ref.Data())); + } + return arg.Get(); +} + +bool IsRowidArgType(const ITypeInfoHelper1& typeHelper, const TType* argType) { + TDataTypeInspector argInspector(typeHelper, argType); + return argInspector && argInspector.GetTypeId() == NUdf::TDataType::Id; +} + +TUnboxedValue MakeRowidFromBytes( + const IValueBuilder* valueBuilder, + const std::array& bytes) +{ + return valueBuilder->NewString(TStringRef( + reinterpret_cast(bytes.data()), + bytes.size())); +} + +TUnboxedValue MakeRowKeyRowidValue( + const IValueBuilder* valueBuilder, ui64 prefix, bool hasPrefix) +{ + return MakeRowidFromBytes( + valueBuilder, + NRowidKeyGen::MakeRowKeyRowidBytes(prefix, Seconds(), hasPrefix)); +} + +TUnboxedValue MakeColumnKeyRowidValue(const IValueBuilder* valueBuilder) { + return MakeRowidFromBytes( + valueBuilder, + NRowidKeyGen::MakeColumnKeyRowidBytes(Seconds())); +} + +enum class EKeyKind { + RowKey, + ColumnKey, +}; + +template +class TNewRowid: public TBoxedValue { +public: + using TTypeAwareMarker = bool; + + explicit TNewRowid(TSourcePosition pos) + : Pos_(pos) + { + } + + static const TStringRef& Name() { + if constexpr (Kind == EKeyKind::RowKey) { + static auto name = TStringRef::Of("newRowKey"); + return name; + } else { + static auto name = TStringRef::Of("newColumnKey"); + return name; + } + } + + static bool DeclareSignature( + const TStringRef& name, + TType* userType, + IFunctionTypeInfoBuilder& builder, + bool typesOnly) + { + if (Name() != name) { + return false; + } + + if (!userType) { + builder.SetError("Missing user type."); + return true; + } + + builder.UserType(userType); + const auto typeHelper = builder.TypeInfoHelper(); + const auto userTypeInspector = TTupleTypeInspector(*typeHelper, userType); + if (!userTypeInspector || userTypeInspector.GetElementsCount() < 1) { + builder.SetError("Invalid user type."); + return true; + } + + const auto argsTypeTuple = userTypeInspector.GetElementType(0); + const auto argsTypeInspector = TTupleTypeInspector(*typeHelper, argsTypeTuple); + if (!argsTypeInspector) { + builder.SetError("Invalid user type - expected tuple."); + return true; + } + + const ui32 argsCount = argsTypeInspector.GetElementsCount(); + if (argsCount > MaxDepArgs) { + builder.SetError(TStringBuilder() << "Too many dependency arguments: " << argsCount); + return true; + } + + auto args = builder.Args(argsCount); + for (ui32 i = 0; i < argsCount; ++i) { + args->Add(argsTypeInspector.GetElementType(i)); + } + args->Done(); + builder.Returns(); + + if (!typesOnly) { + builder.Implementation(new TNewRowid(builder.GetSourcePosition())); + } + return true; + } + +private: + TUnboxedValue Run( + const IValueBuilder* valueBuilder, + const TUnboxedValuePod* args) const final try + { + Y_UNUSED(args); + if constexpr (Kind == EKeyKind::RowKey) { + return MakeRowKeyRowidValue(valueBuilder, 0, false); + } else { + return MakeColumnKeyRowidValue(valueBuilder); + } + } catch (const std::exception& e) { + UdfTerminate((TStringBuilder() << Pos_ << " " << e.what()).c_str()); + } + + TSourcePosition Pos_; +}; + +class TNewRowGroup: public TBoxedValue { +public: + using TTypeAwareMarker = bool; + + TNewRowGroup(TSourcePosition pos, bool prefixFromRowid) + : Pos_(pos) + , PrefixFromRowid_(prefixFromRowid) + { + } + + static const TStringRef& Name() { + static auto name = TStringRef::Of("newRowGroup"); + return name; + } + + static bool DeclareSignature( + const TStringRef& name, + TType* userType, + IFunctionTypeInfoBuilder& builder, + bool typesOnly) + { + if (Name() != name) { + return false; + } + + if (!userType) { + builder.SetError("Missing user type."); + return true; + } + + builder.UserType(userType); + const auto typeHelper = builder.TypeInfoHelper(); + const auto userTypeInspector = TTupleTypeInspector(*typeHelper, userType); + if (!userTypeInspector || userTypeInspector.GetElementsCount() < 1) { + builder.SetError("Invalid user type."); + return true; + } + + const auto argsTypeTuple = userTypeInspector.GetElementType(0); + const auto argsTypeInspector = TTupleTypeInspector(*typeHelper, argsTypeTuple); + if (!argsTypeInspector || argsTypeInspector.GetElementsCount() < 2) { + builder.SetError("newRowGroup requires prefix and count arguments."); + return true; + } + + const ui32 argsCount = argsTypeInspector.GetElementsCount(); + if (argsCount > 2 + MaxDepArgs) { + builder.SetError(TStringBuilder() << "Too many dependency arguments: " << (argsCount - 2)); + return true; + } + + const bool prefixFromRowid = IsRowidArgType(*typeHelper, argsTypeInspector.GetElementType(0)); + auto args = builder.Args(argsCount); + for (ui32 i = 0; i < argsCount; ++i) { + args->Add(argsTypeInspector.GetElementType(i)); + } + args->Done(); + builder.Returns(builder.List()->Item().Build()); + + if (!typesOnly) { + builder.Implementation(new TNewRowGroup(builder.GetSourcePosition(), prefixFromRowid)); + } + return true; + } + +private: + TUnboxedValue Run( + const IValueBuilder* valueBuilder, + const TUnboxedValuePod* args) const final try + { + const ui64 prefix = ReadPrefixArg(args[0], PrefixFromRowid_); + const ui64 count = args[1].Get(); + if (count > NRowidKeyGen::MaxRowGroupCount) { + throw std::runtime_error(TStringBuilder() + << "newRowGroup count exceeds limit " << NRowidKeyGen::MaxRowGroupCount); + } + + std::vector items; + items.reserve(count); + const ui64 epochSeconds = Seconds(); + for (ui64 i = 0; i < count; ++i) { + items.push_back(MakeRowidFromBytes( + valueBuilder, + NRowidKeyGen::MakeRowKeyRowidBytes(prefix, epochSeconds, true))); + } + return valueBuilder->NewList(items.data(), items.size()); + } catch (const std::exception& e) { + UdfTerminate((TStringBuilder() << Pos_ << " " << e.what()).c_str()); + } + + TSourcePosition Pos_; + bool PrefixFromRowid_ = false; +}; + +class TRowidModule: public IUdfModule { +public: + TStringRef Name() const { + return TStringRef::Of("Rowid"); + } + + void CleanupOnTerminate() const final { + } + + void GetAllFunctions(IFunctionsSink& sink) const final { + static const TString newRowKeyPolyArgs = BuildNoPrefixPolyArgs("Unexpected arguments for Rowid::newRowKey"); + static const TString newColumnKeyPolyArgs = BuildNoPrefixPolyArgs("Unexpected arguments for Rowid::newColumnKey"); + static const TString newRowGroupPolyArgs = BuildRowGroupPolyArgs("Unexpected arguments for Rowid::newRowGroup"); + + auto newRowKey = sink.Add(TNewRowid::Name()); + newRowKey->SetTypeAwareness(); + newRowKey->SetPolyArgs(TStringRef(newRowKeyPolyArgs)); + + auto newColumnKey = sink.Add(TNewRowid::Name()); + newColumnKey->SetTypeAwareness(); + newColumnKey->SetPolyArgs(TStringRef(newColumnKeyPolyArgs)); + + auto newRowGroup = sink.Add(TNewRowGroup::Name()); + newRowGroup->SetTypeAwareness(); + newRowGroup->SetPolyArgs(TStringRef(newRowGroupPolyArgs)); + } + + void BuildFunctionTypeInfo( + const TStringRef& name, + TType* userType, + const TStringRef& typeConfig, + ui32 flags, + IFunctionTypeInfoBuilder& builder) const override + { + Y_UNUSED(typeConfig); + try { + const bool typesOnly = (flags & TFlags::TypesOnly); + if (TNewRowid::DeclareSignature(name, userType, builder, typesOnly)) { + return; + } + if (TNewRowid::DeclareSignature(name, userType, builder, typesOnly)) { + return; + } + if (TNewRowGroup::DeclareSignature(name, userType, builder, typesOnly)) { + return; + } + ythrow yexception() << "Unknown function name: " << TStringBuf(name); + } catch (const std::exception& e) { + builder.SetError(CurrentExceptionMessage()); + } + } +}; + +} // namespace + +REGISTER_MODULES(TRowidModule) diff --git a/ydb/library/yql/udfs/common/rowid/rowid_keygen.h b/ydb/library/yql/udfs/common/rowid/rowid_keygen.h new file mode 100644 index 0000000000000..7dd1e7b45bda8 --- /dev/null +++ b/ydb/library/yql/udfs/common/rowid/rowid_keygen.h @@ -0,0 +1,128 @@ +#pragma once + +#include + +#include +#include + +#include +#include +#include + +// Rowid generators for YDB primary keys. +// +// Rowid is a 14-byte opaque value. Internal and external byte orders coincide +// (unlike Uuid, which uses Microsoft GUID mixed-endian layout in YDB). +// +// Layouts follow the pk_generation RFC (same bit fields as UUIDv8 row/column +// keys, but without version/variant nibbles — the random suffix is shorter so +// the whole value fits in 14 bytes). + +namespace NYql::NRowidKeyGen { + +static constexpr ui32 RowidLen = NKikimr::NRowid::ROWID_LEN; + +// Row-key layout: +// [12 prefix][31 timestamp sec][5 random][64 random] +static constexpr ui32 PrefixBits = 12; +static constexpr ui32 TimestampBits = 31; +static constexpr ui64 TimestampModulus = 1ULL << TimestampBits; +static constexpr ui64 PrefixMsbMask = ((1ULL << PrefixBits) - 1) << (64 - PrefixBits); +static constexpr ui64 PrefixParamMask = (1ULL << PrefixBits) - 1; +static constexpr ui32 RowKeyTimestampShift = 64 - PrefixBits - TimestampBits; +static constexpr ui64 RowKeyTimestampMask = ((1ULL << TimestampBits) - 1) << RowKeyTimestampShift; + +// Column-key layout: +// [31 timestamp sec][1 random][80 random] +static constexpr ui32 ColumnKeyTimestampShift = 64 - TimestampBits; +static constexpr ui64 ColumnKeyTimestampMask = ((1ULL << TimestampBits) - 1) << ColumnKeyTimestampShift; + +static constexpr ui64 MaxRowGroupCount = 1'000'000; + +inline ui64 ReadBe64(const ui8* data) { + ui64 value = 0; + for (ui32 i = 0; i < 8; ++i) { + value = (value << 8) | data[i]; + } + return value; +} + +inline void WriteBe64(ui64 value, ui8* data) { + for (int i = 7; i >= 0; --i) { + data[i] = static_cast(value & 0xff); + value >>= 8; + } +} + +inline void FillRandomBytes(ui8* data, size_t size) { + for (size_t offset = 0; offset < size; offset += sizeof(ui64)) { + const ui64 random = RandomNumber(); + std::memcpy(data + offset, &random, std::min(size - offset, sizeof(ui64))); + } +} + +inline ui64 PrefixParamToMsb(ui64 prefix) { + return (prefix & PrefixParamMask) << (64 - PrefixBits); +} + +inline ui64 ExtractPrefixFromRowidBytes(const ui8* data) { + const ui64 msb = ReadBe64(data); + return (msb & PrefixMsbMask) >> (64 - PrefixBits); +} + +inline ui64 GetRowKeyTimestampCode(ui64 epochSeconds) { + return (epochSeconds % TimestampModulus) << RowKeyTimestampShift; +} + +inline ui64 GetColumnKeyTimestampCode(ui64 epochSeconds) { + return (epochSeconds % TimestampModulus) << ColumnKeyTimestampShift; +} + +inline ui64 UpdateMsbRowKey(ui64 msb, ui64 prefix, ui64 epochSeconds, bool hasPrefix) { + const ui64 tsCode = GetRowKeyTimestampCode(epochSeconds); + if (hasPrefix) { + return (msb & ~(PrefixMsbMask | RowKeyTimestampMask)) + | (PrefixParamToMsb(prefix) | (tsCode & RowKeyTimestampMask)); + } + return (msb & ~RowKeyTimestampMask) | (tsCode & RowKeyTimestampMask); +} + +inline ui64 UpdateMsbColumnKey(ui64 msb, ui64 epochSeconds) { + const ui64 tsCode = GetColumnKeyTimestampCode(epochSeconds); + return (msb & ~ColumnKeyTimestampMask) | (tsCode & ColumnKeyTimestampMask); +} + +// Build a row-table Rowid key. +// +// Sort order (memcmp): (1) 12-bit random prefix; (2) 31-bit second-granularity +// timestamp; (3) random suffix. Without an explicit prefix, prefix bits stay +// random. With hasPrefix=true, the prefix is fixed (used by newRowGroup). +inline std::array MakeRowKeyRowidBytes( + ui64 prefix, ui64 epochSeconds, bool hasPrefix) +{ + std::array result{}; + FillRandomBytes(result.data(), result.size()); + + ui64 msb = ReadBe64(result.data()); + msb = UpdateMsbRowKey(msb, prefix, epochSeconds, hasPrefix); + WriteBe64(msb, result.data()); + + return result; +} + +// Build a column-table Rowid key. +// +// Sort order (memcmp): 31-bit second-granularity timestamp first, then random +// suffix. No partition prefix — column tables use hash partitioning. +inline std::array MakeColumnKeyRowidBytes(ui64 epochSeconds) { + std::array result{}; + FillRandomBytes(result.data(), result.size()); + + ui64 msb = ReadBe64(result.data()); + msb = UpdateMsbColumnKey(msb, epochSeconds); + WriteBe64(msb, result.data()); + + return result; +} + +} // namespace NYql::NRowidKeyGen diff --git a/ydb/library/yql/udfs/common/rowid/test/canondata/result.json b/ydb/library/yql/udfs/common/rowid/test/canondata/result.json new file mode 100644 index 0000000000000..94347d0d43394 --- /dev/null +++ b/ydb/library/yql/udfs/common/rowid/test/canondata/result.json @@ -0,0 +1,7 @@ +{ + "test.test[new_rowid]": [ + { + "uri": "file://test.test_new_rowid_/results.txt" + } + ] +} diff --git a/ydb/library/yql/udfs/common/rowid/test/canondata/test.test_new_rowid_/results.txt b/ydb/library/yql/udfs/common/rowid/test/canondata/test.test_new_rowid_/results.txt new file mode 100644 index 0000000000000..ccb512cd1afe0 --- /dev/null +++ b/ydb/library/yql/udfs/common/rowid/test/canondata/test.test_new_rowid_/results.txt @@ -0,0 +1,366 @@ +[ + { + "Write" = [ + { + "Type" = [ + "ListType"; + [ + "StructType"; + [ + [ + "column_key_unique"; + [ + "DataType"; + "Bool" + ] + ] + ] + ] + ]; + "Data" = [ + [ + %false + ] + ] + } + ] + }; + { + "Write" = [ + { + "Type" = [ + "ListType"; + [ + "StructType"; + [ + [ + "row_key_unique"; + [ + "DataType"; + "Bool" + ] + ] + ] + ] + ]; + "Data" = [ + [ + %false + ] + ] + } + ] + }; + { + "Write" = [ + { + "Type" = [ + "ListType"; + [ + "StructType"; + [ + [ + "column_key_dep_unique"; + [ + "DataType"; + "Bool" + ] + ] + ] + ] + ]; + "Data" = [ + [ + %true + ] + ] + } + ] + }; + { + "Write" = [ + { + "Type" = [ + "ListType"; + [ + "StructType"; + [ + [ + "row_key_dep_unique"; + [ + "DataType"; + "Bool" + ] + ] + ] + ] + ]; + "Data" = [ + [ + %true + ] + ] + } + ] + }; + { + "Write" = [ + { + "Type" = [ + "ListType"; + [ + "StructType"; + [ + [ + "column_key_three_dep_unique"; + [ + "DataType"; + "Bool" + ] + ] + ] + ] + ]; + "Data" = [ + [ + %true + ] + ] + } + ] + }; + { + "Write" = [ + { + "Type" = [ + "ListType"; + [ + "StructType"; + [ + [ + "row_key_three_dep_unique"; + [ + "DataType"; + "Bool" + ] + ] + ] + ] + ]; + "Data" = [ + [ + %true + ] + ] + } + ] + }; + { + "Write" = [ + { + "Type" = [ + "ListType"; + [ + "StructType"; + [ + [ + "row_group_count"; + [ + "DataType"; + "Bool" + ] + ] + ] + ] + ]; + "Data" = [ + [ + %true + ] + ] + } + ] + }; + { + "Write" = [ + { + "Type" = [ + "ListType"; + [ + "StructType"; + [ + [ + "row_group_distinct"; + [ + "DataType"; + "Bool" + ] + ] + ] + ] + ]; + "Data" = [ + [ + %true + ] + ] + } + ] + }; + { + "Write" = [ + { + "Type" = [ + "ListType"; + [ + "StructType"; + [ + [ + "row_group_rowid_prefix_count"; + [ + "DataType"; + "Bool" + ] + ] + ] + ] + ]; + "Data" = [ + [ + %true + ] + ] + } + ] + }; + { + "Write" = [ + { + "Type" = [ + "ListType"; + [ + "StructType"; + [ + [ + "row_group_rowid_prefix_distinct"; + [ + "DataType"; + "Bool" + ] + ] + ] + ] + ]; + "Data" = [ + [ + %true + ] + ] + } + ] + }; + { + "Write" = [ + { + "Type" = [ + "ListType"; + [ + "StructType"; + [ + [ + "row_group_dep_unique"; + [ + "DataType"; + "Bool" + ] + ] + ] + ] + ]; + "Data" = [ + [ + %true + ] + ] + } + ] + }; + { + "Write" = [ + { + "Type" = [ + "ListType"; + [ + "StructType"; + [ + [ + "column_key_base64_length"; + [ + "DataType"; + "Bool" + ] + ] + ] + ] + ]; + "Data" = [ + [ + %true + ] + ] + } + ] + }; + { + "Write" = [ + { + "Type" = [ + "ListType"; + [ + "StructType"; + [ + [ + "row_key_base64_length"; + [ + "DataType"; + "Bool" + ] + ] + ] + ] + ]; + "Data" = [ + [ + %true + ] + ] + } + ] + }; + { + "Write" = [ + { + "Type" = [ + "ListType"; + [ + "StructType"; + [ + [ + "prefix_is_uint64"; + [ + "DataType"; + "Bool" + ] + ] + ] + ] + ]; + "Data" = [ + [ + %true + ] + ] + } + ] + } +] \ No newline at end of file diff --git a/ydb/library/yql/udfs/common/rowid/test/cases/new_rowid.sql b/ydb/library/yql/udfs/common/rowid/test/cases/new_rowid.sql new file mode 100644 index 0000000000000..1a3fa4fe892e8 --- /dev/null +++ b/ydb/library/yql/udfs/common/rowid/test/cases/new_rowid.sql @@ -0,0 +1,26 @@ +$p = RandomNumber(1); + +SELECT Rowid::newColumnKey() != Rowid::newColumnKey() AS column_key_unique; +SELECT Rowid::newRowKey() != Rowid::newRowKey() AS row_key_unique; + +SELECT Rowid::newColumnKey(1) != Rowid::newColumnKey(2) AS column_key_dep_unique; +SELECT Rowid::newRowKey(1) != Rowid::newRowKey(2) AS row_key_dep_unique; +SELECT Rowid::newColumnKey(1, 2, 3) != Rowid::newColumnKey(1, 2, 4) AS column_key_three_dep_unique; +SELECT Rowid::newRowKey(1, 2, 3) != Rowid::newRowKey(1, 2, 4) AS row_key_three_dep_unique; + +$group = Rowid::newRowGroup($p, 3ul); +SELECT ListLength($group) = 3ul AS row_group_count; +SELECT Unwrap($group[0]) != Unwrap($group[1]) AND Unwrap($group[1]) != Unwrap($group[2]) AS row_group_distinct; + +$groupFromRowid = Rowid::newRowGroup(Unwrap($group[0]), 2ul); +SELECT ListLength($groupFromRowid) = 2ul AS row_group_rowid_prefix_count; +SELECT Unwrap($groupFromRowid[0]) != Unwrap($groupFromRowid[1]) AS row_group_rowid_prefix_distinct; + +$groupDep = Rowid::newRowGroup($p, 2ul, 1); +$groupDep2 = Rowid::newRowGroup($p, 2ul, 2); +SELECT Unwrap($groupDep[0]) != Unwrap($groupDep2[0]) AS row_group_dep_unique; + +SELECT Length(CAST(Rowid::newColumnKey() AS String)) = 19ul AS column_key_base64_length; +SELECT Length(CAST(Rowid::newRowKey() AS String)) = 19ul AS row_key_base64_length; + +SELECT $p != 0ul OR $p == 0ul AS prefix_is_uint64; diff --git a/ydb/library/yql/udfs/common/rowid/test/ya.make b/ydb/library/yql/udfs/common/rowid/test/ya.make new file mode 100644 index 0000000000000..5824bb263585e --- /dev/null +++ b/ydb/library/yql/udfs/common/rowid/test/ya.make @@ -0,0 +1,11 @@ +YQL_UDF_TEST() + +DEPENDS(ydb/library/yql/udfs/common/rowid) + +SIZE(MEDIUM) + +IF (SANITIZER_TYPE == "memory") + TAG(ya:not_autocheck) # YQL-15385 +ENDIF() + +END() diff --git a/ydb/library/yql/udfs/common/rowid/ut/rowid_sort_order_ut.cpp b/ydb/library/yql/udfs/common/rowid/ut/rowid_sort_order_ut.cpp new file mode 100644 index 0000000000000..d1643243e4bb4 --- /dev/null +++ b/ydb/library/yql/udfs/common/rowid/ut/rowid_sort_order_ut.cpp @@ -0,0 +1,150 @@ +#include + +#include + +#include + +#include +#include + +using namespace NYql::NRowidKeyGen; + +namespace { +constexpr ui64 kTestRowPrefix = 0xA5A; // 12 bits + +std::array MakeColumnKey(ui64 epochSeconds) { + return MakeColumnKeyRowidBytes(epochSeconds); +} + +std::array MakeRowKey(ui64 prefix, ui64 epochSeconds) { + return MakeRowKeyRowidBytes(prefix, epochSeconds, true); +} + +ui64 ExtractRowTimestamp(const std::array& bytes) { + const ui64 msb = ReadBe64(bytes.data()); + return (msb & RowKeyTimestampMask) >> RowKeyTimestampShift; +} + +ui64 ExtractColumnTimestamp(const std::array& bytes) { + const ui64 msb = ReadBe64(bytes.data()); + return (msb & ColumnKeyTimestampMask) >> ColumnKeyTimestampShift; +} +} // namespace + +Y_UNIT_TEST_SUITE(TRowidKeyGenSortOrder) { + Y_UNIT_TEST(RowKeyUsesBottom12PrefixBits) { + SetRandomSeed(1); + const ui64 epochSeconds = 1'700'000'000ull; + const ui64 rawPrefix = 0x12345ull; + const ui64 expectedParam = rawPrefix & PrefixParamMask; + + const auto fromRaw = MakeRowKeyRowidBytes(rawPrefix, epochSeconds, true); + SetRandomSeed(1); + const auto fromBottomBits = MakeRowKeyRowidBytes(expectedParam, epochSeconds, true); + UNIT_ASSERT_EQUAL(fromRaw, fromBottomBits); + + constexpr ui64 kSmallPrefix = 7; + SetRandomSeed(2); + const auto withSmallPrefix = MakeRowKeyRowidBytes(kSmallPrefix, epochSeconds, true); + SetRandomSeed(2); + const auto withZeroPrefix = MakeRowKeyRowidBytes(0, epochSeconds, true); + UNIT_ASSERT(withSmallPrefix != withZeroPrefix); + UNIT_ASSERT_EQUAL(ExtractPrefixFromRowidBytes(withSmallPrefix.data()), kSmallPrefix); + } + + Y_UNIT_TEST(RowKeyEmbedsTimestamp) { + SetRandomSeed(3); + const ui64 epochSeconds = 1'700'000'123ull; + const auto bytes = MakeRowKey(kTestRowPrefix, epochSeconds); + UNIT_ASSERT_EQUAL(ExtractRowTimestamp(bytes), epochSeconds % TimestampModulus); + UNIT_ASSERT_EQUAL(ExtractPrefixFromRowidBytes(bytes.data()), kTestRowPrefix); + } + + Y_UNIT_TEST(ColumnKeyEmbedsTimestamp) { + SetRandomSeed(4); + const ui64 epochSeconds = 1'700'000'456ull; + const auto bytes = MakeColumnKey(epochSeconds); + UNIT_ASSERT_EQUAL(ExtractColumnTimestamp(bytes), epochSeconds % TimestampModulus); + } + + Y_UNIT_TEST(ColumnKeysSortByTimestamp) { + SetRandomSeed(5); + const ui64 earlier = 1'700'000'000ull; + const ui64 later = earlier + 10; + const auto earlierGenerated = MakeColumnKey(earlier); + const auto laterGenerated = MakeColumnKey(later); + UNIT_ASSERT(std::memcmp(earlierGenerated.data(), laterGenerated.data(), RowidLen) < 0); + } + + Y_UNIT_TEST(RowKeysWithSamePrefixSortByTimestamp) { + SetRandomSeed(6); + const ui64 earlier = 1'700'000'000ull; + const ui64 later = earlier + 10; + const auto earlierGenerated = MakeRowKey(kTestRowPrefix, earlier); + const auto laterGenerated = MakeRowKey(kTestRowPrefix, later); + UNIT_ASSERT(std::memcmp(earlierGenerated.data(), laterGenerated.data(), RowidLen) < 0); + } + + Y_UNIT_TEST(ColumnKeySequenceIsSorted) { + SetRandomSeed(7); + const ui64 baseEpochSeconds = 1'700'000'000ull; + std::vector> generated; + for (ui64 i = 0; i < 16; ++i) { + generated.push_back(MakeColumnKey(baseEpochSeconds + i)); + } + auto sorted = generated; + std::sort(sorted.begin(), sorted.end()); + UNIT_ASSERT_EQUAL(generated, sorted); + } + + Y_UNIT_TEST(RowKeySequenceWithFixedPrefixIsSorted) { + SetRandomSeed(8); + const ui64 baseEpochSeconds = 1'700'000'000ull; + std::vector> generated; + for (ui64 i = 0; i < 16; ++i) { + generated.push_back(MakeRowKey(kTestRowPrefix, baseEpochSeconds + i)); + } + auto sorted = generated; + std::sort(sorted.begin(), sorted.end()); + UNIT_ASSERT_EQUAL(generated, sorted); + } + + Y_UNIT_TEST(SameTimestampColumnKeysAreDistinct) { + SetRandomSeed(9); + const ui64 epochSeconds = 1'700'000'000ull; + std::set> unique; + for (int i = 0; i < 32; ++i) { + unique.insert(MakeColumnKey(epochSeconds)); + } + UNIT_ASSERT_EQUAL(unique.size(), 32u); + } + + Y_UNIT_TEST(SameTimestampRowKeysAreDistinct) { + SetRandomSeed(10); + const ui64 epochSeconds = 1'700'000'000ull; + std::set> unique; + for (int i = 0; i < 32; ++i) { + unique.insert(MakeRowKey(kTestRowPrefix, epochSeconds)); + } + UNIT_ASSERT_EQUAL(unique.size(), 32u); + } + + Y_UNIT_TEST(RowKeyWithoutPrefixUsesRandomPrefixBits) { + SetRandomSeed(11); + const ui64 epochSeconds = 1'700'000'000ull; + std::set prefixes; + for (int i = 0; i < 32; ++i) { + const auto bytes = MakeRowKeyRowidBytes(0, epochSeconds, false); + prefixes.insert(ExtractPrefixFromRowidBytes(bytes.data())); + } + UNIT_ASSERT_GT(prefixes.size(), 1u); + } + + Y_UNIT_TEST(RowidLengthIsFourteen) { + SetRandomSeed(12); + const auto columnKey = MakeColumnKey(Seconds()); + const auto rowKey = MakeRowKeyRowidBytes(0, Seconds(), false); + UNIT_ASSERT_EQUAL(columnKey.size(), 14u); + UNIT_ASSERT_EQUAL(rowKey.size(), 14u); + } +} diff --git a/ydb/library/yql/udfs/common/rowid/ut/ya.make b/ydb/library/yql/udfs/common/rowid/ut/ya.make new file mode 100644 index 0000000000000..31ce4c8971a8e --- /dev/null +++ b/ydb/library/yql/udfs/common/rowid/ut/ya.make @@ -0,0 +1,13 @@ +UNITTEST() + +SIZE(SMALL) + +SRCS( + rowid_sort_order_ut.cpp +) + +PEERDIR( + yql/essentials/types/rowid +) + +END() diff --git a/ydb/library/yql/udfs/common/rowid/ya.make b/ydb/library/yql/udfs/common/rowid/ya.make new file mode 100644 index 0000000000000..122a3de476858 --- /dev/null +++ b/ydb/library/yql/udfs/common/rowid/ya.make @@ -0,0 +1,22 @@ +YQL_UDF_YDB(rowid_udf) + +YQL_ABI_VERSION( + 2 + 46 + 0 +) + +SRCS( + rowid.cpp +) + +PEERDIR( + yql/essentials/types/rowid +) + +END() + +RECURSE_FOR_TESTS( + test + ut +) diff --git a/ydb/library/yql/udfs/common/ya.make b/ydb/library/yql/udfs/common/ya.make index 1c67b042b6433..f0761259fa265 100644 --- a/ydb/library/yql/udfs/common/ya.make +++ b/ydb/library/yql/udfs/common/ya.make @@ -3,5 +3,6 @@ RECURSE( hybrid_search knn roaring + rowid ) diff --git a/ydb/public/api/protos/ydb_value.proto b/ydb/public/api/protos/ydb_value.proto index 0625dde6325ec..cc0e9b227b18b 100644 --- a/ydb/public/api/protos/ydb_value.proto +++ b/ydb/public/api/protos/ydb_value.proto @@ -95,6 +95,7 @@ message Type { JSON = 0x1202; UUID = 0x1203; JSON_DOCUMENT = 0x1204; + ROWID = 0x1205; DYNUMBER = 0x1302; } diff --git a/ydb/public/lib/json_value/ydb_json_value.cpp b/ydb/public/lib/json_value/ydb_json_value.cpp index 2ff2a5ae097b7..a30a01b89f30c 100644 --- a/ydb/public/lib/json_value/ydb_json_value.cpp +++ b/ydb/public/lib/json_value/ydb_json_value.cpp @@ -366,6 +366,9 @@ namespace NYdb { case EPrimitiveType::Uuid: Writer.WriteString(Parser.GetUuid().ToString()); break; + case EPrimitiveType::Rowid: + Writer.WriteString(Parser.GetRowid().ToString()); + break; case EPrimitiveType::JsonDocument: Writer.WriteString(Parser.GetJsonDocument()); break; @@ -791,6 +794,10 @@ namespace { EnsureType(jsonValue, NJson::JSON_STRING); ValueBuilder.Uuid(TUuidValue{jsonValue.GetString()}); break; + case EPrimitiveType::Rowid: + EnsureType(jsonValue, NJson::JSON_STRING); + ValueBuilder.Rowid(TRowidValue{jsonValue.GetString()}); + break; case EPrimitiveType::JsonDocument: EnsureType(jsonValue, NJson::JSON_STRING); ValueBuilder.JsonDocument(jsonValue.GetString()); diff --git a/ydb/public/lib/scheme_types/scheme_type_id.h b/ydb/public/lib/scheme_types/scheme_type_id.h index b8a6e2573d65f..e1a73dd3c5703 100644 --- a/ydb/public/lib/scheme_types/scheme_type_id.h +++ b/ydb/public/lib/scheme_types/scheme_type_id.h @@ -56,6 +56,7 @@ static constexpr TTypeId Json = NYql::NProto::Json; static constexpr TTypeId Uuid = NYql::NProto::Uuid; static constexpr TTypeId JsonDocument = NYql::NProto::JsonDocument; +static constexpr TTypeId Rowid = NYql::NProto::Rowid; static constexpr TTypeId DyNumber = NYql::NProto::DyNumber; @@ -87,6 +88,7 @@ static constexpr TTypeId YqlIds[] = { JsonDocument, DyNumber, Uuid, + Rowid, Date32, Datetime64, Timestamp64, @@ -148,6 +150,7 @@ const char *TypeName(TTypeId typeId) { case NTypeIds::JsonDocument: return "JsonDocument"; case NTypeIds::DyNumber: return "DyNumber"; case NTypeIds::Uuid: return "Uuid"; + case NTypeIds::Rowid: return "Rowid"; default: return "Unknown"; } } diff --git a/ydb/public/lib/value/value.cpp b/ydb/public/lib/value/value.cpp index f91674135dc0c..c7b3d2739f2ca 100644 --- a/ydb/public/lib/value/value.cpp +++ b/ydb/public/lib/value/value.cpp @@ -1,6 +1,7 @@ #include "value.h" #include +#include #include @@ -499,6 +500,8 @@ TString TValue::GetDataText() const { NYdb::TUuidValue val(Value.GetLow128(), Value.GetHi128()); return val.ToString(); } + case NScheme::NTypeIds::Rowid: + return NRowid::RowidBytesToBase64(Value.GetBytes()); } return TStringBuilder() << "\"\""; diff --git a/ydb/public/lib/value/ya.make b/ydb/public/lib/value/ya.make index dae7563b03713..2712d3cb8a3ec 100644 --- a/ydb/public/lib/value/ya.make +++ b/ydb/public/lib/value/ya.make @@ -12,6 +12,7 @@ PEERDIR( ydb/library/mkql_proto/protos ydb/public/lib/scheme_types ydb/public/sdk/cpp/src/client/value + yql/essentials/types/rowid ) END() diff --git a/ydb/public/lib/yson_value/ydb_yson_value.cpp b/ydb/public/lib/yson_value/ydb_yson_value.cpp index a75e59c38c3cb..b400156e2ccbd 100644 --- a/ydb/public/lib/yson_value/ydb_yson_value.cpp +++ b/ydb/public/lib/yson_value/ydb_yson_value.cpp @@ -94,6 +94,9 @@ static void PrimitiveValueToYson(EPrimitiveType type, TValueParser& parser, NYso case EPrimitiveType::Uuid: writer.OnStringScalar(parser.GetUuid().ToString()); break; + case EPrimitiveType::Rowid: + writer.OnStringScalar(parser.GetRowid().ToString()); + break; case EPrimitiveType::DyNumber: writer.OnStringScalar(parser.GetDyNumber()); break; diff --git a/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/result/detail/rows_parser_getter.h b/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/result/detail/rows_parser_getter.h index ad6cde3dd4b3f..b86b355fc8aef 100644 --- a/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/result/detail/rows_parser_getter.h +++ b/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/result/detail/rows_parser_getter.h @@ -105,6 +105,7 @@ Y_DEFINE_VALUE_PARSER_GETTER(uint64_t, Uint64); Y_DEFINE_VALUE_PARSER_GETTER(float, Float); Y_DEFINE_VALUE_PARSER_GETTER(double, Double); Y_DEFINE_VALUE_PARSER_GETTER(TUuidValue, Uuid); +Y_DEFINE_VALUE_PARSER_GETTER(TRowidValue, Rowid); Y_DEFINE_VALUE_PARSER_GETTER(TDate32, Date32); Y_DEFINE_VALUE_PARSER_GETTER(TDatetime64, Datetime64); Y_DEFINE_VALUE_PARSER_GETTER(TTimestamp64, Timestamp64); diff --git a/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/value/value.h b/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/value/value.h index 0cab88542e014..fe3f1a3fea608 100644 --- a/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/value/value.h +++ b/ydb/public/sdk/cpp/include/ydb-cpp-sdk/client/value/value.h @@ -8,6 +8,7 @@ #include #include +#include namespace Ydb { class Type; @@ -75,6 +76,7 @@ enum class EPrimitiveType { Json = 0x1202, Uuid = 0x1203, JsonDocument = 0x1204, + Rowid = 0x1205, DyNumber = 0x1302, }; @@ -274,6 +276,17 @@ struct TUuidValue { } Buf_; }; +struct TRowidValue { + static constexpr size_t Size = 14; + + std::string ToString() const; + TRowidValue(const char* rowidBytes, size_t size); + TRowidValue(const Ydb::Value& rowidValueProto); + TRowidValue(const std::string& rowidString); + + char Bytes[Size]; +}; + //! Representation of YDB value. class TValue { friend class TValueParser; @@ -348,6 +361,7 @@ class TValueParser : public TMoveOnly { TDecimalValue GetDecimal() const; TPgValue GetPg() const; TUuidValue GetUuid() const; + TRowidValue GetRowid() const; const std::string& GetJsonDocument() const; const std::string& GetDyNumber() const; @@ -381,6 +395,7 @@ class TValueParser : public TMoveOnly { std::optional GetOptionalJson() const; std::optional GetOptionalDecimal() const; std::optional GetOptionalUuid() const; + std::optional GetOptionalRowid() const; std::optional GetOptionalJsonDocument() const; std::optional GetOptionalDyNumber() const; @@ -463,6 +478,7 @@ class TValueBuilderBase : public TMoveOnly { TDerived& Decimal(const TDecimalValue& value); TDerived& Pg(const TPgValue& value); TDerived& Uuid(const TUuidValue& value); + TDerived& Rowid(const TRowidValue& value); TDerived& JsonDocument(const std::string& value); TDerived& DyNumber(const std::string& value); TDerived& Date32(const std::chrono::sys_time& value); @@ -495,6 +511,7 @@ class TValueBuilderBase : public TMoveOnly { TDerived& OptionalYson(const std::optional& value); TDerived& OptionalJson(const std::optional& value); TDerived& OptionalUuid(const std::optional& value); + TDerived& OptionalRowid(const std::optional& value); TDerived& OptionalJsonDocument(const std::optional& value); TDerived& OptionalDyNumber(const std::optional& value); TDerived& OptionalDate32(const std::optional>& value); @@ -578,3 +595,6 @@ class TValueBuilder : public TValueBuilderBase { template<> void Out(IOutputStream& o, const NYdb::TUuidValue& value); + +template<> +void Out(IOutputStream& o, const NYdb::TRowidValue& value); diff --git a/ydb/public/sdk/cpp/src/client/value/value.cpp b/ydb/public/sdk/cpp/src/client/value/value.cpp index 6b1055aea44a3..461ef473e6f61 100644 --- a/ydb/public/sdk/cpp/src/client/value/value.cpp +++ b/ydb/public/sdk/cpp/src/client/value/value.cpp @@ -15,6 +15,8 @@ #include #include +#include + #include #include #include @@ -1040,6 +1042,30 @@ std::string TUuidValue::ToString() const { return s.Str(); } +TRowidValue::TRowidValue(const char* rowidBytes, size_t size) { + static_assert(Size == NKikimr::NRowid::ROWID_LEN); + if (size != Size) { + ThrowFatalError(TStringBuilder() << "Invalid Rowid bytes size: " << size); + } + std::memcpy(Bytes, rowidBytes, Size); +} + +TRowidValue::TRowidValue(const Ydb::Value& valueProto) + : TRowidValue(valueProto.bytes_value().data(), valueProto.bytes_value().size()) +{} + +TRowidValue::TRowidValue(const std::string& rowidString) { + static_assert(Size == NKikimr::NRowid::ROWID_LEN); + if (!NKikimr::NRowid::ParseRowidBase64(rowidString, Bytes)) { + ThrowFatalError(TStringBuilder() << "Unable to parse string as rowid"); + } +} + +std::string TRowidValue::ToString() const { + static_assert(Size == NKikimr::NRowid::ROWID_LEN); + return NKikimr::NRowid::RowidBytesToBase64(TStringBuf(Bytes, Size)); +} + //////////////////////////////////////////////////////////////////////////////// class TValue::TImpl { @@ -1290,6 +1316,11 @@ class TValueParser::TImpl { return TUuidValue(GetProto()); } + TRowidValue GetRowid() const { + CheckPrimitive(NYdb::EPrimitiveType::Rowid); + return TRowidValue(GetProto()); + } + const std::string& GetJsonDocument() const { CheckPrimitive(NYdb::EPrimitiveType::JsonDocument); return GetProto().text_value(); @@ -1662,6 +1693,8 @@ class TValueParser::TImpl { return Ydb::Value::kTextValue; case NYdb::EPrimitiveType::Uuid: return Ydb::Value::kLow128; + case NYdb::EPrimitiveType::Rowid: + return Ydb::Value::kBytesValue; default: FatalError(TStringBuilder() << "Unexpected primitive type: " << primitiveTypeId); return Ydb::Value::kBytesValue; @@ -1819,6 +1852,10 @@ TUuidValue TValueParser::GetUuid() const { return Impl_->GetUuid(); } +TRowidValue TValueParser::GetRowid() const { + return Impl_->GetRowid(); +} + const std::string& TValueParser::GetJsonDocument() const { return Impl_->GetJsonDocument(); } @@ -1959,6 +1996,10 @@ std::optional TValueParser::GetOptionalUuid() const { RET_OPT_VALUE(TUuidValue, Uuid); } +std::optional TValueParser::GetOptionalRowid() const { + RET_OPT_VALUE(TRowidValue, Rowid); +} + std::optional TValueParser::GetOptionalJsonDocument() const { RET_OPT_VALUE(std::string, JsonDocument); } @@ -2294,6 +2335,11 @@ class TValueBuilderImpl { GetValue().set_high_128(value.Buf_.Halfs[1]); } + void Rowid(const TRowidValue& value) { + FillPrimitiveType(EPrimitiveType::Rowid); + GetValue().set_bytes_value(value.Bytes, TRowidValue::Size); + } + void JsonDocument(const std::string& value) { FillPrimitiveType(EPrimitiveType::JsonDocument); GetValue().set_text_value(TStringType{value}); @@ -3082,6 +3128,12 @@ TDerived& TValueBuilderBase::Uuid(const TUuidValue& value) { return static_cast(*this); } +template +TDerived& TValueBuilderBase::Rowid(const TRowidValue& value) { + Impl_->Rowid(value); + return static_cast(*this); +} + template TDerived& TValueBuilderBase::JsonDocument(const std::string& value) { Impl_->JsonDocument(value); @@ -3267,6 +3319,11 @@ TDerived& TValueBuilderBase::OptionalUuid(const std::optional +TDerived& TValueBuilderBase::OptionalRowid(const std::optional& value) { + SET_OPT_VALUE_FROM_OPTIONAL(Rowid); +} + template TDerived& TValueBuilderBase::OptionalJsonDocument(const std::optional& value) { SET_OPT_VALUE_FROM_OPTIONAL(JsonDocument); @@ -3496,3 +3553,8 @@ template<> void Out(IOutputStream& o, const NYdb::TUuidValue& value) { o << value.ToString(); } + +template<> +void Out(IOutputStream& o, const NYdb::TRowidValue& value) { + o << value.ToString(); +} diff --git a/ydb/public/sdk/cpp/src/client/value/ya.make b/ydb/public/sdk/cpp/src/client/value/ya.make index 063865e23a6ca..a72e8486a664e 100644 --- a/ydb/public/sdk/cpp/src/client/value/ya.make +++ b/ydb/public/sdk/cpp/src/client/value/ya.make @@ -14,6 +14,7 @@ PEERDIR( ydb/public/sdk/cpp/src/client/types/fatal_error_handlers ydb/public/sdk/cpp/src/library/decimal ydb/public/sdk/cpp/src/library/uuid + yql/essentials/types/rowid ) END() diff --git a/ydb/public/sdk/cpp/tests/unit/client/value/value_ut.cpp b/ydb/public/sdk/cpp/tests/unit/client/value/value_ut.cpp index edc9a1de9d3df..af238a8279e57 100644 --- a/ydb/public/sdk/cpp/tests/unit/client/value/value_ut.cpp +++ b/ydb/public/sdk/cpp/tests/unit/client/value/value_ut.cpp @@ -1466,4 +1466,24 @@ TEST(YdbValue, IncorrectUuid) { ASSERT_THROW(TUuidValue("5ca32-c22841b-11e8-adc0-fa7ae01bbebc"), TContractViolation); } +TEST(YdbValue, CorrectRowid) { + std::string rowidStr = "YWJjZGVmZ2hpamtsbW4"; + TRowidValue rowid(rowidStr); + ASSERT_EQ(rowidStr, rowid.ToString()); + + auto value = TValueBuilder().Rowid(rowid).Build(); + CheckProtoValue(value.GetProto(), "bytes_value: \"abcdefghijklmn\"\n"); + + TValueParser parser(value); + ASSERT_EQ(EPrimitiveType::Rowid, parser.GetPrimitiveType()); + ASSERT_EQ(rowidStr, parser.GetRowid().ToString()); +} + +TEST(YdbValue, IncorrectRowid) { + ASSERT_THROW(TRowidValue(""), TContractViolation); + ASSERT_THROW(TRowidValue("YWJjZGVmZ2hpamtsbW4="), TContractViolation); + ASSERT_THROW(TRowidValue("YWJjZGVmZ2hpamtsbW*"), TContractViolation); + ASSERT_THROW(TRowidValue("abcdefghijklmn", 13), TContractViolation); +} + } // namespace NYdb diff --git a/yql/essentials/ast/yql_type_string.cpp b/yql/essentials/ast/yql_type_string.cpp index 3b8af1628c567..eb0e52739ca88 100644 --- a/yql/essentials/ast/yql_type_string.cpp +++ b/yql/essentials/ast/yql_type_string.cpp @@ -93,6 +93,7 @@ enum EToken { TOKEN_ERROR = -60, TOKEN_LINEAR = -61, TOKEN_DYNAMICLINEAR = -62, + TOKEN_ROWID = -63, // identifiers TOKEN_IDENTIFIER = -100, @@ -149,6 +150,7 @@ EToken TokenTypeFromStr(TStringBuf str) {TStringBuf("TzDatetime"), TOKEN_TZDATETIME}, {TStringBuf("TzTimestamp"), TOKEN_TZTIMESTAMP}, {TStringBuf("Uuid"), TOKEN_UUID}, + {TStringBuf("Rowid"), TOKEN_ROWID}, {TStringBuf("Flow"), TOKEN_FLOW}, {TStringBuf("Set"), TOKEN_SET}, {TStringBuf("Enum"), TOKEN_ENUM}, @@ -235,6 +237,7 @@ class TTypeParser { case TOKEN_TZDATETIME: case TOKEN_TZTIMESTAMP: case TOKEN_UUID: + case TOKEN_ROWID: case TOKEN_JSON_DOCUMENT: case TOKEN_DYNUMBER: case TOKEN_DATE32: diff --git a/yql/essentials/core/sql_types/simple_types.cpp b/yql/essentials/core/sql_types/simple_types.cpp index c6cd40e37aaaa..48eb25568792e 100644 --- a/yql/essentials/core/sql_types/simple_types.cpp +++ b/yql/essentials/core/sql_types/simple_types.cpp @@ -53,6 +53,7 @@ const std::unordered_map SimpleTypes = { {"utf8", {"Utf8", "Utf8", "Data"}}, {"uuid", {"Uuid", "Uuid", "Data"}}, + {"rowid", {"Rowid", "Rowid", "Data"}}, {"yson", {"Yson", "Yson", "Data"}}, {"json", {"Json", "Json", "Data"}}, {"jsondocument", {"JsonDocument", "JsonDocument", "Data"}}, diff --git a/yql/essentials/core/type_ann/type_ann_core.cpp b/yql/essentials/core/type_ann/type_ann_core.cpp index f4d65411f267b..4e93eeb86a6aa 100644 --- a/yql/essentials/core/type_ann/type_ann_core.cpp +++ b/yql/essentials/core/type_ann/type_ann_core.cpp @@ -983,6 +983,13 @@ namespace NTypeAnnImpl { ctx.Expr.AddError(TIssue(ctx.Expr.GetPosition(input->Pos()), TStringBuilder() << "Bad atom format for type: " << input->Content() << ", value: " << TString(input->Head().Content()).Quote())); + return IGraphTransformer::TStatus::Error; + } + } else if (input->Content() == "Rowid") { + if (input->Head().Content().size() != NKikimr::NUdf::ROWID_SIZE) { + ctx.Expr.AddError(TIssue(ctx.Expr.GetPosition(input->Pos()), TStringBuilder() << "Bad atom format for type: " + << input->Content() << ", value: " << TString(input->Head().Content()).Quote())); + return IGraphTransformer::TStatus::Error; } } else if (input->Content() == "JsonDocument") { diff --git a/yql/essentials/core/yql_default_valid_value.cpp b/yql/essentials/core/yql_default_valid_value.cpp index e4b05425e61ad..5052095f84203 100644 --- a/yql/essentials/core/yql_default_valid_value.cpp +++ b/yql/essentials/core/yql_default_valid_value.cpp @@ -315,6 +315,11 @@ class TValidValueNodeVisitor: public TTypeAnnotationVisitor { Result_ = Ctx_.Builder(Pos_).Callable("Uuid").Atom(0, TStringBuf(UUID_VALID_LITERAL, sizeof(UUID_VALID_LITERAL) - 1)).Seal().Build(); break; } + case NUdf::EDataSlot::Rowid: { + constexpr char ROWID_VALID_LITERAL[] = "\x00\x00\x00\x00\x00\x00\x00\x00\x00\x00\x00\x00\x00\x00"; + Result_ = Ctx_.Builder(Pos_).Callable("Rowid").Atom(0, TStringBuf(ROWID_VALID_LITERAL, sizeof(ROWID_VALID_LITERAL) - 1)).Seal().Build(); + break; + } case NUdf::EDataSlot::Date: Result_ = Ctx_.Builder(Pos_).Callable("Date").Atom(0, "0").Seal().Build(); break; diff --git a/yql/essentials/core/yql_opt_utils.cpp b/yql/essentials/core/yql_opt_utils.cpp index ccad22335fce2..6542ec99d30fa 100644 --- a/yql/essentials/core/yql_opt_utils.cpp +++ b/yql/essentials/core/yql_opt_utils.cpp @@ -1923,6 +1923,7 @@ ui64 GetTypeWeight(const TTypeAnnotationNode& type) { case NUdf::EDataSlot::TzTimestamp: return 9; case NUdf::EDataSlot::Decimal: return 15; + case NUdf::EDataSlot::Rowid: return 14; case NUdf::EDataSlot::Uuid: return 16; default: return 32; diff --git a/yql/essentials/minikql/computation/mkql_block_impl.cpp b/yql/essentials/minikql/computation/mkql_block_impl.cpp index b9afd7bb6c958..343c06637140d 100644 --- a/yql/essentials/minikql/computation/mkql_block_impl.cpp +++ b/yql/essentials/minikql/computation/mkql_block_impl.cpp @@ -132,12 +132,13 @@ arrow::Datum DoConvertScalar(TType* type, const T& value, arrow::MemoryPool& poo case NUdf::EDataSlot::Yson: case NUdf::EDataSlot::Json: case NUdf::EDataSlot::JsonDocument: - case NUdf::EDataSlot::DyNumber: { + case NUdf::EDataSlot::DyNumber: + case NUdf::EDataSlot::Rowid: { const auto& str = value.AsStringRef(); std::shared_ptr buffer(ARROW_RESULT(arrow::AllocateBuffer(str.Size(), &pool))); std::memcpy(buffer->mutable_data(), str.Data(), str.Size()); std::shared_ptr scalar; - if (slot == NUdf::EDataSlot::String || slot == NUdf::EDataSlot::Yson || slot == NUdf::EDataSlot::JsonDocument || slot == NUdf::EDataSlot::DyNumber) { + if (slot == NUdf::EDataSlot::String || slot == NUdf::EDataSlot::Yson || slot == NUdf::EDataSlot::JsonDocument || slot == NUdf::EDataSlot::DyNumber || slot == NUdf::EDataSlot::Rowid) { scalar = std::make_shared(buffer, arrow::binary()); } else { // NOTE: Do not use |arrow::BinaryScalar| for utf8 and json types directly. diff --git a/yql/essentials/minikql/computation/mkql_computation_node_pack.cpp b/yql/essentials/minikql/computation/mkql_computation_node_pack.cpp index b7c45fa496596..669b3fc02e185 100644 --- a/yql/essentials/minikql/computation/mkql_computation_node_pack.cpp +++ b/yql/essentials/minikql/computation/mkql_computation_node_pack.cpp @@ -395,6 +395,9 @@ NUdf::TUnboxedValue UnpackFromChunkedBuffer(const TType* type, TChunkedInputBuff case NUdf::EDataSlot::Uuid: { return UnpackString(buf, 16); } + case NUdf::EDataSlot::Rowid: { + return UnpackString(buf, NUdf::ROWID_SIZE); + } case NUdf::EDataSlot::Decimal: { return NUdf::TUnboxedValuePod(UnpackDecimal(buf)); } @@ -721,6 +724,11 @@ void PackImpl(const TType* type, TBuf& buffer, const NUdf::TUnboxedValuePod& val PackBlob(ref.Data(), ref.Size(), buffer); break; } + case NUdf::EDataSlot::Rowid: { + auto ref = value.AsStringRef(); + PackBlob(ref.Data(), ref.Size(), buffer); + break; + } case NUdf::EDataSlot::TzDate: { PackData(value.Get(), buffer); PackData(value.GetTimezoneId(), buffer); diff --git a/yql/essentials/minikql/computation/presort.cpp b/yql/essentials/minikql/computation/presort.cpp index f0ef7ccf56c29..d5d5a75187d64 100644 --- a/yql/essentials/minikql/computation/presort.cpp +++ b/yql/essentials/minikql/computation/presort.cpp @@ -18,41 +18,62 @@ namespace NKikimr::NMiniKQL { namespace NDetail { constexpr size_t UuidSize = 16; +constexpr size_t RowidSize = NYql::NUdf::ROWID_SIZE; -template -Y_FORCE_INLINE void EncodeUuid(TVector& output, const char* data) { - output.resize(output.size() + UuidSize); - auto ptr = output.end() - UuidSize; +template +Y_FORCE_INLINE void EncodeFixedBytes(TVector& output, const char* data) { + output.resize(output.size() + Size); + auto ptr = output.end() - Size; if (Desc) { - for (size_t i = 0; i < UuidSize; ++i) { + for (size_t i = 0; i < Size; ++i) { *ptr++ = ui8(*data++) ^ 0xFF; } } else { - std::memcpy(ptr, data, UuidSize); + std::memcpy(ptr, data, Size); } } template -Y_FORCE_INLINE TStringBuf DecodeUuid(TStringBuf& input, TVector& value) { - EnsureInputSize(input, UuidSize); +Y_FORCE_INLINE void EncodeUuid(TVector& output, const char* data) { + EncodeFixedBytes(output, data); +} + +template +Y_FORCE_INLINE void EncodeRowid(TVector& output, const char* data) { + EncodeFixedBytes(output, data); +} + +template +Y_FORCE_INLINE TStringBuf DecodeFixedBytes(TStringBuf& input, TVector& value) { + EnsureInputSize(input, Size); auto data = input.data(); - input.Skip(UuidSize); + input.Skip(Size); - value.resize(UuidSize); + value.resize(Size); auto ptr = value.begin(); if (Desc) { - for (size_t i = 0; i < UuidSize; ++i) { + for (size_t i = 0; i < Size; ++i) { *ptr++ = ui8(*data++) ^ 0xFF; } } else { - std::memcpy(ptr, data, UuidSize); + std::memcpy(ptr, data, Size); } return TStringBuf((const char*)value.begin(), (const char*)value.end()); } +template +Y_FORCE_INLINE TStringBuf DecodeUuid(TStringBuf& input, TVector& value) { + return DecodeFixedBytes(input, value); +} + +template +Y_FORCE_INLINE TStringBuf DecodeRowid(TStringBuf& input, TVector& value) { + return DecodeFixedBytes(input, value); +} + template Y_FORCE_INLINE void EncodeTzUnsigned(TVector& output, TUnsigned value, ui16 tzId) { constexpr size_t size = sizeof(TUnsigned); @@ -172,6 +193,9 @@ Y_FORCE_INLINE void Encode(TVector& output, NUdf::EDataSlot slot, const NUd case NUdf::EDataSlot::Uuid: EncodeUuid(output, value.AsStringRef().Data()); break; + case NUdf::EDataSlot::Rowid: + EncodeRowid(output, value.AsStringRef().Data()); + break; case NUdf::EDataSlot::TzDate: EncodeTzUnsigned(output, value.Get(), value.GetTimezoneId()); break; @@ -254,6 +278,10 @@ Y_FORCE_INLINE NUdf::TUnboxedValue Decode(TStringBuf& input, NUdf::EDataSlot slo buffer.clear(); return MakeString(NUdf::TStringRef(DecodeUuid(input, buffer))); + case NUdf::EDataSlot::Rowid: + buffer.clear(); + return MakeString(NUdf::TStringRef(DecodeRowid(input, buffer))); + case NUdf::EDataSlot::TzDate: { ui16 date; ui16 tzId; diff --git a/yql/essentials/minikql/invoke_builtins/mkql_builtins_compare.h b/yql/essentials/minikql/invoke_builtins/mkql_builtins_compare.h index 061e66632d7f1..72156b69e870d 100644 --- a/yql/essentials/minikql/invoke_builtins/mkql_builtins_compare.h +++ b/yql/essentials/minikql/invoke_builtins/mkql_builtins_compare.h @@ -642,6 +642,7 @@ void RegisterCompareStrings(IBuiltinFunctionRegistry& registry, const std::strin RegisterCompareCustomOpt, NUdf::TDataType, TFunc, TArgs>(registry, name); if constexpr (WithSpecial) { RegisterCompareCustomOpt, NUdf::TDataType, TFunc, TArgs>(registry, name); + RegisterCompareCustomOpt, NUdf::TDataType, TFunc, TArgs>(registry, name); RegisterCompareCustomOpt, NUdf::TDataType, TFunc, TArgs>(registry, name); } } @@ -653,6 +654,7 @@ void RegisterAggrCompareStrings(IBuiltinFunctionRegistry& registry, const std::s RegisterAggrCompareCustomOpt, TFunc, TArgs>(registry, name); RegisterAggrCompareCustomOpt, TFunc, TArgs>(registry, name); RegisterAggrCompareCustomOpt, TFunc, TArgs>(registry, name); + RegisterAggrCompareCustomOpt, TFunc, TArgs>(registry, name); RegisterAggrCompareCustomOpt, TFunc, TArgs>(registry, name); } diff --git a/yql/essentials/minikql/invoke_builtins/mkql_builtins_max.cpp b/yql/essentials/minikql/invoke_builtins/mkql_builtins_max.cpp index 56d6623d3d8bd..0ed8a897156ea 100644 --- a/yql/essentials/minikql/invoke_builtins/mkql_builtins_max.cpp +++ b/yql/essentials/minikql/invoke_builtins/mkql_builtins_max.cpp @@ -210,6 +210,7 @@ void RegisterMax(IBuiltinFunctionRegistry& registry) { RegisterCustomSameTypesFunction, TCustomMax, TBinaryArgsOpt>(registry, "Max"); RegisterCustomSameTypesFunction, TCustomMax, TBinaryArgsOpt>(registry, "Max"); RegisterCustomSameTypesFunction, TCustomMax, TBinaryArgsOpt>(registry, "Max"); + RegisterCustomSameTypesFunction, TCustomMax, TBinaryArgsOpt>(registry, "Max"); } void RegisterAggrMax(IBuiltinFunctionRegistry& registry) { @@ -225,6 +226,7 @@ void RegisterAggrMax(IBuiltinFunctionRegistry& registry) { RegisterCustomAggregateFunction, TCustomMax, TBinaryArgsSameOpt>(registry, "AggrMax"); RegisterCustomAggregateFunction, TCustomMax, TBinaryArgsSameOpt>(registry, "AggrMax"); RegisterCustomAggregateFunction, TCustomMax, TBinaryArgsSameOpt>(registry, "AggrMax"); + RegisterCustomAggregateFunction, TCustomMax, TBinaryArgsSameOpt>(registry, "AggrMax"); } } // namespace NMiniKQL diff --git a/yql/essentials/minikql/invoke_builtins/mkql_builtins_min.cpp b/yql/essentials/minikql/invoke_builtins/mkql_builtins_min.cpp index f2153fe069f88..b88862d2450d2 100644 --- a/yql/essentials/minikql/invoke_builtins/mkql_builtins_min.cpp +++ b/yql/essentials/minikql/invoke_builtins/mkql_builtins_min.cpp @@ -209,6 +209,7 @@ void RegisterMin(IBuiltinFunctionRegistry& registry) { RegisterCustomSameTypesFunction, TCustomMin, TBinaryArgsOpt>(registry, "Min"); RegisterCustomSameTypesFunction, TCustomMin, TBinaryArgsOpt>(registry, "Min"); RegisterCustomSameTypesFunction, TCustomMin, TBinaryArgsOpt>(registry, "Min"); + RegisterCustomSameTypesFunction, TCustomMin, TBinaryArgsOpt>(registry, "Min"); } void RegisterAggrMin(IBuiltinFunctionRegistry& registry) { @@ -224,6 +225,7 @@ void RegisterAggrMin(IBuiltinFunctionRegistry& registry) { RegisterCustomAggregateFunction, TCustomMin, TBinaryArgsSameOpt>(registry, "AggrMin"); RegisterCustomAggregateFunction, TCustomMin, TBinaryArgsSameOpt>(registry, "AggrMin"); RegisterCustomAggregateFunction, TCustomMin, TBinaryArgsSameOpt>(registry, "AggrMin"); + RegisterCustomAggregateFunction, TCustomMin, TBinaryArgsSameOpt>(registry, "AggrMin"); } } // namespace NMiniKQL diff --git a/yql/essentials/minikql/mkql_node_serialization.cpp b/yql/essentials/minikql/mkql_node_serialization.cpp index b705f9a0f60ec..7010f72828124 100644 --- a/yql/essentials/minikql/mkql_node_serialization.cpp +++ b/yql/essentials/minikql/mkql_node_serialization.cpp @@ -809,6 +809,11 @@ class TWriter { Owner_.WriteMany(v.Data(), v.Size()); break; } + case NUdf::TDataType::Id: { + const auto v = value.AsStringRef(); + Owner_.WriteMany(v.Data(), v.Size()); + break; + } case NUdf::TDataType::Id: Owner_.WriteMany(static_cast(value.GetRawPtr()), sizeof(NYql::NDecimal::TInt128) - 1U); break; @@ -1838,6 +1843,11 @@ class TReader { value = Env_.NewStringValue(NUdf::TStringRef(buffer, 16)); break; } + case NUdf::TDataType::Id: { + const char* buffer = ReadMany(NUdf::ROWID_SIZE); + value = Env_.NewStringValue(NUdf::TStringRef(buffer, NUdf::ROWID_SIZE)); + break; + } case NUdf::TDataType::Id: { value = NUdf::TUnboxedValuePod::Zero(); const char* buffer = ReadMany(sizeof(NYql::NDecimal::TInt128) - 1U); diff --git a/yql/essentials/minikql/mkql_type_builder.cpp b/yql/essentials/minikql/mkql_type_builder.cpp index 4c354629b010a..991fa934039bf 100644 --- a/yql/essentials/minikql/mkql_type_builder.cpp +++ b/yql/essentials/minikql/mkql_type_builder.cpp @@ -1560,6 +1560,10 @@ bool ConvertArrowTypeImpl(NUdf::EDataSlot slot, std::shared_ptr type = arrow::fixed_size_binary(UuidBinarySize); return true; } + case NUdf::EDataSlot::Rowid: { + type = arrow::binary(); + return true; + } case NUdf::EDataSlot::Decimal: { type = arrow::fixed_size_binary(sizeof(NYql::NUdf::TUnboxedValuePod)); return true; @@ -2708,6 +2712,7 @@ size_t CalcMaxBlockItemSize(const TType* type) { case NUdf::EDataSlot::Decimal: { return sizeof(NYql::NDecimal::TInt128); } + case NUdf::EDataSlot::Rowid: case NUdf::EDataSlot::DyNumber: return sizeof(arrow::BinaryType::offset_type); } diff --git a/yql/essentials/minikql/mkql_type_ops.cpp b/yql/essentials/minikql/mkql_type_ops.cpp index 107de0808ed1e..0fefe662082d1 100644 --- a/yql/essentials/minikql/mkql_type_ops.cpp +++ b/yql/essentials/minikql/mkql_type_ops.cpp @@ -11,6 +11,7 @@ #include #include #include +#include #include #include @@ -175,6 +176,8 @@ bool IsValidValue(NUdf::EDataSlot type, const NUdf::TUnboxedValuePod& value) { return bool(value) && NDom::IsValidJson(value.AsStringRef()); case NUdf::EDataSlot::Uuid: return bool(value) && value.AsStringRef().Size() == 16; + case NUdf::EDataSlot::Rowid: + return bool(value) && value.AsStringRef().Size() == NUdf::ROWID_SIZE; case NUdf::EDataSlot::DyNumber: return NDyNumber::IsValidDyNumber(value.AsStringRef()); case NUdf::EDataSlot::JsonDocument: @@ -537,6 +540,11 @@ NUdf::TUnboxedValuePod ValueToString(NUdf::EDataSlot type, NUdf::TUnboxedValuePo break; } + case NUdf::EDataSlot::Rowid: { + NRowid::RowidBytesToBase64(value.AsStringRef(), out); + break; + } + case NUdf::EDataSlot::Date: if (!WriteDate(out, value.Get())) { return NUdf::TUnboxedValuePod(); @@ -1535,6 +1543,14 @@ NUdf::TUnboxedValuePod ParseUuid(NUdf::TStringRef buf, bool shortForm) { return MakeString(NUdf::TStringRef(reinterpret_cast(dw.data()), sizeof(dw))); } +NUdf::TUnboxedValuePod ParseRowid(NUdf::TStringRef buf) { + std::array bytes; + if (!NRowid::ParseRowidBase64(buf, bytes.data())) { + return NUdf::TUnboxedValuePod(); + } + return MakeString(NUdf::TStringRef(bytes.data(), bytes.size())); +} + bool ParseUuid(NUdf::TStringRef buf, void* out, bool shortForm) { std::array dw; @@ -2491,6 +2507,9 @@ bool IsValidStringValue(NUdf::EDataSlot type, NUdf::TStringRef buf) { case NUdf::EDataSlot::Uuid: return NUuid::IsValidUuid(buf); + case NUdf::EDataSlot::Rowid: + return NRowid::IsValidRowidBase64(buf); + case NUdf::EDataSlot::DyNumber: return NDyNumber::IsValidDyNumberString(buf); @@ -2609,6 +2628,9 @@ NUdf::TUnboxedValuePod ValueFromString(NUdf::EDataSlot type, NUdf::TStringRef bu case NUdf::EDataSlot::Uuid: return ParseUuid(buf); + case NUdf::EDataSlot::Rowid: + return ParseRowid(buf); + case NUdf::EDataSlot::Date: return ParseDate(buf); @@ -2741,6 +2763,7 @@ NUdf::TUnboxedValuePod SimpleValueFromYson(NUdf::EDataSlot type, NUdf::TStringRe case NUdf::EDataSlot::TzTimestamp: case NUdf::EDataSlot::Decimal: case NUdf::EDataSlot::Uuid: + case NUdf::EDataSlot::Rowid: Y_ABORT("TODO"); default:; @@ -2856,6 +2879,7 @@ NUdf::TUnboxedValuePod SimpleValueFromYson(NUdf::EDataSlot type, NUdf::TStringRe case NUdf::EDataSlot::TzTimestamp64: case NUdf::EDataSlot::Decimal: case NUdf::EDataSlot::Uuid: + case NUdf::EDataSlot::Rowid: case NUdf::EDataSlot::DyNumber: case NUdf::EDataSlot::JsonDocument: Y_ABORT("TODO"); diff --git a/yql/essentials/minikql/ya.make b/yql/essentials/minikql/ya.make index 0cfd197043cd2..5298fae174d80 100644 --- a/yql/essentials/minikql/ya.make +++ b/yql/essentials/minikql/ya.make @@ -81,6 +81,7 @@ PEERDIR( yql/essentials/types/binary_json yql/essentials/types/dynumber yql/essentials/types/uuid + yql/essentials/types/rowid yql/essentials/utils yql/essentials/utils/memory_profiling ) diff --git a/yql/essentials/parser/pg_wrapper/comp_factory.cpp b/yql/essentials/parser/pg_wrapper/comp_factory.cpp index 27601a078b3ac..10b19efc1cd61 100644 --- a/yql/essentials/parser/pg_wrapper/comp_factory.cpp +++ b/yql/essentials/parser/pg_wrapper/comp_factory.cpp @@ -3877,6 +3877,8 @@ ui32 ConvertToPgType(NUdf::EDataSlot slot) { return JSONOID; case NUdf::EDataSlot::Uuid: return UUIDOID; + case NUdf::EDataSlot::Rowid: + return BYTEAOID; case NUdf::EDataSlot::Date: return DATEOID; case NUdf::EDataSlot::Datetime: diff --git a/yql/essentials/providers/common/codec/yql_codec.cpp b/yql/essentials/providers/common/codec/yql_codec.cpp index e3e735f061ee1..9500426f9d1d7 100644 --- a/yql/essentials/providers/common/codec/yql_codec.cpp +++ b/yql/essentials/providers/common/codec/yql_codec.cpp @@ -90,6 +90,7 @@ void WriteYsonValueImpl(NResult::TYsonResultWriter& writer, const NUdf::TUnboxed case NUdf::EDataSlot::String: case NUdf::EDataSlot::Uuid: + case NUdf::EDataSlot::Rowid: writer.OnStringScalar(value.AsStringRef()); return; @@ -332,6 +333,7 @@ NYT::TNode DataValueToNode(const NKikimr::NUdf::TUnboxedValuePod& value, NKikimr case NUdf::TDataType::Id: case NUdf::TDataType::Id: case NUdf::TDataType::Id: + case NUdf::TDataType::Id: return NYT::TNode(TString(value.AsStringRef())); case NUdf::TDataType::Id: return NYT::NodeFromYsonString(value.AsStringRef()); @@ -472,6 +474,7 @@ TString DataValueToString(const NKikimr::NUdf::TUnboxedValuePod& value, const TD case NUdf::EDataSlot::Utf8: case NUdf::EDataSlot::Json: case NUdf::EDataSlot::Uuid: + case NUdf::EDataSlot::Rowid: case NUdf::EDataSlot::Yson: return ToString((TStringBuf)value.AsStringRef()); case NUdf::EDataSlot::Decimal: { @@ -871,7 +874,8 @@ NUdf::TUnboxedValue ReadYsonValue(TType* type, case NUdf::TDataType::Id: case NUdf::TDataType::Id: case NUdf::TDataType::Id: - case NUdf::TDataType::Id: { + case NUdf::TDataType::Id: + case NUdf::TDataType::Id: { return ReadYsonStringInResultFormat(cmd, buf); } diff --git a/yql/essentials/providers/common/codec/yql_json_codec.cpp b/yql/essentials/providers/common/codec/yql_json_codec.cpp index 97eb4a0f1643f..ca36c816712b4 100644 --- a/yql/essentials/providers/common/codec/yql_json_codec.cpp +++ b/yql/essentials/providers/common/codec/yql_json_codec.cpp @@ -198,6 +198,7 @@ void WriteValueToJson(TJsonWriter& writer, const NKikimr::NUdf::TUnboxedValuePod case NUdf::TDataType::Id: writer.Write(value.AsStringRef()); break; + case NUdf::TDataType::Id: case NUdf::TDataType::Id: case NUdf::TDataType::Id: case NUdf::TDataType::Id: diff --git a/yql/essentials/providers/common/schema/skiff/yql_skiff_schema.cpp b/yql/essentials/providers/common/schema/skiff/yql_skiff_schema.cpp index ec0414954ce16..90088f9cdb341 100644 --- a/yql/essentials/providers/common/schema/skiff/yql_skiff_schema.cpp +++ b/yql/essentials/providers/common/schema/skiff/yql_skiff_schema.cpp @@ -94,6 +94,7 @@ struct TSkiffTypeLoader { case NUdf::EDataSlot::Utf8: case NUdf::EDataSlot::Json: case NUdf::EDataSlot::Uuid: + case NUdf::EDataSlot::Rowid: case NUdf::EDataSlot::DyNumber: case NUdf::EDataSlot::JsonDocument: return NYT::TNode()("wire_type", "string32"); diff --git a/yql/essentials/public/types/yql_types.proto b/yql/essentials/public/types/yql_types.proto index 32dfcbd1b4ac5..d96748402b71d 100644 --- a/yql/essentials/public/types/yql_types.proto +++ b/yql/essentials/public/types/yql_types.proto @@ -21,6 +21,7 @@ enum TypeIds { Json = 0x1202; Uuid = 0x1203; JsonDocument = 0x1204; + Rowid = 0x1205; Date = 0x0030; Datetime = 0x0031; Timestamp = 0x0032; diff --git a/yql/essentials/public/udf/arrow/dispatch_traits.h b/yql/essentials/public/udf/arrow/dispatch_traits.h index 50e9a526e76b4..0568194c9d304 100644 --- a/yql/essentials/public/udf/arrow/dispatch_traits.h +++ b/yql/essentials/public/udf/arrow/dispatch_traits.h @@ -247,6 +247,9 @@ std::unique_ptr DispatchByArrowTraits(const ITypeInfo Y_ENSURE(false, "Unsupported data slot"); } } + case NUdf::EDataSlot::Rowid: + // Stored as 14 raw bytes in MiniKQL; Arrow uses variable binary for now. + return MakeStringArrowTraitsImpl(isOptional, type, std::forward(args)...); case NUdf::EDataSlot::DyNumber: return MakeStringArrowTraitsImpl(isOptional, type, std::forward(args)...); } diff --git a/yql/essentials/public/udf/udf_data_type.cpp b/yql/essentials/public/udf/udf_data_type.cpp index db96a88e7ac0c..8bdcd9a73590f 100644 --- a/yql/essentials/public/udf/udf_data_type.cpp +++ b/yql/essentials/public/udf/udf_data_type.cpp @@ -33,52 +33,53 @@ ui8 GetDecimalWidth(EDataTypeFeatures features) { #define UN {ECastOptions::Impossible | ECastOptions::Undefined} const std::array, DataSlotCount>, DataSlotCount> CastResultsTable = {{ - // Bool, Int8 ----integrals---- Uint64 Floats, Strings, YJsons, Uuid, DateTimes, Interval, TzDateTimes, Decimal, DyNumber, JsonDocument - {{OK, OK, OK, OK, OK, OK, OK, OK, OK, OK, OK, OK, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO}}, // Bool + // Bool, Int8 ----integrals---- Uint64 Floats, Strings, YJsons, Uuid, DateTimes, Interval, TzDateTimes, Decimal, DyNumber, JsonDocument, ... Tz*, Rowid + {{OK, OK, OK, OK, OK, OK, OK, OK, OK, OK, OK, OK, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO}}, // Bool - {{LD, OK, MF, OK, MF, OK, MF, OK, MF, OK, OK, OK, OK, NO, NO, NO, MF, MF, MF, OK, MF, MF, MF, UN, NO, NO, OK, OK, OK, OK, NO, NO, NO}}, // Int8 - {{LD, MF, OK, OK, OK, OK, OK, OK, OK, OK, OK, OK, OK, NO, NO, NO, OK, OK, OK, OK, OK, OK, OK, UN, NO, NO, OK, OK, OK, OK, NO, NO, NO}}, // Uint8 - {{LD, MF, MF, OK, MF, OK, MF, OK, MF, OK, OK, OK, OK, NO, NO, NO, MF, MF, MF, OK, MF, MF, MF, UN, NO, NO, OK, OK, OK, OK, NO, NO, NO}}, // Int16 - {{LD, MF, MF, MF, OK, OK, OK, OK, OK, OK, OK, OK, OK, NO, NO, NO, MF, OK, OK, OK, MF, OK, OK, UN, NO, NO, OK, OK, OK, OK, NO, NO, NO}}, // Uint16 - {{LD, MF, MF, MF, MF, OK, MF, OK, MF, OK, LD, OK, OK, NO, NO, NO, MF, MF, MF, OK, MF, MF, MF, UN, NO, NO, MF, OK, OK, OK, NO, NO, NO}}, // Int32 - {{LD, MF, MF, MF, MF, MF, OK, OK, OK, OK, LD, OK, OK, NO, NO, NO, MF, MF, OK, OK, MF, MF, OK, UN, NO, NO, MF, OK, OK, OK, NO, NO, NO}}, // Uint32 - {{LD, MF, MF, MF, MF, MF, MF, OK, MF, LD, LD, OK, OK, NO, NO, NO, MF, MF, MF, MF, MF, MF, MF, UN, NO, NO, MF, MF, MF, MF, NO, NO, NO}}, // Int64 - {{LD, MF, MF, MF, MF, MF, MF, MF, OK, LD, LD, OK, OK, NO, NO, NO, MF, MF, MF, MF, MF, MF, MF, UN, NO, NO, MF, MF, MF, MF, NO, NO, NO}}, // Uint64 + {{LD, OK, MF, OK, MF, OK, MF, OK, MF, OK, OK, OK, OK, NO, NO, NO, MF, MF, MF, OK, MF, MF, MF, UN, NO, NO, OK, OK, OK, OK, NO, NO, NO, NO}}, // Int8 + {{LD, MF, OK, OK, OK, OK, OK, OK, OK, OK, OK, OK, OK, NO, NO, NO, OK, OK, OK, OK, OK, OK, OK, UN, NO, NO, OK, OK, OK, OK, NO, NO, NO, NO}}, // Uint8 + {{LD, MF, MF, OK, MF, OK, MF, OK, MF, OK, OK, OK, OK, NO, NO, NO, MF, MF, MF, OK, MF, MF, MF, UN, NO, NO, OK, OK, OK, OK, NO, NO, NO, NO}}, // Int16 + {{LD, MF, MF, MF, OK, OK, OK, OK, OK, OK, OK, OK, OK, NO, NO, NO, MF, OK, OK, OK, MF, OK, OK, UN, NO, NO, OK, OK, OK, OK, NO, NO, NO, NO}}, // Uint16 + {{LD, MF, MF, MF, MF, OK, MF, OK, MF, OK, LD, OK, OK, NO, NO, NO, MF, MF, MF, OK, MF, MF, MF, UN, NO, NO, MF, OK, OK, OK, NO, NO, NO, NO}}, // Int32 + {{LD, MF, MF, MF, MF, MF, OK, OK, OK, OK, LD, OK, OK, NO, NO, NO, MF, MF, OK, OK, MF, MF, OK, UN, NO, NO, MF, OK, OK, OK, NO, NO, NO, NO}}, // Uint32 + {{LD, MF, MF, MF, MF, MF, MF, OK, MF, LD, LD, OK, OK, NO, NO, NO, MF, MF, MF, MF, MF, MF, MF, UN, NO, NO, MF, MF, MF, MF, NO, NO, NO, NO}}, // Int64 + {{LD, MF, MF, MF, MF, MF, MF, MF, OK, LD, LD, OK, OK, NO, NO, NO, MF, MF, MF, MF, MF, MF, MF, UN, NO, NO, MF, MF, MF, MF, NO, NO, NO, NO}}, // Uint64 - {{FL, FL, FL, FL, FL, FL, FL, FL, FL, OK, LD, OK, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO}}, // Double - {{FL, FL, FL, FL, FL, FL, FL, FL, FL, OK, OK, OK, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO}}, // Float + {{FL, FL, FL, FL, FL, FL, FL, FL, FL, OK, LD, OK, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO}}, // Double + {{FL, FL, FL, FL, FL, FL, FL, FL, FL, OK, OK, OK, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO}}, // Float - {{MF, MF, MF, MF, MF, MF, MF, MF, MF, FL, FL, OK, MF, MF, MF, MF, MF, MF, MF, MF, MF, MF, MF, FL, FL, MF, MF, MF, MF, MF, MF, MF, MF}}, // String - {{MF, MF, MF, MF, MF, MF, MF, MF, MF, FL, FL, OK, OK, MF, MF, MF, MF, MF, MF, MF, MF, MF, MF, FL, FL, MF, MF, MF, MF, MF, MF, MF, MF}}, // Utf8 + {{MF, MF, MF, MF, MF, MF, MF, MF, MF, FL, FL, OK, MF, MF, MF, MF, MF, MF, MF, MF, MF, MF, MF, FL, FL, MF, MF, MF, MF, MF, MF, MF, MF, MF}}, // String + {{MF, MF, MF, MF, MF, MF, MF, MF, MF, FL, FL, OK, OK, MF, MF, MF, MF, MF, MF, MF, MF, MF, MF, FL, FL, MF, MF, MF, MF, MF, MF, MF, MF, MF}}, // Utf8 - {{NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, NO, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO}}, // Yson - {{NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, OK, NO, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, NO, NO, NO, NO, NO, NO, NO}}, // Json + {{NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, NO, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO}}, // Yson + {{NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, OK, NO, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, NO, NO, NO, NO, NO, NO, NO, NO}}, // Json - {{NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, OK, NO, NO, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO}}, // Uuid + {{NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, OK, NO, NO, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO}}, // Uuid - {{NO, MF, MF, MF, OK, OK, OK, OK, OK, OK, OK, OK, OK, NO, NO, NO, OK, OK, OK, NO, OK, OK, OK, NO, NO, NO, OK, OK, OK, NO, OK, OK, OK}}, // Date - {{NO, MF, MF, MF, MF, MF, OK, OK, OK, OK, LD, OK, OK, NO, NO, NO, LD, OK, OK, NO, LD, OK, OK, NO, NO, NO, LD, OK, OK, NO, LD, OK, OK}}, // Datetime - {{NO, MF, MF, MF, MF, MF, MF, OK, OK, LD, LD, OK, OK, NO, NO, NO, LD, LD, OK, NO, LD, LD, OK, NO, NO, NO, LD, LD, OK, NO, LD, LD, OK}}, // Timestamp + {{NO, MF, MF, MF, OK, OK, OK, OK, OK, OK, OK, OK, OK, NO, NO, NO, OK, OK, OK, NO, OK, OK, OK, NO, NO, NO, OK, OK, OK, NO, OK, OK, OK, NO}}, // Date + {{NO, MF, MF, MF, MF, MF, OK, OK, OK, OK, LD, OK, OK, NO, NO, NO, LD, OK, OK, NO, LD, OK, OK, NO, NO, NO, LD, OK, OK, NO, LD, OK, OK, NO}}, // Datetime + {{NO, MF, MF, MF, MF, MF, MF, OK, OK, LD, LD, OK, OK, NO, NO, NO, LD, LD, OK, NO, LD, LD, OK, NO, NO, NO, LD, LD, OK, NO, LD, LD, OK, NO}}, // Timestamp - {{NO, MF, MF, MF, MF, MF, MF, OK, MF, LD, LD, OK, OK, NO, NO, NO, NO, NO, NO, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, NO, NO, NO}}, // Interval + {{NO, MF, MF, MF, MF, MF, MF, OK, MF, LD, LD, OK, OK, NO, NO, NO, NO, NO, NO, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, NO, NO, NO, NO}}, // Interval - {{NO, MF, MF, MF, OK, OK, OK, OK, OK, OK, OK, OK, OK, NO, NO, NO, LD, LD, LD, NO, OK, OK, OK, NO, NO, NO, LD, LD, LD, NO, OK, OK, OK}}, // TzDate - {{NO, MF, MF, MF, MF, MF, OK, OK, OK, OK, LD, OK, OK, NO, NO, NO, LD, LD, LD, NO, LD, OK, OK, NO, NO, NO, LD, LD, LD, NO, LD, OK, OK}}, // TzDatetime - {{NO, MF, MF, MF, MF, MF, MF, OK, OK, LD, LD, OK, OK, NO, NO, NO, LD, LD, LD, NO, LD, LD, OK, NO, NO, NO, LD, LD, LD, NO, LD, LD, OK}}, // TzTimestamp + {{NO, MF, MF, MF, OK, OK, OK, OK, OK, OK, OK, OK, OK, NO, NO, NO, LD, LD, LD, NO, OK, OK, OK, NO, NO, NO, LD, LD, LD, NO, OK, OK, OK, NO}}, // TzDate + {{NO, MF, MF, MF, MF, MF, OK, OK, OK, OK, LD, OK, OK, NO, NO, NO, LD, LD, LD, NO, LD, OK, OK, NO, NO, NO, LD, LD, LD, NO, LD, OK, OK, NO}}, // TzDatetime + {{NO, MF, MF, MF, MF, MF, MF, OK, OK, LD, LD, OK, OK, NO, NO, NO, LD, LD, LD, NO, LD, LD, OK, NO, NO, NO, LD, LD, LD, NO, LD, LD, OK, NO}}, // TzTimestamp - {{NO, UN, UN, UN, UN, UN, UN, UN, UN, LD, LD, OK, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, UN, NO, NO, NO, NO, NO, NO, NO, NO, NO}}, // Decimal - {{NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, NO, NO, NO, NO, NO, NO, NO, NO}}, // DyNumber - {{NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, OK, NO, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, NO, NO, NO, NO, NO, NO, NO}}, // JsonDocument + {{NO, UN, UN, UN, UN, UN, UN, UN, UN, LD, LD, OK, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, UN, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO}}, // Decimal + {{NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO}}, // DyNumber + {{NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, OK, NO, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, NO, NO, NO, NO, NO, NO, NO, NO}}, // JsonDocument - {{NO, MF, MF, MF, MF, OK, MF, OK, MF, LD, OK, OK, OK, NO, NO, NO, MF, MF, MF, NO, MF, MF, MF, NO, NO, NO, OK, OK, OK, NO, OK, OK, OK}}, // Date32 - {{NO, MF, MF, MF, MF, MF, MF, OK, MF, LD, LD, OK, OK, NO, NO, NO, FL, MF, MF, NO, FL, MF, MF, NO, NO, NO, LD, OK, OK, NO, LD, OK, OK}}, // Datetime64 - {{NO, MF, MF, MF, MF, MF, MF, OK, MF, LD, LD, OK, OK, NO, NO, NO, FL, FL, MF, NO, FL, FL, MF, NO, NO, NO, LD, LD, OK, NO, LD, LD, OK}}, // Timestamp64 + {{NO, MF, MF, MF, MF, OK, MF, OK, MF, LD, OK, OK, OK, NO, NO, NO, MF, MF, MF, NO, MF, MF, MF, NO, NO, NO, OK, OK, OK, NO, OK, OK, OK, NO}}, // Date32 + {{NO, MF, MF, MF, MF, MF, MF, OK, MF, LD, LD, OK, OK, NO, NO, NO, FL, MF, MF, NO, FL, MF, MF, NO, NO, NO, LD, OK, OK, NO, LD, OK, OK, NO}}, // Datetime64 + {{NO, MF, MF, MF, MF, MF, MF, OK, MF, LD, LD, OK, OK, NO, NO, NO, FL, FL, MF, NO, FL, FL, MF, NO, NO, NO, LD, LD, OK, NO, LD, LD, OK, NO}}, // Timestamp64 - {{NO, MF, MF, MF, MF, MF, MF, OK, MF, LD, LD, OK, OK, NO, NO, NO, NO, NO, NO, MF, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, NO, NO, NO}}, // Interval64 + {{NO, MF, MF, MF, MF, MF, MF, OK, MF, LD, LD, OK, OK, NO, NO, NO, NO, NO, NO, MF, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, NO, NO, NO, NO}}, // Interval64 - {{NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, OK, NO, NO, NO, FL, FL, FL, NO, MF, MF, MF, NO, NO, NO, LD, LD, LD, NO, OK, OK, OK}}, // TzDate32 - {{NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, OK, NO, NO, NO, FL, FL, FL, NO, FL, MF, MF, NO, NO, NO, LD, LD, LD, NO, LD, OK, OK}}, // TzDatetime64 - {{NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, OK, NO, NO, NO, FL, FL, FL, NO, FL, FL, MF, NO, NO, NO, LD, LD, LD, NO, LD, LD, OK}}, // TzTimestamp64 + {{NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, OK, NO, NO, NO, FL, FL, FL, NO, MF, MF, MF, NO, NO, NO, LD, LD, LD, NO, OK, OK, OK, NO}}, // TzDate32 + {{NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, OK, NO, NO, NO, FL, FL, FL, NO, FL, MF, MF, NO, NO, NO, LD, LD, LD, NO, LD, OK, OK, NO}}, // TzDatetime64 + {{NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, OK, NO, NO, NO, FL, FL, FL, NO, FL, FL, MF, NO, NO, NO, LD, LD, LD, NO, LD, LD, OK, NO}}, // TzTimestamp64 + {{NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK, OK, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, NO, OK}}, // Rowid }}; } // namespace diff --git a/yql/essentials/public/udf/udf_data_type.h b/yql/essentials/public/udf/udf_data_type.h index b9bfec856fa6d..756df44a6eb83 100644 --- a/yql/essentials/public/udf/udf_data_type.h +++ b/yql/essentials/public/udf/udf_data_type.h @@ -141,6 +141,7 @@ class TYson {}; class TJson {}; class TUuid {}; class TJsonDocument {}; +class TRowid {}; class TDate {}; class TDatetime {}; @@ -175,6 +176,7 @@ constexpr i64 MAX_INTERVAL64 = MAX_TIMESTAMP64 - MIN_TIMESTAMP64; constexpr i32 MIN_YEAR32 = -144169; // inclusive constexpr i32 MAX_YEAR32 = 148108; // non-inclusive constexpr size_t UUID_SIZE = 16; +constexpr size_t ROWID_SIZE = 14; #define UDF_TYPE_ID_MAP(XX) \ XX(Bool, NYql::NProto::Bool, bool, CommonType, bool, 0) \ @@ -209,7 +211,8 @@ constexpr size_t UUID_SIZE = 16; XX(Interval64, NYql::NProto::Interval64, TInterval64, CommonType | TimeIntervalType | ExtDateType, i64, 0) \ XX(TzDate32, NYql::NProto::TzDate32, TTzDate32, CommonType | TzDateType | ExtDateType, i32, 0) \ XX(TzDatetime64, NYql::NProto::TzDatetime64, TTzDatetime64, CommonType | TzDateType | ExtDateType, i64, 0) \ - XX(TzTimestamp64, NYql::NProto::TzTimestamp64, TTzTimestamp64, CommonType | TzDateType | ExtDateType, i64, 0) + XX(TzTimestamp64, NYql::NProto::TzTimestamp64, TTzTimestamp64, CommonType | TzDateType | ExtDateType, i64, 0) \ + XX(Rowid, NYql::NProto::Rowid, TRowid, CommonType, TRowid, 0) #define UDF_TYPE_ID(xName, xTypeId, xType, xFeatures, xLayoutType, xParamsCount) \ template <> \ diff --git a/yql/essentials/public/udf/udf_type_ops.h b/yql/essentials/public/udf/udf_type_ops.h index c74898107d950..4852f3a907aa2 100644 --- a/yql/essentials/public/udf/udf_type_ops.h +++ b/yql/essentials/public/udf/udf_type_ops.h @@ -132,6 +132,11 @@ inline THashType GetValueHash(const TUnboxedValuePod& value) { return GetStringHash(value); } +template <> +inline THashType GetValueHash(const TUnboxedValuePod& value) { + return GetStringHash(value); +} + template <> inline THashType GetValueHash(const TUnboxedValuePod&) { Y_ABORT("Yson isn't hashable."); @@ -359,6 +364,11 @@ inline int CompareValues(const TUnboxedValuePod& lhs, const TUn return CompareStrings(lhs, rhs); } +template <> +inline int CompareValues(const TUnboxedValuePod& lhs, const TUnboxedValuePod& rhs) { + return CompareStrings(lhs, rhs); +} + template <> inline int CompareValues(const TUnboxedValuePod&, const TUnboxedValuePod&) { Y_ABORT("Yson isn't comparable."); @@ -564,6 +574,11 @@ inline bool EquateValues(const TUnboxedValuePod& lhs, const TUn return EquateStrings(lhs, rhs); } +template <> +inline bool EquateValues(const TUnboxedValuePod& lhs, const TUnboxedValuePod& rhs) { + return EquateStrings(lhs, rhs); +} + template <> inline bool EquateValues(const TUnboxedValuePod&, const TUnboxedValuePod&) { Y_ABORT("Yson isn't comparable."); diff --git a/yql/essentials/types/rowid/rowid.cpp b/yql/essentials/types/rowid/rowid.cpp new file mode 100644 index 0000000000000..831b324d23dfe --- /dev/null +++ b/yql/essentials/types/rowid/rowid.cpp @@ -0,0 +1,56 @@ +#include "rowid.h" + +#include + +#include + +namespace NKikimr::NRowid { +namespace { + +bool DecodeRowidBase64(TStringBuf buf, char* out) { + if (buf.size() != ROWID_BASE64_LEN) { + return false; + } + + try { + constexpr size_t paddedLen = ROWID_BASE64_LEN + 1; + char padded[paddedLen]; + std::memcpy(padded, buf.data(), ROWID_BASE64_LEN); + padded[ROWID_BASE64_LEN] = '='; + + char decoded[Base64DecodeBufSize(paddedLen)]; + const size_t n = Base64StrictDecode(decoded, padded, padded + paddedLen); + if (n != ROWID_LEN) { + return false; + } + std::memcpy(out, decoded, ROWID_LEN); + return true; + } catch (...) { + return false; + } +} + +} // namespace + +bool IsValidRowidBase64(TStringBuf buf) { + char decoded[ROWID_LEN]; + return DecodeRowidBase64(buf, decoded); +} + +TString RowidBytesToBase64(TStringBuf in) { + Y_ABORT_UNLESS(in.size() == ROWID_LEN); + return Base64EncodeNoPadding(in); +} + +void RowidBytesToBase64(TStringBuf in, IOutputStream& out) { + Y_ABORT_UNLESS(in.size() == ROWID_LEN); + char encoded[Base64EncodeBufSize(ROWID_LEN)]; + const auto encodedBuf = Base64EncodeNoPadding(in, encoded); + out.Write(encodedBuf.data(), encodedBuf.size()); +} + +bool ParseRowidBase64(TStringBuf buf, char* out) { + return DecodeRowidBase64(buf, out); +} + +} // namespace NKikimr::NRowid diff --git a/yql/essentials/types/rowid/rowid.h b/yql/essentials/types/rowid/rowid.h new file mode 100644 index 0000000000000..291e52f7ac869 --- /dev/null +++ b/yql/essentials/types/rowid/rowid.h @@ -0,0 +1,31 @@ +#pragma once + +#include +#include +#include + +#include +#include + +class IOutputStream; + +namespace NKikimr::NRowid { + +static constexpr ui32 ROWID_LEN = 14; +// Base64 without padding for 14 bytes is 19 characters (14*8/6 rounded up). +static constexpr ui32 ROWID_BASE64_LEN = 19; + +inline bool IsValidRowidBytes(TStringBuf buf) { + return buf.size() == ROWID_LEN; +} + +bool IsValidRowidBase64(TStringBuf buf); +TString RowidBytesToBase64(TStringBuf in); +void RowidBytesToBase64(TStringBuf in, IOutputStream& out); +bool ParseRowidBase64(TStringBuf buf, char* out /*[ROWID_LEN]*/); + +inline bool ParseRowidBase64ToArray(TStringBuf buf, std::array& out) { + return ParseRowidBase64(buf, out.data()); +} + +} // namespace NKikimr::NRowid diff --git a/yql/essentials/types/rowid/ya.make b/yql/essentials/types/rowid/ya.make new file mode 100644 index 0000000000000..44736db1b2b8a --- /dev/null +++ b/yql/essentials/types/rowid/ya.make @@ -0,0 +1,11 @@ +LIBRARY() + +SRCS( + rowid.cpp +) + +PEERDIR( + library/cpp/string_utils/base64 +) + +END() diff --git a/yql/essentials/types/ya.make b/yql/essentials/types/ya.make index fc082816a95f8..4354efef9c015 100644 --- a/yql/essentials/types/ya.make +++ b/yql/essentials/types/ya.make @@ -1,6 +1,7 @@ RECURSE( binary_json dynumber + rowid uuid ) diff --git a/yt/yql/providers/yt/lib/expr_traits/yql_expr_traits.cpp b/yt/yql/providers/yt/lib/expr_traits/yql_expr_traits.cpp index a995230f70a07..2d93403f830ab 100644 --- a/yt/yql/providers/yt/lib/expr_traits/yql_expr_traits.cpp +++ b/yt/yql/providers/yt/lib/expr_traits/yql_expr_traits.cpp @@ -251,6 +251,7 @@ namespace NYql { case EDataSlot::Json: case EDataSlot::Decimal: case EDataSlot::Uuid: + case EDataSlot::Rowid: case EDataSlot::TzDate: case EDataSlot::TzDatetime: case EDataSlot::TzTimestamp: