eldenmoon commented on code in PR #66204: URL: https://github.com/apache/doris/pull/66204#discussion_r3671857195
########## be/src/storage/segment/variant/v2/variant_assembler_value.cpp: ########## @@ -0,0 +1,1225 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +#include <array> +#include <cstring> +#include <limits> +#include <utility> + +#include "common/check.h" +#include "common/exception.h" +#include "core/assert_cast.h" +#include "core/column/column_decimal.h" +#include "core/column/column_nullable.h" +#include "core/column/column_vector.h" +#include "core/column/variant_v2/column_variant_v2_typed_column.h" +#include "core/data_type/data_type_array.h" +#include "core/data_type/data_type_factory.hpp" +#include "core/data_type/data_type_nullable.h" +#include "core/typeid_cast.h" +#include "core/value/variant/variant_parquet_encoding.h" +#include "exec/common/format_ip.h" +#include "exprs/function/parse/variant_jsonb_parse.h" +#include "storage/segment/variant/v2/variant_assembler_internal.h" +#include "util/utf8_check.h" + +namespace doris::segment_v2::variant_v2::variant_assembler_internal { +namespace { + +class CellCursor { +public: + explicit CellCursor(StringRef cell) + : _current(reinterpret_cast<const uint8_t*>(cell.data)), _remaining(cell.size) { + if (cell.data == nullptr && cell.size != 0) { + throw Exception(ErrorCode::CORRUPTION, + "Variant storage cell has a null pointer for {} bytes", cell.size); + } + } + + template <typename T> + T read(std::string_view description) { + require(sizeof(T), description); + T value; + std::memcpy(&value, _current, sizeof(T)); + _current += sizeof(T); + _remaining -= sizeof(T); + return value; + } + + StringRef read_bytes(size_t size, std::string_view description) { + require(size, description); + const StringRef result {reinterpret_cast<const char*>(_current), size}; + _current += size; + _remaining -= size; + return result; + } + + size_t remaining() const noexcept { return _remaining; } + bool empty() const noexcept { return _remaining == 0; } + +private: + void require(size_t size, std::string_view description) const { + if (size > _remaining) { + throw Exception(ErrorCode::CORRUPTION, + "Truncated Variant storage cell while reading {}: need {} bytes, " + "have {}", + description, size, _remaining); + } + } + + const uint8_t* _current; + size_t _remaining; +}; + +template <typename Integer> +uint32_t decimal_digits(Integer value) { + uint32_t result = 0; + do { + value /= 10; + ++result; + } while (value != 0); + return result; +} + +template <typename Integer> +void validate_decimal(Integer value, uint8_t precision, uint8_t scale, uint8_t maximum, + std::string_view description) { + if (precision == 0 || precision > maximum || scale > precision) { + throw Exception(ErrorCode::CORRUPTION, + "Variant storage {} precision/scale {}/{} is invalid", description, + precision, scale); + } + if (decimal_digits(value) > precision) { + throw Exception(ErrorCode::CORRUPTION, + "Variant storage {} value exceeds declared precision {}", description, + precision); + } +} + +void validate_depth(uint32_t depth) { + if (depth > VARIANT_MAX_NESTING_DEPTH) { + throw Exception(ErrorCode::CORRUPTION, + "Variant storage cell exceeds maximum nesting depth {}", + VARIANT_MAX_NESTING_DEPTH); + } +} + +void validate_container_depth(uint32_t depth) { + if (depth >= VARIANT_MAX_NESTING_DEPTH) { + throw Exception(ErrorCode::CORRUPTION, + "Variant storage container exceeds maximum nesting depth {}", + VARIANT_MAX_NESTING_DEPTH); + } +} + +template <typename DateValue> +int32_t storage_date_days(DateValue value, std::string_view description) { + if (!value.is_valid_date()) { + throw Exception(ErrorCode::CORRUPTION, "Invalid {} in Variant storage cell", description); + } + return variant_days_since_epoch(value, 0, description); +} + +template <typename DateTimeValue> +int64_t storage_timestamp_micros(DateTimeValue value, std::string_view description) { + if (!value.is_valid_date()) { + throw Exception(ErrorCode::CORRUPTION, "Invalid {} in Variant storage cell", description); + } + return variant_timestamp_micros(value, 0, description); +} + +VecDateTimeValue read_legacy_temporal(CellCursor& cursor, FieldType type) { + const bool is_date = type == FieldType::OLAP_FIELD_TYPE_DATE; + DORIS_CHECK(is_date || type == FieldType::OLAP_FIELD_TYPE_DATETIME); + const auto value = cursor.read<VecDateTimeValue>(is_date ? "legacy DATE" : "legacy DATETIME"); + if (is_date) { + static_cast<void>(storage_date_days(value, "legacy DATE")); + } else { + static_cast<void>(storage_timestamp_micros(value, "legacy DATETIME")); + } + return value; +} + +struct LegacyDecimalCell { + uint8_t precision; + uint8_t scale; + __int128 value; +}; + +LegacyDecimalCell read_legacy_decimal(CellCursor& cursor) { + LegacyDecimalCell result { + .precision = cursor.read<uint8_t>("legacy DecimalV2 precision"), + .scale = cursor.read<uint8_t>("legacy DecimalV2 scale"), + .value = cursor.read<__int128>("legacy DecimalV2 value"), + }; + if (result.precision == 0 || result.precision > DecimalV2Value::PRECISION || + result.scale > result.precision) { + throw Exception(ErrorCode::CORRUPTION, + "Variant storage legacy DecimalV2 precision/scale {}/{} is invalid", + result.precision, result.scale); + } + // DecimalV2 always stores a fixed nine-digit fractional coefficient. Preserve the serialized + // SerDe metadata for typed identity, but do not use it to rescale the coefficient. + validate_decimal(result.value, DecimalV2Value::PRECISION, DecimalV2Value::SCALE, + DecimalV2Value::PRECISION, "legacy DecimalV2"); + return result; +} + +void add_largeint(VariantBatchBuilder::Row& output, __int128 value) { + output.add_largeint(value); +} + +void add_decimal256(VariantBatchBuilder::Row& output, wide::Int256 value, uint8_t scale) { + const std::string text = Decimal256 {value}.to_string(scale); + output.add_string(StringRef(text)); +} + +void add_ipv4(VariantBatchBuilder::Row& output, IPv4 value) { + std::array<char, IPV4_MAX_TEXT_LENGTH + 1> buffer {}; + char* end = buffer.data(); + const auto* address = reinterpret_cast<const unsigned char*>(&value); + format_ipv4(address, end); + output.add_string({buffer.data(), static_cast<size_t>(end - buffer.data())}); +} + +void add_ipv6(VariantBatchBuilder::Row& output, IPv6 value) { + std::array<char, IPV6_MAX_TEXT_LENGTH + 1> buffer {}; + char* end = buffer.data(); + format_ipv6(reinterpret_cast<unsigned char*>(&value), end); + output.add_string({buffer.data(), static_cast<size_t>(end - buffer.data())}); +} + +StringRef read_storage_string(CellCursor& cursor) { + const size_t size = cursor.read<size_t>("string size"); + const StringRef value = cursor.read_bytes(size, "string payload"); + if (value.size != 0 && !validate_utf8(value.data, value.size)) { + throw Exception(ErrorCode::CORRUPTION, "Variant storage string is not valid UTF-8"); + } + return value; +} + +// Keep the exhaustive storage FieldType decoder together so validation and byte consumption for +// every wire tag remain auditable in one dispatch table. +// NOLINTNEXTLINE(readability-function-size) +void decode_storage_value(CellCursor& cursor, VariantBatchBuilder::Row& output, uint32_t depth) { + validate_depth(depth); + const auto type = static_cast<FieldType>(cursor.read<uint8_t>("field type")); + switch (type) { + case FieldType::OLAP_FIELD_TYPE_NONE: + output.add_null(); + return; + case FieldType::OLAP_FIELD_TYPE_BOOL: { + const uint8_t value = cursor.read<uint8_t>("boolean"); + if (value > 1) { + throw Exception(ErrorCode::CORRUPTION, "Invalid Variant storage boolean byte {}", + value); + } + output.add_bool(value != 0); + return; + } + case FieldType::OLAP_FIELD_TYPE_TINYINT: + output.add_int(cursor.read<int8_t>("tinyint")); + return; + case FieldType::OLAP_FIELD_TYPE_SMALLINT: + output.add_int(cursor.read<int16_t>("smallint")); + return; + case FieldType::OLAP_FIELD_TYPE_INT: + output.add_int(cursor.read<int32_t>("int")); + return; + case FieldType::OLAP_FIELD_TYPE_BIGINT: + output.add_int(cursor.read<int64_t>("bigint")); + return; + case FieldType::OLAP_FIELD_TYPE_LARGEINT: + add_largeint(output, cursor.read<__int128>("largeint")); + return; + case FieldType::OLAP_FIELD_TYPE_FLOAT: + output.add_float(cursor.read<float>("float")); + return; + case FieldType::OLAP_FIELD_TYPE_DOUBLE: + output.add_double(cursor.read<double>("double")); + return; + case FieldType::OLAP_FIELD_TYPE_STRING: + output.add_string(read_storage_string(cursor)); + return; + case FieldType::OLAP_FIELD_TYPE_JSONB: { + const size_t size = cursor.read<size_t>("JSONB size"); + jsonb_to_variant(cursor.read_bytes(size, "JSONB payload"), output, depth); + return; + } + case FieldType::OLAP_FIELD_TYPE_ARRAY: { + validate_container_depth(depth); + const size_t count = cursor.read<size_t>("array element count"); + if (count > cursor.remaining()) { + throw Exception(ErrorCode::CORRUPTION, + "Variant storage array count {} exceeds remaining {} bytes", count, + cursor.remaining()); + } + auto array = output.start_array(); + for (size_t index = 0; index < count; ++index) { + decode_storage_value(cursor, output, depth + 1); + } + array.finish(); + return; + } + case FieldType::OLAP_FIELD_TYPE_IPV4: + add_ipv4(output, cursor.read<IPv4>("IPv4")); + return; + case FieldType::OLAP_FIELD_TYPE_IPV6: + add_ipv6(output, cursor.read<IPv6>("IPv6")); + return; + case FieldType::OLAP_FIELD_TYPE_DATE: { + const auto value = read_legacy_temporal(cursor, type); + output.add_date(storage_date_days(value, "legacy DATE")); + return; + } + case FieldType::OLAP_FIELD_TYPE_DATETIME: { + const auto value = read_legacy_temporal(cursor, type); + output.add_timestamp_micros(storage_timestamp_micros(value, "legacy DATETIME"), false); + return; + } + case FieldType::OLAP_FIELD_TYPE_DATEV2: { + const UInt32 raw = cursor.read<UInt32>("DateV2"); + const auto value = binary_cast<UInt32, DateV2Value<DateV2ValueType>>(raw); + output.add_date(storage_date_days(value, "DATEV2")); + return; + } + case FieldType::OLAP_FIELD_TYPE_DATETIMEV2: + case FieldType::OLAP_FIELD_TYPE_TIMESTAMPTZ: { + const uint8_t scale = cursor.read<uint8_t>("timestamp scale"); + if (scale > 6) { + throw Exception(ErrorCode::CORRUPTION, "Variant storage timestamp scale {} exceeds 6", + scale); + } + const UInt64 raw = cursor.read<UInt64>("timestamp value"); + if (type == FieldType::OLAP_FIELD_TYPE_DATETIMEV2) { + const auto value = binary_cast<UInt64, DateV2Value<DateTimeV2ValueType>>(raw); + output.add_timestamp_micros(storage_timestamp_micros(value, "DATETIMEV2"), false); + } else { + const auto value = binary_cast<UInt64, TimestampTzValue>(raw); + output.add_timestamp_micros(storage_timestamp_micros(value, "TIMESTAMPTZ"), true); + } + return; + } + case FieldType::OLAP_FIELD_TYPE_DECIMAL: { + const auto value = read_legacy_decimal(cursor); + output.add_decimal(value.value, DecimalV2Value::SCALE, 16); + return; + } + case FieldType::OLAP_FIELD_TYPE_DECIMAL32: { + const uint8_t precision = cursor.read<uint8_t>("Decimal32 precision"); + const uint8_t scale = cursor.read<uint8_t>("Decimal32 scale"); + const int32_t value = cursor.read<int32_t>("Decimal32 value"); + validate_decimal(value, precision, scale, 9, "Decimal32"); + output.add_decimal(value, scale, 4); + return; + } + case FieldType::OLAP_FIELD_TYPE_DECIMAL64: { + const uint8_t precision = cursor.read<uint8_t>("Decimal64 precision"); + const uint8_t scale = cursor.read<uint8_t>("Decimal64 scale"); + const int64_t value = cursor.read<int64_t>("Decimal64 value"); + validate_decimal(value, precision, scale, 18, "Decimal64"); + output.add_decimal(value, scale, 8); + return; + } + case FieldType::OLAP_FIELD_TYPE_DECIMAL128I: { + const uint8_t precision = cursor.read<uint8_t>("Decimal128 precision"); + const uint8_t scale = cursor.read<uint8_t>("Decimal128 scale"); + const __int128 value = cursor.read<__int128>("Decimal128 value"); + validate_decimal(value, precision, scale, 38, "Decimal128"); + output.add_decimal(value, scale, 16); + return; + } + case FieldType::OLAP_FIELD_TYPE_DECIMAL256: { + const uint8_t precision = cursor.read<uint8_t>("Decimal256 precision"); + const uint8_t scale = cursor.read<uint8_t>("Decimal256 scale"); + const wide::Int256 value = cursor.read<wide::Int256>("Decimal256 value"); + validate_decimal(value, precision, scale, 76, "Decimal256"); + add_decimal256(output, value, scale); + return; + } + default: + throw Exception(ErrorCode::CORRUPTION, "Unknown Variant storage FieldType {}", + static_cast<uint8_t>(type)); + } +} + +template <typename ColumnType> +bool matches(const IColumn* column) { + return check_and_get_column<ColumnType>(column) != nullptr; +} + +bool supported_column_shape(PrimitiveType type, const IColumn* column) { + switch (type) { + case TYPE_BOOLEAN: + return matches<ColumnUInt8>(column); + case TYPE_TINYINT: + return matches<ColumnInt8>(column); + case TYPE_SMALLINT: + return matches<ColumnInt16>(column); + case TYPE_INT: + return matches<ColumnInt32>(column); + case TYPE_BIGINT: + return matches<ColumnInt64>(column); + case TYPE_LARGEINT: + return matches<ColumnInt128>(column); + case TYPE_FLOAT: + return matches<ColumnFloat32>(column); + case TYPE_DOUBLE: + return matches<ColumnFloat64>(column); + case TYPE_DECIMALV2: + return matches<ColumnDecimal128V2>(column); + case TYPE_DECIMAL32: + return matches<ColumnDecimal32>(column); + case TYPE_DECIMAL64: + return matches<ColumnDecimal64>(column); + case TYPE_DECIMAL128I: + return matches<ColumnDecimal128V3>(column); + case TYPE_DECIMAL256: + return matches<ColumnDecimal256>(column); + case TYPE_DATE: + return matches<ColumnDate>(column); + case TYPE_DATEV2: + return matches<ColumnDateV2>(column); + case TYPE_DATETIME: + return matches<ColumnDateTime>(column); + case TYPE_DATETIMEV2: + return matches<ColumnDateTimeV2>(column); + case TYPE_TIMESTAMPTZ: + return matches<ColumnTimeStampTz>(column); + case TYPE_CHAR: + case TYPE_VARCHAR: + case TYPE_STRING: + case TYPE_JSONB: + return matches<ColumnString>(column); + case TYPE_IPV4: + return matches<ColumnIPv4>(column); + case TYPE_IPV6: + return matches<ColumnIPv6>(column); + case TYPE_ARRAY: + return matches<ColumnArray>(column); + default: + return false; + } +} + +// This is a flat PrimitiveType dispatch; splitting it would scatter the concrete-column contract. +// NOLINTNEXTLINE(readability-function-size) +void append_materialized_scalar(const PreparedColumn& column, size_t row, + VariantBatchBuilder::Row& output, uint32_t depth) { + const uint8_t scale = column.scale; + switch (column.primitive) { + case TYPE_BOOLEAN: + output.add_bool(assert_cast<const ColumnUInt8&, TypeCheckOnRelease::DISABLE>(*column.data) + .get_data()[row] != 0); + return; + case TYPE_TINYINT: + output.add_int(assert_cast<const ColumnInt8&, TypeCheckOnRelease::DISABLE>(*column.data) + .get_data()[row]); + return; + case TYPE_SMALLINT: + output.add_int(assert_cast<const ColumnInt16&, TypeCheckOnRelease::DISABLE>(*column.data) + .get_data()[row]); + return; + case TYPE_INT: + output.add_int(assert_cast<const ColumnInt32&, TypeCheckOnRelease::DISABLE>(*column.data) + .get_data()[row]); + return; + case TYPE_BIGINT: + output.add_int(assert_cast<const ColumnInt64&, TypeCheckOnRelease::DISABLE>(*column.data) + .get_data()[row]); + return; + case TYPE_LARGEINT: + add_largeint(output, + assert_cast<const ColumnInt128&, TypeCheckOnRelease::DISABLE>(*column.data) + .get_data()[row]); + return; + case TYPE_FLOAT: + output.add_float( + assert_cast<const ColumnFloat32&, TypeCheckOnRelease::DISABLE>(*column.data) + .get_data()[row]); + return; + case TYPE_DOUBLE: + output.add_double( + assert_cast<const ColumnFloat64&, TypeCheckOnRelease::DISABLE>(*column.data) + .get_data()[row]); + return; + case TYPE_DECIMALV2: + output.add_decimal( + assert_cast<const ColumnDecimal128V2&, TypeCheckOnRelease::DISABLE>(*column.data) + .get_data()[row] + .value(), + scale, 16); + return; + case TYPE_DECIMAL32: + output.add_decimal( + assert_cast<const ColumnDecimal32&, TypeCheckOnRelease::DISABLE>(*column.data) + .get_data()[row] + .value, + scale, 4); + return; + case TYPE_DECIMAL64: + output.add_decimal( + assert_cast<const ColumnDecimal64&, TypeCheckOnRelease::DISABLE>(*column.data) + .get_data()[row] + .value, + scale, 8); + return; + case TYPE_DECIMAL128I: + output.add_decimal( + assert_cast<const ColumnDecimal128V3&, TypeCheckOnRelease::DISABLE>(*column.data) + .get_data()[row] + .value, + scale, 16); + return; + case TYPE_DECIMAL256: + add_decimal256( + output, + assert_cast<const ColumnDecimal256&, TypeCheckOnRelease::DISABLE>(*column.data) + .get_data()[row] + .value, + scale); + return; + case TYPE_DATE: + output.add_date(storage_date_days( + assert_cast<const ColumnDate&, TypeCheckOnRelease::DISABLE>(*column.data) + .get_data()[row], + "DATE")); + return; + case TYPE_DATEV2: + output.add_date(storage_date_days( + assert_cast<const ColumnDateV2&, TypeCheckOnRelease::DISABLE>(*column.data) + .get_data()[row], + "DATEV2")); + return; + case TYPE_DATETIME: + output.add_timestamp_micros( + storage_timestamp_micros( + assert_cast<const ColumnDateTime&, TypeCheckOnRelease::DISABLE>( + *column.data) + .get_data()[row], + "DATETIME"), + false); + return; + case TYPE_DATETIMEV2: + output.add_timestamp_micros( + storage_timestamp_micros( + assert_cast<const ColumnDateTimeV2&, TypeCheckOnRelease::DISABLE>( + *column.data) + .get_data()[row], + "DATETIMEV2"), + false); + return; + case TYPE_TIMESTAMPTZ: + output.add_timestamp_micros( + storage_timestamp_micros( + assert_cast<const ColumnTimeStampTz&, TypeCheckOnRelease::DISABLE>( + *column.data) + .get_data()[row], + "TIMESTAMPTZ"), + true); + return; + case TYPE_CHAR: + case TYPE_VARCHAR: + case TYPE_STRING: + output.add_string( + assert_cast<const ColumnString&, TypeCheckOnRelease::DISABLE>(*column.data) + .get_data_at(row)); + return; + case TYPE_JSONB: + jsonb_to_variant(assert_cast<const ColumnString&, TypeCheckOnRelease::DISABLE>(*column.data) + .get_data_at(row), + output, depth); + return; + case TYPE_IPV4: + add_ipv4(output, assert_cast<const ColumnIPv4&, TypeCheckOnRelease::DISABLE>(*column.data) + .get_data()[row]); + return; + case TYPE_IPV6: + add_ipv6(output, assert_cast<const ColumnIPv6&, TypeCheckOnRelease::DISABLE>(*column.data) + .get_data()[row]); + return; + default: + throw Exception(ErrorCode::CORRUPTION, "Unsupported materialized Variant type {}", + column.type->get_name()); + } +} + +size_t fixed_payload_size(FieldType type) { + switch (type) { + case FieldType::OLAP_FIELD_TYPE_BOOL: + case FieldType::OLAP_FIELD_TYPE_TINYINT: + return 1; + case FieldType::OLAP_FIELD_TYPE_SMALLINT: + return 2; + case FieldType::OLAP_FIELD_TYPE_INT: + case FieldType::OLAP_FIELD_TYPE_FLOAT: + case FieldType::OLAP_FIELD_TYPE_IPV4: + case FieldType::OLAP_FIELD_TYPE_DATEV2: + return 4; + case FieldType::OLAP_FIELD_TYPE_BIGINT: + case FieldType::OLAP_FIELD_TYPE_DOUBLE: + return 8; + case FieldType::OLAP_FIELD_TYPE_LARGEINT: + case FieldType::OLAP_FIELD_TYPE_IPV6: + return 16; + default: + return 0; + } +} + +void require_exact(CellCursor& cursor) { + if (!cursor.empty()) { + throw Exception(ErrorCode::CORRUPTION, "Variant storage cell has {} trailing bytes", + cursor.remaining()); + } +} + +// Keep typed eligibility validation adjacent for all storage scalar tags. +// NOLINTNEXTLINE(readability-function-size) +bool inspect_typed_cell(CellCursor& cursor, CellSignature* signature) { + signature->type = static_cast<FieldType>(cursor.read<uint8_t>("field type")); + signature->typed = true; + if (signature->type == FieldType::OLAP_FIELD_TYPE_BOOL) { + const uint8_t value = cursor.read<uint8_t>("boolean"); + if (value > 1) { + throw Exception(ErrorCode::CORRUPTION, "Invalid Variant storage boolean byte {}", + value); + } + require_exact(cursor); + return true; + } + if (signature->type == FieldType::OLAP_FIELD_TYPE_DATEV2) { + const UInt32 raw = cursor.read<UInt32>("DateV2"); + const auto value = binary_cast<UInt32, DateV2Value<DateV2ValueType>>(raw); + static_cast<void>(storage_date_days(value, "DATEV2")); + require_exact(cursor); + return true; + } + if (signature->type == FieldType::OLAP_FIELD_TYPE_DATE || + signature->type == FieldType::OLAP_FIELD_TYPE_DATETIME) { + static_cast<void>(read_legacy_temporal(cursor, signature->type)); + require_exact(cursor); + return true; + } + if (signature->type == FieldType::OLAP_FIELD_TYPE_DECIMAL) { + const auto value = read_legacy_decimal(cursor); + signature->precision = value.precision; + signature->scale = value.scale; + require_exact(cursor); + return true; + } + const size_t fixed = fixed_payload_size(signature->type); + if (fixed != 0) { + static_cast<void>(cursor.read_bytes(fixed, "fixed scalar payload")); + require_exact(cursor); + return true; + } + switch (signature->type) { + case FieldType::OLAP_FIELD_TYPE_STRING: + static_cast<void>(read_storage_string(cursor)); + require_exact(cursor); + return true; + case FieldType::OLAP_FIELD_TYPE_DATETIMEV2: + case FieldType::OLAP_FIELD_TYPE_TIMESTAMPTZ: + signature->scale = cursor.read<uint8_t>("timestamp scale"); + if (signature->scale > 6) { + throw Exception(ErrorCode::CORRUPTION, "Variant storage timestamp scale {} exceeds 6", + signature->scale); + } + if (signature->type == FieldType::OLAP_FIELD_TYPE_DATETIMEV2) { + const auto value = binary_cast<UInt64, DateV2Value<DateTimeV2ValueType>>( + cursor.read<UInt64>("timestamp value")); + static_cast<void>(storage_timestamp_micros(value, "DATETIMEV2")); + } else { + const auto value = + binary_cast<UInt64, TimestampTzValue>(cursor.read<UInt64>("timestamp value")); + static_cast<void>(storage_timestamp_micros(value, "TIMESTAMPTZ")); + } + require_exact(cursor); + return true; + case FieldType::OLAP_FIELD_TYPE_DECIMAL32: + case FieldType::OLAP_FIELD_TYPE_DECIMAL64: + case FieldType::OLAP_FIELD_TYPE_DECIMAL128I: + case FieldType::OLAP_FIELD_TYPE_DECIMAL256: { + signature->precision = cursor.read<uint8_t>("decimal precision"); + signature->scale = cursor.read<uint8_t>("decimal scale"); + uint8_t maximum = 0; + if (signature->type == FieldType::OLAP_FIELD_TYPE_DECIMAL32) { + maximum = 9; + } else if (signature->type == FieldType::OLAP_FIELD_TYPE_DECIMAL64) { + maximum = 18; + } else if (signature->type == FieldType::OLAP_FIELD_TYPE_DECIMAL128I) { + maximum = 38; + } else { + maximum = 76; + } + if (signature->precision == 0 || signature->precision > maximum || + signature->scale > signature->precision) { + throw Exception(ErrorCode::CORRUPTION, + "Variant storage decimal precision/scale {}/{} is invalid", + signature->precision, signature->scale); + } + if (signature->type == FieldType::OLAP_FIELD_TYPE_DECIMAL32) { + validate_decimal(cursor.read<int32_t>("Decimal32 value"), signature->precision, + signature->scale, maximum, "Decimal32"); + } else if (signature->type == FieldType::OLAP_FIELD_TYPE_DECIMAL64) { + validate_decimal(cursor.read<int64_t>("Decimal64 value"), signature->precision, + signature->scale, maximum, "Decimal64"); + } else if (signature->type == FieldType::OLAP_FIELD_TYPE_DECIMAL128I) { + validate_decimal(cursor.read<__int128>("Decimal128 value"), signature->precision, + signature->scale, maximum, "Decimal128"); + } else { + validate_decimal(cursor.read<wide::Int256>("Decimal256 value"), signature->precision, + signature->scale, maximum, "Decimal256"); + } + require_exact(cursor); + return true; + } + case FieldType::OLAP_FIELD_TYPE_NONE: + case FieldType::OLAP_FIELD_TYPE_JSONB: + case FieldType::OLAP_FIELD_TYPE_ARRAY: + signature->typed = false; + return false; + default: + throw Exception(ErrorCode::CORRUPTION, "Unknown Variant storage FieldType {}", + static_cast<uint8_t>(signature->type)); + } +} + +template <typename Column> +Column& typed_output(IColumn* output) { + return assert_cast<Column&, TypeCheckOnRelease::DISABLE>(*output); +} + +// This mirrors inspect_typed_cell as one exhaustive scalar publication table. +// NOLINTNEXTLINE(readability-function-size) +void append_cell_to_typed(StringRef cell, const CellSignature& signature, IColumn* output) { + CellCursor cursor(cell); + const auto type = static_cast<FieldType>(cursor.read<uint8_t>("field type")); + if (type != signature.type) { + throw Exception(ErrorCode::CORRUPTION, + "Variant binary type changed while publishing typed output"); + } + switch (type) { + case FieldType::OLAP_FIELD_TYPE_BOOL: { + const uint8_t value = cursor.read<uint8_t>("boolean"); + if (value > 1) { + throw Exception(ErrorCode::CORRUPTION, "Invalid Variant storage boolean byte {}", + value); + } + typed_output<ColumnUInt8>(output).insert_value(value); + break; + } + case FieldType::OLAP_FIELD_TYPE_TINYINT: + typed_output<ColumnInt8>(output).insert_value(cursor.read<int8_t>("tinyint")); + break; + case FieldType::OLAP_FIELD_TYPE_SMALLINT: + typed_output<ColumnInt16>(output).insert_value(cursor.read<int16_t>("smallint")); + break; + case FieldType::OLAP_FIELD_TYPE_INT: + typed_output<ColumnInt32>(output).insert_value(cursor.read<int32_t>("int")); + break; + case FieldType::OLAP_FIELD_TYPE_BIGINT: + typed_output<ColumnInt64>(output).insert_value(cursor.read<int64_t>("bigint")); + break; + case FieldType::OLAP_FIELD_TYPE_LARGEINT: + typed_output<ColumnInt128>(output).insert_value(cursor.read<__int128>("largeint")); + break; + case FieldType::OLAP_FIELD_TYPE_FLOAT: + typed_output<ColumnFloat32>(output).insert_value(cursor.read<float>("float")); + break; + case FieldType::OLAP_FIELD_TYPE_DOUBLE: + typed_output<ColumnFloat64>(output).insert_value(cursor.read<double>("double")); + break; + case FieldType::OLAP_FIELD_TYPE_STRING: { + const StringRef value = read_storage_string(cursor); + typed_output<ColumnString>(output).insert_data(value.data, value.size); + break; + } + case FieldType::OLAP_FIELD_TYPE_IPV4: + typed_output<ColumnIPv4>(output).insert_value(cursor.read<IPv4>("IPv4")); + break; + case FieldType::OLAP_FIELD_TYPE_IPV6: + typed_output<ColumnIPv6>(output).insert_value(cursor.read<IPv6>("IPv6")); + break; + case FieldType::OLAP_FIELD_TYPE_DATE: { + const auto value = read_legacy_temporal(cursor, type); + typed_output<ColumnDate>(output).insert_value(value); + break; + } + case FieldType::OLAP_FIELD_TYPE_DATETIME: { + const auto value = read_legacy_temporal(cursor, type); + typed_output<ColumnDateTime>(output).insert_value(value); + break; + } + case FieldType::OLAP_FIELD_TYPE_DATEV2: { + const UInt32 raw = cursor.read<UInt32>("DateV2"); + const auto value = binary_cast<UInt32, DateV2Value<DateV2ValueType>>(raw); + static_cast<void>(storage_date_days(value, "DATEV2")); + typed_output<ColumnDateV2>(output).insert_value(raw); + break; + } + case FieldType::OLAP_FIELD_TYPE_DATETIMEV2: { + const uint8_t scale = cursor.read<uint8_t>("DateTimeV2 scale"); + if (scale != signature.scale || scale > 6) { + throw Exception(ErrorCode::CORRUPTION, + "Variant binary DateTimeV2 scale changed while publishing"); + } + const UInt64 raw = cursor.read<UInt64>("DateTimeV2"); + const auto value = binary_cast<UInt64, DateV2Value<DateTimeV2ValueType>>(raw); + static_cast<void>(storage_timestamp_micros(value, "DATETIMEV2")); + typed_output<ColumnDateTimeV2>(output).insert_value(raw); + break; + } + case FieldType::OLAP_FIELD_TYPE_TIMESTAMPTZ: { + const uint8_t scale = cursor.read<uint8_t>("TimestampTz scale"); + if (scale != signature.scale || scale > 6) { + throw Exception(ErrorCode::CORRUPTION, + "Variant binary TimestampTz scale changed while publishing"); + } + const UInt64 raw = cursor.read<UInt64>("TimestampTz"); + const auto value = binary_cast<UInt64, TimestampTzValue>(raw); + static_cast<void>(storage_timestamp_micros(value, "TIMESTAMPTZ")); + typed_output<ColumnTimeStampTz>(output).insert_value(value); + break; + } + case FieldType::OLAP_FIELD_TYPE_DECIMAL: { + const auto value = read_legacy_decimal(cursor); + if (value.precision != signature.precision || value.scale != signature.scale) { + throw Exception(ErrorCode::CORRUPTION, + "Variant binary legacy DecimalV2 signature changed while publishing"); + } + typed_output<ColumnDecimal128V2>(output).insert_value(DecimalV2Value {value.value}); + break; + } + case FieldType::OLAP_FIELD_TYPE_DECIMAL32: { + const uint8_t precision = cursor.read<uint8_t>("Decimal32 precision"); + const uint8_t scale = cursor.read<uint8_t>("Decimal32 scale"); + const int32_t value = cursor.read<int32_t>("Decimal32"); + if (precision != signature.precision || scale != signature.scale) { + throw Exception(ErrorCode::CORRUPTION, + "Variant binary Decimal32 signature changed while publishing"); + } + validate_decimal(value, precision, scale, 9, "Decimal32"); + typed_output<ColumnDecimal32>(output).insert_value(Decimal32 {value}); + break; + } + case FieldType::OLAP_FIELD_TYPE_DECIMAL64: { + const uint8_t precision = cursor.read<uint8_t>("Decimal64 precision"); + const uint8_t scale = cursor.read<uint8_t>("Decimal64 scale"); + const int64_t value = cursor.read<int64_t>("Decimal64"); + if (precision != signature.precision || scale != signature.scale) { + throw Exception(ErrorCode::CORRUPTION, + "Variant binary Decimal64 signature changed while publishing"); + } + validate_decimal(value, precision, scale, 18, "Decimal64"); + typed_output<ColumnDecimal64>(output).insert_value(Decimal64 {value}); + break; + } + case FieldType::OLAP_FIELD_TYPE_DECIMAL128I: { + const uint8_t precision = cursor.read<uint8_t>("Decimal128 precision"); + const uint8_t scale = cursor.read<uint8_t>("Decimal128 scale"); + const __int128 value = cursor.read<__int128>("Decimal128"); + if (precision != signature.precision || scale != signature.scale) { + throw Exception(ErrorCode::CORRUPTION, + "Variant binary Decimal128 signature changed while publishing"); + } + validate_decimal(value, precision, scale, 38, "Decimal128"); + typed_output<ColumnDecimal128V3>(output).insert_value(Decimal128V3 {value}); + break; + } + case FieldType::OLAP_FIELD_TYPE_DECIMAL256: { + const uint8_t precision = cursor.read<uint8_t>("Decimal256 precision"); + const uint8_t scale = cursor.read<uint8_t>("Decimal256 scale"); + const wide::Int256 value = cursor.read<wide::Int256>("Decimal256"); + if (precision != signature.precision || scale != signature.scale) { + throw Exception(ErrorCode::CORRUPTION, + "Variant binary Decimal256 signature changed while publishing"); + } + validate_decimal(value, precision, scale, 76, "Decimal256"); + typed_output<ColumnDecimal256>(output).insert_value(Decimal256 {value}); + break; + } + default: + throw Exception(ErrorCode::CORRUPTION, "Unsupported typed Variant storage FieldType {}", + static_cast<uint8_t>(type)); + } + require_exact(cursor); +} + +void append_default_to_typed(const CellSignature& signature, IColumn* output) { Review Comment: 为什么不直接output->insert_default?? -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
