From 0f1904f40e6b8ed17b16a953dfc4ccb4c5e94261 Mon Sep 17 00:00:00 2001 From: Divjot Arora Date: Wed, 19 Aug 2026 22:01:59 +0000 Subject: [PATCH 1/2] Allow TIMESTAMP logical type to annotate FIXED_LEN_BYTE_ARRAY(12) --- cpp/src/parquet/arrow/arrow_schema_test.cc | 6 ++ cpp/src/parquet/arrow/schema_internal.cc | 2 + cpp/src/parquet/reader_test.cc | 62 ++++++++++++++++++++ cpp/src/parquet/schema_test.cc | 6 ++ cpp/src/parquet/statistics.cc | 54 +++++++++++++++++ cpp/src/parquet/statistics_test.cc | 67 ++++++++++++++++++++++ cpp/src/parquet/types.cc | 11 +++- 7 files changed, 206 insertions(+), 2 deletions(-) diff --git a/cpp/src/parquet/arrow/arrow_schema_test.cc b/cpp/src/parquet/arrow/arrow_schema_test.cc index 894f68900280..5d887f0e7429 100644 --- a/cpp/src/parquet/arrow/arrow_schema_test.cc +++ b/cpp/src/parquet/arrow/arrow_schema_test.cc @@ -263,6 +263,12 @@ TEST_F(TestConvertParquetSchema, ParquetAnnotatedFields) { ::arrow::fixed_size_binary(16)}, {"float16", LogicalType::Float16(), ParquetType::FIXED_LEN_BYTE_ARRAY, 2, ::arrow::float16()}, + {"timestamp_flba12_ms", LogicalType::Timestamp(true, LogicalType::TimeUnit::MILLIS), + ParquetType::FIXED_LEN_BYTE_ARRAY, 12, ::arrow::fixed_size_binary(12)}, + {"timestamp_flba12_us", LogicalType::Timestamp(true, LogicalType::TimeUnit::MICROS), + ParquetType::FIXED_LEN_BYTE_ARRAY, 12, ::arrow::fixed_size_binary(12)}, + {"timestamp_flba12_ns", LogicalType::Timestamp(true, LogicalType::TimeUnit::NANOS), + ParquetType::FIXED_LEN_BYTE_ARRAY, 12, ::arrow::fixed_size_binary(12)}, {"none", LogicalType::None(), ParquetType::BOOLEAN, -1, ::arrow::boolean()}, {"none", LogicalType::None(), ParquetType::INT32, -1, ::arrow::int32()}, {"none", LogicalType::None(), ParquetType::INT64, -1, ::arrow::int64()}, diff --git a/cpp/src/parquet/arrow/schema_internal.cc b/cpp/src/parquet/arrow/schema_internal.cc index 2e8cf764b27f..d4633ebbd6e6 100644 --- a/cpp/src/parquet/arrow/schema_internal.cc +++ b/cpp/src/parquet/arrow/schema_internal.cc @@ -207,6 +207,8 @@ Result> FromFLBA( return ::arrow::extension::uuid(); } + return ::arrow::fixed_size_binary(physical_length); + case LogicalType::Type::TIMESTAMP: return ::arrow::fixed_size_binary(physical_length); default: return Status::NotImplemented("Unhandled logical_type ", logical_type.ToString(), diff --git a/cpp/src/parquet/reader_test.cc b/cpp/src/parquet/reader_test.cc index cdeee116fbc3..f9a30255b595 100644 --- a/cpp/src/parquet/reader_test.cc +++ b/cpp/src/parquet/reader_test.cc @@ -138,6 +138,8 @@ std::string byte_stream_split_extended() { return data_file("byte_stream_split_extended.gzip.parquet"); } +std::string flba12_timestamp() { return data_file("flba12_timestamp.parquet"); } + template std::vector ReadColumnValues(ParquetFileReader* file_reader, int row_group, int column, int64_t expected_values_read) { @@ -1769,6 +1771,66 @@ TEST(TestByteStreamSplit, ExtendedIntegrationFile) { } #endif // ARROW_WITH_ZLIB +TEST(TestFileReader, TestFlba12Timestamp) { + auto file = ParquetFileReader::OpenFile(flba12_timestamp()); + + const int64_t kNumRows = 6; + // Row indices of the minimum (year 0001) and maximum (year 9999) values. + const int kMinRow = 5; + const int kMaxRow = 4; + + auto metadata = file->metadata(); + ASSERT_EQ(kNumRows, metadata->num_rows()); + ASSERT_EQ(3, metadata->num_columns()); + ASSERT_EQ(1, metadata->num_row_groups()); + + const struct { + const char* name; + LogicalType::TimeUnit::unit unit; + } columns[] = { + {"timestamp_millis", LogicalType::TimeUnit::MILLIS}, + {"timestamp_micros", LogicalType::TimeUnit::MICROS}, + {"timestamp_nanos", LogicalType::TimeUnit::NANOS}, + }; + + auto rg_reader = file->RowGroup(0); + for (int c = 0; c < 3; ++c) { + const auto* descr = metadata->schema()->Column(c); + ASSERT_EQ(columns[c].name, descr->name()); + ASSERT_EQ(Type::FIXED_LEN_BYTE_ARRAY, descr->physical_type()); + ASSERT_EQ(12, descr->type_length()); + ASSERT_EQ(SortOrder::SIGNED, descr->sort_order()); + ASSERT_EQ(ColumnOrder::TYPE_DEFINED_ORDER, descr->column_order().get_order()); + + const auto& logical_type = descr->logical_type(); + ASSERT_EQ(LogicalType::Type::TIMESTAMP, logical_type->type()); + const auto& ts = + ::arrow::internal::checked_cast(*logical_type); + ASSERT_TRUE(ts.is_adjusted_to_utc()); + ASSERT_EQ(columns[c].unit, ts.time_unit()); + + std::string min_value, max_value; + { + auto col_reader = + checked_pointer_cast>(rg_reader->Column(c)); + std::vector values(kNumRows); + int64_t values_read = 0; + int64_t levels_read = + col_reader->ReadBatch(kNumRows, nullptr, nullptr, values.data(), &values_read); + ASSERT_EQ(kNumRows, levels_read); + ASSERT_EQ(kNumRows, values_read); + min_value.assign(reinterpret_cast(values[kMinRow].ptr), 12); + max_value.assign(reinterpret_cast(values[kMaxRow].ptr), 12); + } + + auto stats = rg_reader->metadata()->ColumnChunk(c)->statistics(); + ASSERT_NE(nullptr, stats); + ASSERT_TRUE(stats->HasMinMax()); + ASSERT_EQ(min_value, stats->EncodeMin()); + ASSERT_EQ(max_value, stats->EncodeMax()); + } +} + struct PageIndexReaderParam { std::vector row_group_indices; std::vector column_indices; diff --git a/cpp/src/parquet/schema_test.cc b/cpp/src/parquet/schema_test.cc index 704b2da79c1b..f690a0bcafbb 100644 --- a/cpp/src/parquet/schema_test.cc +++ b/cpp/src/parquet/schema_test.cc @@ -1417,6 +1417,12 @@ TEST(TestLogicalTypeOperation, LogicalTypeApplicability) { for (const InapplicableType& t : inapplicable_types) { ASSERT_FALSE(logical_type->is_applicable(t.physical_type, t.physical_length)); } + + // TIMESTAMP is applicable to INT64 and FLBA(12). + logical_type = LogicalType::Timestamp(true, LogicalType::TimeUnit::MILLIS); + ASSERT_TRUE(logical_type->is_applicable(Type::INT64)); + ASSERT_TRUE(logical_type->is_applicable(Type::FIXED_LEN_BYTE_ARRAY, 12)); + ASSERT_FALSE(logical_type->is_applicable(Type::FIXED_LEN_BYTE_ARRAY, 8)); } TEST(TestLogicalTypeOperation, DecimalLogicalTypeApplicability) { diff --git a/cpp/src/parquet/statistics.cc b/cpp/src/parquet/statistics.cc index d43998ef78f2..ebbe7318ec9e 100644 --- a/cpp/src/parquet/statistics.cc +++ b/cpp/src/parquet/statistics.cc @@ -463,6 +463,56 @@ struct RebindLogical { using c_type = DType::c_type; }; +// Tag type for FLBA(12) timestamps. +struct Flba12TimestampType {}; + +// Max / min representable signed 96-bit two's-complement, little-endian. +constexpr uint8_t kFlba12SignedMax[12] = {0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, + 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0x7F}; +constexpr uint8_t kFlba12SignedMin[12] = {0x00, 0x00, 0x00, 0x00, 0x00, 0x00, + 0x00, 0x00, 0x00, 0x00, 0x00, 0x80}; + +template <> +struct CompareHelper { + using T = FLBA; + + // Seed for the running minimum is the maximum value; for the maximum, the minimum + // value. + static T DefaultMin() { return T{kFlba12SignedMax}; } + static T DefaultMax() { return T{kFlba12SignedMin}; } + + static T Coalesce(T val, T fallback) { return val.ptr == nullptr ? fallback : val; } + + // Signed little-endian comparison. + // Differing signs: negative (MSB >= 0x80) is smaller. + // Same sign: unsigned scan and comparison. + static inline bool Compare(int /*type_length*/, const T& a, const T& b) { + const bool a_neg = (a.ptr[11] & 0x80) != 0; + const bool b_neg = (b.ptr[11] & 0x80) != 0; + if (a_neg != b_neg) return a_neg; + for (int i = 11; i >= 0; --i) if (a.ptr[i] != b.ptr[i]) return a.ptr[i] < b.ptr[i]; + return false; + } + + static T Min(int type_length, const T& a, const T& b) { + if (a.ptr == nullptr) return b; + if (b.ptr == nullptr) return a; + return Compare(type_length, a, b) ? a : b; + } + + static T Max(int type_length, const T& a, const T& b) { + if (a.ptr == nullptr) return b; + if (b.ptr == nullptr) return a; + return Compare(type_length, a, b) ? b : a; + } +}; + +template <> +struct RebindLogical { + using DType = FLBAType; + using c_type = DType::c_type; +}; + template class TypedComparatorImpl : virtual public TypedComparator::DType> { @@ -1014,6 +1064,10 @@ std::shared_ptr DoMakeComparator(Type::type physical_type, return std::make_shared>( type_length); } + if (logical_type == LogicalType::Type::TIMESTAMP) { + return std::make_shared>( + type_length); + } return std::make_shared>(type_length); default: ParquetException::NYI("Signed Compare not implemented"); diff --git a/cpp/src/parquet/statistics_test.cc b/cpp/src/parquet/statistics_test.cc index bf0961c4fc6a..af9d96a0ade8 100644 --- a/cpp/src/parquet/statistics_test.cc +++ b/cpp/src/parquet/statistics_test.cc @@ -175,6 +175,73 @@ TEST(Comparison, SignedFLBA) { } } +TEST(Comparison, SignedLittleEndianFLBA12Timestamp) { + NodePtr node = + PrimitiveNode::Make("ts_flba12", Repetition::REQUIRED, + LogicalType::Timestamp(true, LogicalType::TimeUnit::NANOS), + Type::FIXED_LEN_BYTE_ARRAY, 12); + ColumnDescriptor descr(node, 0, 0); + ASSERT_EQ(SortOrder::SIGNED, descr.sort_order()); + auto comparator = MakeComparator(&descr); + + // large_neg: 0x80 00…00 (most negative 96-bit LE value, sign byte=0x80) + // neg256: 0x00 FF FF…FF (= -256 in LE two's complement) + // minus_one: FF FF…FF (= -1) + // zero: 00 00…00 + // plus_one: 01 00…00 + // large_pos: 7F FF…FF (= near max positive) + std::vector large_neg_bytes = {0x00, 0x00, 0x00, 0x00, 0x00, 0x00, + 0x00, 0x00, 0x00, 0x00, 0x00, 0x80}; + std::vector neg256_bytes = {0x00, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, + 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF}; + std::vector minus_one_bytes = {0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, + 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF}; + std::vector zero_bytes = {0x00, 0x00, 0x00, 0x00, 0x00, 0x00, + 0x00, 0x00, 0x00, 0x00, 0x00, 0x00}; + std::vector plus_one_bytes = {0x01, 0x00, 0x00, 0x00, 0x00, 0x00, + 0x00, 0x00, 0x00, 0x00, 0x00, 0x00}; + std::vector large_pos_bytes = {0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, + 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0x7F}; + + std::vector vals = {FLBA(large_neg_bytes.data()), FLBA(neg256_bytes.data()), + FLBA(minus_one_bytes.data()), FLBA(zero_bytes.data()), + FLBA(plus_one_bytes.data()), FLBA(large_pos_bytes.data())}; + + for (size_t x = 0; x < vals.size(); x++) { + EXPECT_FALSE(comparator->Compare(vals[x], vals[x])) << x; + for (size_t y = x + 1; y < vals.size(); y++) { + EXPECT_TRUE(comparator->Compare(vals[x], vals[y])) << x << " < " << y; + EXPECT_FALSE(comparator->Compare(vals[y], vals[x])) << y << " < " << x; + } + } +} + +TEST(Comparison, SignedLittleEndianFLBA12TimestampMinMax) { + // Guards the DefaultMin/DefaultMax accumulator seeds: over positive-only values. + NodePtr node = + PrimitiveNode::Make("ts_flba12", Repetition::REQUIRED, + LogicalType::Timestamp(true, LogicalType::TimeUnit::NANOS), + Type::FIXED_LEN_BYTE_ARRAY, 12); + ColumnDescriptor descr(node, 0, 0); + auto comparator = MakeComparator(&descr); + + std::vector plus_one = {0x01, 0x00, 0x00, 0x00, 0x00, 0x00, + 0x00, 0x00, 0x00, 0x00, 0x00, 0x00}; + std::vector plus_five = {0x05, 0x00, 0x00, 0x00, 0x00, 0x00, + 0x00, 0x00, 0x00, 0x00, 0x00, 0x00}; + std::vector large_pos = {0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, + 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0x7F}; + std::vector vals = {FLBA(plus_five.data()), FLBA(large_pos.data()), + FLBA(plus_one.data())}; + + auto min_max = comparator->GetMinMax(vals.data(), vals.size()); + // min == plus_one and max == large_pos. + EXPECT_FALSE(comparator->Compare(min_max.first, FLBA(plus_one.data()))); + EXPECT_FALSE(comparator->Compare(FLBA(plus_one.data()), min_max.first)); + EXPECT_FALSE(comparator->Compare(min_max.second, FLBA(large_pos.data()))); + EXPECT_FALSE(comparator->Compare(FLBA(large_pos.data()), min_max.second)); +} + TEST(Comparison, UnsignedFLBA) { int size = 10; auto comparator = diff --git a/cpp/src/parquet/types.cc b/cpp/src/parquet/types.cc index cc3199f367af..d901ec87d30d 100644 --- a/cpp/src/parquet/types.cc +++ b/cpp/src/parquet/types.cc @@ -1376,10 +1376,12 @@ LogicalType::TimeUnit::unit TimeLogicalType::time_unit() const { } class LogicalType::Impl::Timestamp final : public LogicalType::Impl::Compatible, - public LogicalType::Impl::SimpleApplicable { + public LogicalType::Impl::Applicable { public: friend class TimestampLogicalType; + bool is_applicable(parquet::Type::type primitive_type, + int32_t primitive_length = -1) const override; bool is_serialized() const override; bool is_compatible(ConvertedType::type converted_type, schema::DecimalMetadata converted_decimal_metadata) const override; @@ -1400,7 +1402,6 @@ class LogicalType::Impl::Timestamp final : public LogicalType::Impl::Compatible, Timestamp(bool adjusted, LogicalType::TimeUnit::unit unit, bool is_from_converted_type, bool force_set_converted_type) : LogicalType::Impl(LogicalType::Type::TIMESTAMP, SortOrder::SIGNED), - LogicalType::Impl::SimpleApplicable(parquet::Type::INT64), adjusted_(adjusted), unit_(unit), is_from_converted_type_(is_from_converted_type), @@ -1411,6 +1412,12 @@ class LogicalType::Impl::Timestamp final : public LogicalType::Impl::Compatible, bool force_set_converted_type_ = false; }; +bool LogicalType::Impl::Timestamp::is_applicable(parquet::Type::type primitive_type, + int32_t primitive_length) const { + return primitive_type == parquet::Type::INT64 || + (primitive_type == parquet::Type::FIXED_LEN_BYTE_ARRAY && primitive_length == 12); +} + bool LogicalType::Impl::Timestamp::is_serialized() const { return !is_from_converted_type_; } From 97b3d4edf4a74a919834992168e879d8e6ea786f Mon Sep 17 00:00:00 2001 From: Divjot Arora Date: Thu, 20 Aug 2026 21:57:37 +0000 Subject: [PATCH 2/2] address comments --- .../parquet/arrow/arrow_reader_writer_test.cc | 66 +++++++++++++++++++ cpp/src/parquet/arrow/arrow_schema_test.cc | 23 +++++++ cpp/src/parquet/arrow/reader_internal.cc | 48 ++++++++++++++ cpp/src/parquet/arrow/schema_internal.cc | 5 ++ cpp/src/parquet/properties.h | 29 +++++++- cpp/src/parquet/reader_test.cc | 7 ++ cpp/src/parquet/statistics.cc | 21 +++--- 7 files changed, 190 insertions(+), 9 deletions(-) diff --git a/cpp/src/parquet/arrow/arrow_reader_writer_test.cc b/cpp/src/parquet/arrow/arrow_reader_writer_test.cc index 2bdbc38b3647..983210f75a6a 100644 --- a/cpp/src/parquet/arrow/arrow_reader_writer_test.cc +++ b/cpp/src/parquet/arrow/arrow_reader_writer_test.cc @@ -2166,6 +2166,72 @@ TEST(TestArrowReadWrite, CoerceTimestampsLosePrecision) { allow_truncation_to_micros)); } +TEST(TestArrowReadWrite, FlbaTimestampConversionValues) { + auto node = + PrimitiveNode::Make("ts", Repetition::REQUIRED, + LogicalType::Timestamp(true, LogicalType::TimeUnit::MICROS), + ParquetType::FIXED_LEN_BYTE_ARRAY, /*length=*/12); + auto file_schema = std::static_pointer_cast( + GroupNode::Make("schema", Repetition::REQUIRED, {node})); + + // Little-endian 96-bit values: 1,000,000 (fits int64) and 2^64 (overflows int64). + uint8_t in_range[12] = {0x40, 0x42, 0x0f, 0, 0, 0, 0, 0, 0, 0, 0, 0}; + uint8_t overflow[12] = {0, 0, 0, 0, 0, 0, 0, 0, 1, 0, 0, 0}; + FLBA values[2] = {FLBA(in_range), FLBA(overflow)}; + + auto sink = CreateOutputStream(); + auto writer = ParquetFileWriter::Open(sink, file_schema); + RowGroupWriter* rg_writer = writer->AppendRowGroup(); + auto* col_writer = dynamic_cast*>(rg_writer->NextColumn()); + ASSERT_NE(col_writer, nullptr); + col_writer->WriteBatch(2, nullptr, nullptr, values); + col_writer->Close(); + rg_writer->Close(); + writer->Close(); + ASSERT_OK_AND_ASSIGN(auto buffer, sink->Finish()); + + auto read_table = [&buffer](ArrowReaderProperties props, + std::shared_ptr* out) -> ::arrow::Status { + FileReaderBuilder builder; + RETURN_NOT_OK(builder.Open(std::make_shared(buffer))); + std::unique_ptr reader; + RETURN_NOT_OK(builder.properties(props)->Build(&reader)); + return reader->ReadTable(out); + }; + + // Default: raw, lossless FixedSizeBinary(12). + { + std::shared_ptr
table; + ASSERT_OK(read_table(ArrowReaderProperties(), &table)); + ASSERT_EQ(::arrow::Type::FIXED_SIZE_BINARY, table->schema()->field(0)->type()->id()); + } + + // Convert, error on overflow (default policy): the 2^64 row fails the read. + { + ArrowReaderProperties props; + props.set_convert_flba_timestamps(true); + std::shared_ptr
table; + ASSERT_RAISES(Invalid, read_table(props, &table)); + } + + // Convert, clamp on overflow: in-range value is exact; overflow clamps to + // INT64_MAX. + { + ArrowReaderProperties props; + props.set_convert_flba_timestamps(true); + props.set_flba_timestamp_clamp_on_overflow(true); + std::shared_ptr
table; + ASSERT_OK(read_table(props, &table)); + ASSERT_EQ(*::arrow::timestamp(TimeUnit::MICRO, "UTC"), + *table->schema()->field(0)->type()); + auto ts = + std::static_pointer_cast<::arrow::TimestampArray>(table->column(0)->chunk(0)); + ASSERT_EQ(2, ts->length()); + ASSERT_EQ(1000000, ts->Value(0)); + ASSERT_EQ(INT64_MAX, ts->Value(1)); + } +} + TEST(TestArrowReadWrite, ImplicitSecondToMillisecondTimestampCoercion) { using ::arrow::ArrayFromVector; using ::arrow::field; diff --git a/cpp/src/parquet/arrow/arrow_schema_test.cc b/cpp/src/parquet/arrow/arrow_schema_test.cc index 5d887f0e7429..07e8811c7368 100644 --- a/cpp/src/parquet/arrow/arrow_schema_test.cc +++ b/cpp/src/parquet/arrow/arrow_schema_test.cc @@ -312,6 +312,29 @@ TEST_F(TestConvertParquetSchema, DuplicateFieldNames) { ASSERT_NO_FATAL_FAILURE(CheckFlatSchema(::arrow::schema(arrow_fields))); } +TEST_F(TestConvertParquetSchema, FlbaTimestampConversion) { + auto make_fields = [] { + std::vector fields; + fields.push_back( + PrimitiveNode::Make("ts", Repetition::REQUIRED, + LogicalType::Timestamp(true, LogicalType::TimeUnit::MICROS), + ParquetType::FIXED_LEN_BYTE_ARRAY, /*length=*/12)); + return fields; + }; + + // Should output the raw FLBA value. + ASSERT_OK(ConvertSchema(make_fields())); + ASSERT_NO_FATAL_FAILURE(CheckFlatSchema( + ::arrow::schema({::arrow::field("ts", ::arrow::fixed_size_binary(12), false)}))); + + // Should convert to an Arrow timestamp. + ArrowReaderProperties props; + props.set_convert_flba_timestamps(true); + ASSERT_OK(ConvertSchema(make_fields(), /*key_value_metadata=*/{}, props)); + ASSERT_NO_FATAL_FAILURE(CheckFlatSchema(::arrow::schema({::arrow::field( + "ts", ::arrow::timestamp(::arrow::TimeUnit::MICRO, "UTC"), false)}))); +} + TEST_F(TestConvertParquetSchema, ParquetKeyValueMetadata) { std::vector parquet_fields; std::vector> arrow_fields; diff --git a/cpp/src/parquet/arrow/reader_internal.cc b/cpp/src/parquet/arrow/reader_internal.cc index 12f36fe39cf8..7563fae82174 100644 --- a/cpp/src/parquet/arrow/reader_internal.cc +++ b/cpp/src/parquet/arrow/reader_internal.cc @@ -28,6 +28,7 @@ #include #include "arrow/array.h" +#include "arrow/builder.h" #include "arrow/compute/api.h" #include "arrow/datum.h" #include "arrow/io/memory.h" @@ -855,6 +856,49 @@ Status TransferHalfFloat(RecordReader* reader, MemoryPool* pool, return Status::OK(); } +// Read a TIMESTAMP-annotated FLBA(12) column as a 64-bit Arrow timestamp. Values that do +// not fit in the 64 bit range either error or clamp to min/max int64, depending on +// configuration. +Status TransferFlbaTimestamp(RecordReader* reader, MemoryPool* pool, + const std::shared_ptr& field, Datum* out, + bool clamp_on_overflow) { + static const auto binary_type = ::arrow::fixed_size_binary(12); + std::shared_ptr chunked_array; + RETURN_NOT_OK( + TransferBinary(reader, pool, field->WithType(binary_type), &chunked_array)); + + ::arrow::TimestampBuilder builder(field->type(), pool); + RETURN_NOT_OK(builder.Reserve(chunked_array->length())); + for (const auto& chunk : chunked_array->chunks()) { + const auto& values = checked_cast(*chunk); + for (int64_t i = 0; i < values.length(); ++i) { + if (values.IsNull(i)) { + builder.UnsafeAppendNull(); + continue; + } + const uint8_t* bytes = values.GetValue(i); + const uint64_t low = bit_util::FromLittleEndian(SafeLoadAs(bytes)); + const uint32_t high = bit_util::FromLittleEndian(SafeLoadAs(bytes + 8)); + const int64_t low_signed = static_cast(low); + // Fits in int64 iff the high part is a pure sign-extension of the low part. + if (static_cast(high) != (low_signed < 0 ? -1 : 0)) { + if (!clamp_on_overflow) { + return Status::Invalid( + "FLBA(12) TIMESTAMP value does not fit in a 64-bit Arrow timestamp"); + } + const bool negative = (bytes[11] & 0x80) != 0; + builder.UnsafeAppend(negative ? INT64_MIN : INT64_MAX); + } else { + builder.UnsafeAppend(low_signed); + } + } + } + std::shared_ptr<::arrow::Array> array; + RETURN_NOT_OK(builder.Finish(&array)); + *out = array; + return Status::OK(); +} + } // namespace #define TRANSFER_INT32(ENUM, ArrowType) \ @@ -966,6 +1010,10 @@ Status TransferColumnData(RecordReader* reader, if (descr->physical_type() == ::parquet::Type::INT96) { RETURN_NOT_OK( TransferInt96(reader, pool, value_field, &result, timestamp_type.unit())); + } else if (descr->physical_type() == ::parquet::Type::FIXED_LEN_BYTE_ARRAY) { + RETURN_NOT_OK(TransferFlbaTimestamp( + reader, pool, value_field, &result, + ctx->reader_properties->flba_timestamp_clamp_on_overflow())); } else { switch (timestamp_type.unit()) { case ::arrow::TimeUnit::MILLI: diff --git a/cpp/src/parquet/arrow/schema_internal.cc b/cpp/src/parquet/arrow/schema_internal.cc index d4633ebbd6e6..68f80089a891 100644 --- a/cpp/src/parquet/arrow/schema_internal.cc +++ b/cpp/src/parquet/arrow/schema_internal.cc @@ -209,6 +209,11 @@ Result> FromFLBA( return ::arrow::fixed_size_binary(physical_length); case LogicalType::Type::TIMESTAMP: + // If configured, convert to a potentially lossy Arrow timestamp. Otherwise, return + // the raw lossless FLBA value. + if (physical_length == 12 && reader_properties.convert_flba_timestamps()) { + return MakeArrowTimestamp(logical_type); + } return ::arrow::fixed_size_binary(physical_length); default: return Status::NotImplemented("Unhandled logical_type ", logical_type.ToString(), diff --git a/cpp/src/parquet/properties.h b/cpp/src/parquet/properties.h index e2244a1176e3..a9984e7d897f 100644 --- a/cpp/src/parquet/properties.h +++ b/cpp/src/parquet/properties.h @@ -1157,7 +1157,9 @@ class PARQUET_EXPORT ArrowReaderProperties { list_type_(kArrowDefaultListType), arrow_extensions_enabled_(false), should_load_statistics_(false), - smallest_decimal_enabled_(false) {} + smallest_decimal_enabled_(false), + convert_flba_timestamps_(false), + flba_timestamp_clamp_on_overflow_(false) {} /// \brief Set whether to use the IO thread pool to parse columns in parallel. /// @@ -1295,6 +1297,29 @@ class PARQUET_EXPORT ArrowReaderProperties { /// this setting will be ignored. bool smallest_decimal_enabled() const { return smallest_decimal_enabled_; } + /// \brief Set whether to infer Arrow timestamps from Parquet FLBA types. + /// + /// When enabled, Parquet FLBA(12) TIMESTAMP columns are read as Arrow timestamps. + /// VAlues that do not fit in 64 bit timestamps are handled per + /// flba_timestamp_clamp_on_overflow(). When disabled, Parquet FLBA(12) TIMESTAMP + /// columns are read as FixedSizeBinary(12). + void set_convert_flba_timestamps(bool convert) { convert_flba_timestamps_ = convert; } + /// \brief Whether FLBA(12) TIMESTAMP columns are read as Arrow timestamps. + bool convert_flba_timestamps() const { return convert_flba_timestamps_; } + + /// \brief Set how out-of-range values are handled when convert_flba_timestamps() is + /// enabled. + /// + /// When true, Parquet FLBA(12) TIMESTAMP values that do not fit in 64 bit timestamps + /// are clamped to min/max INT64. When false, such values raise an error. + void set_flba_timestamp_clamp_on_overflow(bool clamp) { + flba_timestamp_clamp_on_overflow_ = clamp; + } + /// \brief Whether out-of-range FLBA(12) timestamps clamp (true) or error (false). + bool flba_timestamp_clamp_on_overflow() const { + return flba_timestamp_clamp_on_overflow_; + } + private: bool use_threads_; std::unordered_set read_dict_indices_; @@ -1308,6 +1333,8 @@ class PARQUET_EXPORT ArrowReaderProperties { bool arrow_extensions_enabled_; bool should_load_statistics_; bool smallest_decimal_enabled_; + bool convert_flba_timestamps_; + bool flba_timestamp_clamp_on_overflow_; }; /// EXPERIMENTAL: Constructs the default ArrowReaderProperties diff --git a/cpp/src/parquet/reader_test.cc b/cpp/src/parquet/reader_test.cc index f9a30255b595..0fc366c0745e 100644 --- a/cpp/src/parquet/reader_test.cc +++ b/cpp/src/parquet/reader_test.cc @@ -1821,6 +1821,13 @@ TEST(TestFileReader, TestFlba12Timestamp) { ASSERT_EQ(kNumRows, values_read); min_value.assign(reinterpret_cast(values[kMinRow].ptr), 12); max_value.assign(reinterpret_cast(values[kMaxRow].ptr), 12); + + auto comparator = MakeComparator(descr); + auto min_max = comparator->GetMinMax(values.data(), kNumRows); + ASSERT_EQ(min_value, + std::string(reinterpret_cast(min_max.first.ptr), 12)); + ASSERT_EQ(max_value, + std::string(reinterpret_cast(min_max.second.ptr), 12)); } auto stats = rg_reader->metadata()->ColumnChunk(c)->statistics(); diff --git a/cpp/src/parquet/statistics.cc b/cpp/src/parquet/statistics.cc index ebbe7318ec9e..71222db79f1f 100644 --- a/cpp/src/parquet/statistics.cc +++ b/cpp/src/parquet/statistics.cc @@ -29,6 +29,7 @@ #include "arrow/type_traits.h" #include "arrow/util/bit_run_reader.h" #include "arrow/util/checked_cast.h" +#include "arrow/util/endian.h" #include "arrow/util/float16.h" #include "arrow/util/logging_internal.h" #include "arrow/util/ubsan.h" @@ -41,10 +42,12 @@ using arrow::default_memory_pool; using arrow::MemoryPool; +using arrow::bit_util::FromLittleEndian; using arrow::internal::checked_cast; using arrow::util::Float16; using arrow::util::SafeCopy; using arrow::util::SafeLoad; +using arrow::util::SafeLoadAs; namespace parquet { namespace { @@ -483,15 +486,17 @@ struct CompareHelper { static T Coalesce(T val, T fallback) { return val.ptr == nullptr ? fallback : val; } - // Signed little-endian comparison. - // Differing signs: negative (MSB >= 0x80) is smaller. - // Same sign: unsigned scan and comparison. + // FLBA(12) TIMESTAMP is a signed 96-bit little-endian value. Compare the + // most-significant 32 bits signed, then the low 64 bits unsigned. static inline bool Compare(int /*type_length*/, const T& a, const T& b) { - const bool a_neg = (a.ptr[11] & 0x80) != 0; - const bool b_neg = (b.ptr[11] & 0x80) != 0; - if (a_neg != b_neg) return a_neg; - for (int i = 11; i >= 0; --i) if (a.ptr[i] != b.ptr[i]) return a.ptr[i] < b.ptr[i]; - return false; + const int32_t a_hi = + SafeCopy(FromLittleEndian(SafeLoadAs(a.ptr + 8))); + const int32_t b_hi = + SafeCopy(FromLittleEndian(SafeLoadAs(b.ptr + 8))); + if (a_hi != b_hi) return a_hi < b_hi; + const uint64_t a_lo = FromLittleEndian(SafeLoadAs(a.ptr)); + const uint64_t b_lo = FromLittleEndian(SafeLoadAs(b.ptr)); + return a_lo < b_lo; } static T Min(int type_length, const T& a, const T& b) {