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) {