This is an automated email from the ASF dual-hosted git repository.
pitrou pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow.git
The following commit(s) were added to refs/heads/main by this push:
new 464ae94d81 GH-50915: [FORMAT] Allow TIMESTAMP logical type to annotate
FIXED_LEN_BYTE_ARRAY(12) (#50916)
464ae94d81 is described below
commit 464ae94d8180ad6cf2c68359985541f0e1240782
Author: Divjot Arora <[email protected]>
AuthorDate: Mon Sep 28 17:21:35 2026 +0200
GH-50915: [FORMAT] Allow TIMESTAMP logical type to annotate
FIXED_LEN_BYTE_ARRAY(12) (#50916)
### Rationale for this change
See https://github.com/apache/parquet-format/issues/600 for rationale.
### What changes are included in this PR?
This PR adds support for using `TimestampType` to annotate
`FIXED_LEN_BYTE_ARRAY(12)` values. It also adds functionality to convert
FLBA(12) values to Arrow INT64 timestamps with the following flags/logic to
handle overflow:
1. The conversion is guarded behind the `convert_flba_timestamps` property
(default true). If false, conversion fails regardless of value.
2. If the FLBA(12) value overflows max int64 or underflows min int64, the
`flba_timestamp_clamp_on_overflow` property (default false) is consulted. If
false, conversion fails. If true, the value is clamped to min/max int64.
### Are these changes tested?
Yes, via unit tests and an e2e test that reads the file added in
parquet-testing (https://github.com/apache/parquet-testing/pull/123).
### Are there any user-facing changes?
No
* GitHub Issue: #50915
Authored-by: Divjot Arora <[email protected]>
Signed-off-by: Antoine Pitrou <[email protected]>
---
cpp/src/parquet/arrow/arrow_reader_writer_test.cc | 117 ++++++++++++++++++++++
cpp/src/parquet/arrow/arrow_schema_test.cc | 32 ++++++
cpp/src/parquet/arrow/reader_internal.cc | 77 +++++++++++++-
cpp/src/parquet/arrow/schema_internal.cc | 37 ++++---
cpp/src/parquet/arrow/schema_internal.h | 4 +
cpp/src/parquet/properties.h | 29 +++++-
cpp/src/parquet/reader_test.cc | 69 +++++++++++++
cpp/src/parquet/schema_test.cc | 6 ++
cpp/src/parquet/statistics.cc | 59 +++++++++++
cpp/src/parquet/statistics_test.cc | 67 +++++++++++++
cpp/src/parquet/types.cc | 12 ++-
cpp/submodules/parquet-testing | 2 +-
12 files changed, 492 insertions(+), 19 deletions(-)
diff --git a/cpp/src/parquet/arrow/arrow_reader_writer_test.cc
b/cpp/src/parquet/arrow/arrow_reader_writer_test.cc
index 2bdbc38b36..d858962d47 100644
--- a/cpp/src/parquet/arrow/arrow_reader_writer_test.cc
+++ b/cpp/src/parquet/arrow/arrow_reader_writer_test.cc
@@ -2166,6 +2166,123 @@ 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>(
+ GroupNode::Make("schema", Repetition::REQUIRED, {node}));
+
+ // Little-endian 96-bit values: 1,000,000 and -1,000,000 (both fit int64),
+ // 2^64 (overflows INT64_MAX), and -2^64 (underflows INT64_MIN).
+ uint8_t pos_in_range[12] = {0x40, 0x42, 0x0f, 0, 0, 0, 0, 0, 0, 0, 0, 0};
+ uint8_t neg_in_range[12] = {0xc0, 0xbd, 0xf0, 0xff, 0xff, 0xff,
+ 0xff, 0xff, 0xff, 0xff, 0xff, 0xff};
+ uint8_t overflow[12] = {0, 0, 0, 0, 0, 0, 0, 0, 1, 0, 0, 0};
+ uint8_t neg_overflow[12] = {0, 0, 0, 0, 0, 0, 0, 0, 0xff, 0xff, 0xff, 0xff};
+ FLBA values[4] = {FLBA(pos_in_range), FLBA(neg_in_range), FLBA(overflow),
+ FLBA(neg_overflow)};
+
+ auto sink = CreateOutputStream();
+ auto writer = ParquetFileWriter::Open(sink, file_schema);
+ RowGroupWriter* rg_writer = writer->AppendRowGroup();
+ auto* col_writer =
dynamic_cast<TypedColumnWriter<FLBAType>*>(rg_writer->NextColumn());
+ ASSERT_NE(col_writer, nullptr);
+ col_writer->WriteBatch(4, 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) -> Result<std::shared_ptr<Table>>
{
+ FileReaderBuilder builder;
+ RETURN_NOT_OK(builder.Open(std::make_shared<BufferReader>(buffer)));
+ std::unique_ptr<FileReader> reader;
+ RETURN_NOT_OK(builder.properties(props)->Build(&reader));
+ return reader->ReadTable();
+ };
+
+ // Convert, error on overflow (default): the out-of-range rows fail the read.
+ {
+ ArrowReaderProperties props;
+ EXPECT_RAISES_WITH_MESSAGE_THAT(
+ Invalid, ::testing::HasSubstr("does not fit in a 64-bit Arrow
timestamp"),
+ read_table(props));
+ }
+
+ // Conversion disabled: raw, lossless FixedSizeBinary(12).
+ {
+ ArrowReaderProperties props;
+ props.set_convert_flba_timestamps(false);
+ ASSERT_OK_AND_ASSIGN(auto table, read_table(props));
+ ASSERT_OK(table->ValidateFull());
+ ASSERT_EQ(::arrow::Type::FIXED_SIZE_BINARY,
table->schema()->field(0)->type()->id());
+ const auto& raw =
+ checked_cast<const
::arrow::FixedSizeBinaryArray&>(*table->column(0)->chunk(0));
+ for (int64_t i = 0; i < raw.length(); ++i) {
+ ASSERT_EQ(std::string_view(reinterpret_cast<const char*>(values[i].ptr),
12),
+ raw.GetView(i));
+ }
+ }
+
+ // Convert, clamp on overflow: in-range value is exact; positive overflow
clamps
+ // to INT64_MAX and negative overflow clamps to INT64_MIN.
+ {
+ ArrowReaderProperties props;
+ props.set_flba_timestamp_clamp_on_overflow(true);
+ ASSERT_OK_AND_ASSIGN(auto table, read_table(props));
+ ASSERT_OK(table->ValidateFull());
+ 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(4, ts->length());
+ ASSERT_EQ(1000000, ts->Value(0));
+ ASSERT_EQ(-1000000, ts->Value(1));
+ ASSERT_EQ(INT64_MAX, ts->Value(2));
+ ASSERT_EQ(INT64_MIN, ts->Value(3));
+ }
+}
+
+TEST(TestArrowReadWrite, FlbaTimestampIntegration) {
+ ArrowReaderProperties props;
+ props.set_flba_timestamp_clamp_on_overflow(true);
+ ASSERT_OK_AND_ASSIGN(
+ auto reader,
+ FileReader::Make(::arrow::default_memory_pool(),
+ ParquetFileReader::OpenFile(
+ test::get_data_file("flba12_timestamp.parquet"),
false),
+ props));
+ ASSERT_OK_AND_ASSIGN(auto actual, reader->ReadTable());
+ ASSERT_OK(actual->ValidateFull());
+
+ auto expected_schema = ::arrow::schema({
+ ::arrow::field("timestamp_millis", ::arrow::timestamp(TimeUnit::MILLI,
"UTC")),
+ ::arrow::field("timestamp_micros", ::arrow::timestamp(TimeUnit::MICRO,
"UTC")),
+ ::arrow::field("timestamp_nanos", ::arrow::timestamp(TimeUnit::NANO,
"UTC")),
+ });
+ std::shared_ptr<Array> expected_millis;
+ ::arrow::ArrayFromVector<::arrow::TimestampType, int64_t>(
+ expected_schema->field(0)->type(),
+ {0, 1000, -1000, 9223372036000, 253402300799000, -62135596800000},
+ &expected_millis);
+ std::shared_ptr<Array> expected_micros;
+ ::arrow::ArrayFromVector<::arrow::TimestampType, int64_t>(
+ expected_schema->field(1)->type(),
+ {0, 1000000, -1000000, 9223372036000000, 253402300799000000,
-62135596800000000},
+ &expected_micros);
+ std::shared_ptr<Array> expected_nanos;
+ ::arrow::ArrayFromVector<::arrow::TimestampType, int64_t>(
+ expected_schema->field(2)->type(),
+ {0, 1000000000, -1000000000, 9223372036000000000, INT64_MAX, INT64_MIN},
+ &expected_nanos);
+ auto expected =
+ Table::Make(expected_schema, {expected_millis, expected_micros,
expected_nanos});
+ ASSERT_NO_FATAL_FAILURE(::arrow::AssertTablesEqual(*expected, *actual));
+}
+
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 27c302fe0d..6fbde99246 100644
--- a/cpp/src/parquet/arrow/arrow_schema_test.cc
+++ b/cpp/src/parquet/arrow/arrow_schema_test.cc
@@ -263,6 +263,15 @@ 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::timestamp(::arrow::TimeUnit::MILLI, "UTC")},
+ {"timestamp_flba12_us", LogicalType::Timestamp(true,
LogicalType::TimeUnit::MICROS),
+ ParquetType::FIXED_LEN_BYTE_ARRAY, 12,
+ ::arrow::timestamp(::arrow::TimeUnit::MICRO, "UTC")},
+ {"timestamp_flba12_ns", LogicalType::Timestamp(true,
LogicalType::TimeUnit::NANOS),
+ ParquetType::FIXED_LEN_BYTE_ARRAY, 12,
+ ::arrow::timestamp(::arrow::TimeUnit::NANO, "UTC")},
{"none", LogicalType::None(), ParquetType::BOOLEAN, -1,
::arrow::boolean()},
{"none", LogicalType::None(), ParquetType::INT32, -1, ::arrow::int32()},
{"none", LogicalType::None(), ParquetType::INT64, -1, ::arrow::int64()},
@@ -306,6 +315,29 @@ TEST_F(TestConvertParquetSchema, DuplicateFieldNames) {
ASSERT_NO_FATAL_FAILURE(CheckFlatSchema(::arrow::schema(arrow_fields)));
}
+TEST_F(TestConvertParquetSchema, FlbaTimestampConversion) {
+ auto make_fields = [] {
+ std::vector<NodePtr> 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 convert to an Arrow timestamp.
+ ASSERT_OK(ConvertSchema(make_fields()));
+ ASSERT_NO_FATAL_FAILURE(CheckFlatSchema(::arrow::schema({::arrow::field(
+ "ts", ::arrow::timestamp(::arrow::TimeUnit::MICRO, "UTC"), false)})));
+
+ // Should output the raw FLBA value.
+ ArrowReaderProperties props;
+ props.set_convert_flba_timestamps(false);
+ ASSERT_OK(ConvertSchema(make_fields(), /*key_value_metadata=*/{}, props));
+ ASSERT_NO_FATAL_FAILURE(CheckFlatSchema(
+ ::arrow::schema({::arrow::field("ts", ::arrow::fixed_size_binary(12),
false)})));
+}
+
TEST_F(TestConvertParquetSchema, ParquetKeyValueMetadata) {
std::vector<NodePtr> parquet_fields;
std::vector<std::shared_ptr<Field>> arrow_fields;
diff --git a/cpp/src/parquet/arrow/reader_internal.cc
b/cpp/src/parquet/arrow/reader_internal.cc
index b5207aca25..0131beeada 100644
--- a/cpp/src/parquet/arrow/reader_internal.cc
+++ b/cpp/src/parquet/arrow/reader_internal.cc
@@ -18,9 +18,9 @@
#include "parquet/arrow/reader_internal.h"
#include <algorithm>
-#include <climits>
#include <cstdint>
#include <cstring>
+#include <limits>
#include <memory>
#include <string>
#include <string_view>
@@ -46,6 +46,7 @@
#include "arrow/util/int_util_overflow.h"
#include "arrow/util/logging_internal.h"
#include "arrow/util/ubsan.h"
+#include "arrow/visit_data_inline.h"
#include "parquet/arrow/reader.h"
#include "parquet/arrow/schema.h"
@@ -853,6 +854,68 @@ Status TransferHalfFloat(RecordReader* reader, MemoryPool*
pool,
return Status::OK();
}
+// Decode a little-endian 96-bit FLBA(12) TIMESTAMP value into a 64-bit Arrow
timestamp.
+// Values that do not fit in the int64 range either error or clamp to the
minimum or
+// maximum int64 value, depending on clamp_on_overflow.
+inline Result<int64_t> FlbaTimestampToInt64(const uint8_t* bytes,
+ bool clamp_on_overflow) {
+ const uint64_t low = bit_util::FromLittleEndian(SafeLoadAs<uint64_t>(bytes));
+ const uint32_t high = bit_util::FromLittleEndian(SafeLoadAs<uint32_t>(bytes
+ 8));
+ const int32_t high_signed = static_cast<int32_t>(high);
+ const int64_t low_signed = static_cast<int64_t>(low);
+ const int32_t sign_extension = (low_signed < 0) ? -1 : 0;
+ // Fits in int64 iff the high part is a pure sign-extension of the low part.
+ if (ARROW_PREDICT_FALSE(high_signed != sign_extension)) {
+ if (!clamp_on_overflow) {
+ return Status::Invalid(
+ "FLBA(12) TIMESTAMP value does not fit in a 64-bit Arrow timestamp");
+ }
+ return high_signed < 0 ? std::numeric_limits<int64_t>::min()
+ : std::numeric_limits<int64_t>::max();
+ }
+ return low_signed;
+}
+
+// Read a TIMESTAMP-annotated FLBA(12) column as a 64-bit Arrow timestamp.
+Result<Datum> TransferFlbaTimestamp(RecordReader* reader, MemoryPool* pool,
+ const std::shared_ptr<Field>& field,
+ bool clamp_on_overflow) {
+ auto binary_reader = dynamic_cast<BinaryRecordReader*>(reader);
+ DCHECK(binary_reader);
+ ::arrow::ArrayVector chunks = binary_reader->GetBuilderChunks();
+
+ for (size_t i = 0; i < chunks.size(); ++i) {
+ const auto& values = checked_cast<const
::arrow::FixedSizeBinaryArray&>(*chunks[i]);
+ const int64_t length = values.length();
+ ARROW_ASSIGN_OR_RAISE(auto data,
+ ::arrow::AllocateBuffer(length * sizeof(int64_t),
pool));
+ auto out_ptr = reinterpret_cast<int64_t*>(data->mutable_data());
+
+ int64_t j = 0;
+ RETURN_NOT_OK(::arrow::VisitArraySpanInline<::arrow::FixedSizeBinaryType>(
+ ::arrow::ArraySpan(*values.data()),
+ [&](std::string_view v) {
+ ARROW_ASSIGN_OR_RAISE(
+ out_ptr[j++],
+ FlbaTimestampToInt64(reinterpret_cast<const uint8_t*>(v.data()),
+ clamp_on_overflow));
+ return Status::OK();
+ },
+ [&]() {
+ out_ptr[j++] = 0;
+ return ::arrow::Status::OK();
+ }));
+
+ chunks[i] = std::make_shared<::arrow::TimestampArray>(
+ field->type(), length, std::move(data), values.null_bitmap(),
+ values.null_count());
+ }
+ if (!field->nullable()) {
+ ReconstructChunksWithoutNulls(&chunks);
+ }
+ return Datum(std::make_shared<ChunkedArray>(std::move(chunks),
field->type()));
+}
+
} // namespace
#define TRANSFER_INT32(ENUM, ArrowType)
\
@@ -964,6 +1027,18 @@ 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) {
+ // Validate that the provided Arrow timestamp unit matches the Parquet
unit.
+ DCHECK(descr->logical_type()->is_timestamp());
+ const auto& ts_logical =
+ checked_cast<const TimestampLogicalType&>(*descr->logical_type());
+ ARROW_ASSIGN_OR_RAISE(auto expected_unit,
+
ArrowTimeUnitFromParquet(ts_logical.time_unit()));
+ DCHECK_EQ(timestamp_type.unit(), expected_unit);
+ ARROW_ASSIGN_OR_RAISE(
+ result, TransferFlbaTimestamp(
+ reader, pool, value_field,
+
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 2e8cf764b2..d5be8d26c9 100644
--- a/cpp/src/parquet/arrow/schema_internal.cc
+++ b/cpp/src/parquet/arrow/schema_internal.cc
@@ -37,6 +37,20 @@ using ::arrow::Result;
using ::arrow::Status;
using ::arrow::internal::checked_cast;
+Result<::arrow::TimeUnit::type> ArrowTimeUnitFromParquet(
+ LogicalType::TimeUnit::unit unit) {
+ switch (unit) {
+ case LogicalType::TimeUnit::MILLIS:
+ return ::arrow::TimeUnit::MILLI;
+ case LogicalType::TimeUnit::MICROS:
+ return ::arrow::TimeUnit::MICRO;
+ case LogicalType::TimeUnit::NANOS:
+ return ::arrow::TimeUnit::NANO;
+ default:
+ return Status::TypeError("Unrecognized Parquet time unit");
+ }
+}
+
namespace {
Result<std::shared_ptr<ArrowType>> MakeArrowDecimal(const LogicalType&
logical_type,
@@ -105,20 +119,9 @@ Result<std::shared_ptr<ArrowType>>
MakeArrowTimestamp(const LogicalType& logical
const auto& timestamp = checked_cast<const
TimestampLogicalType&>(logical_type);
const bool utc_normalized = timestamp.is_adjusted_to_utc();
static const char* utc_timezone = "UTC";
- switch (timestamp.time_unit()) {
- case LogicalType::TimeUnit::MILLIS:
- return (utc_normalized ? ::arrow::timestamp(::arrow::TimeUnit::MILLI,
utc_timezone)
- : ::arrow::timestamp(::arrow::TimeUnit::MILLI));
- case LogicalType::TimeUnit::MICROS:
- return (utc_normalized ? ::arrow::timestamp(::arrow::TimeUnit::MICRO,
utc_timezone)
- : ::arrow::timestamp(::arrow::TimeUnit::MICRO));
- case LogicalType::TimeUnit::NANOS:
- return (utc_normalized ? ::arrow::timestamp(::arrow::TimeUnit::NANO,
utc_timezone)
- : ::arrow::timestamp(::arrow::TimeUnit::NANO));
- default:
- return Status::TypeError("Unrecognized time unit in timestamp
logical_type: ",
- logical_type.ToString());
- }
+ ARROW_ASSIGN_OR_RAISE(auto unit,
ArrowTimeUnitFromParquet(timestamp.time_unit()));
+ return utc_normalized ? ::arrow::timestamp(unit, utc_timezone)
+ : ::arrow::timestamp(unit);
}
Result<std::shared_ptr<ArrowType>> FromByteArray(
@@ -207,6 +210,12 @@ Result<std::shared_ptr<ArrowType>> FromFLBA(
return ::arrow::extension::uuid();
}
+ return ::arrow::fixed_size_binary(physical_length);
+ case LogicalType::Type::TIMESTAMP:
+ // If configured, convert to a potentially lossy Arrow timestamp.
+ 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/arrow/schema_internal.h
b/cpp/src/parquet/arrow/schema_internal.h
index 09ad891aad..50eb4bb403 100644
--- a/cpp/src/parquet/arrow/schema_internal.h
+++ b/cpp/src/parquet/arrow/schema_internal.h
@@ -18,6 +18,7 @@
#pragma once
#include "arrow/result.h"
+#include "arrow/type.h"
#include "arrow/type_fwd.h"
#include "parquet/schema.h"
@@ -25,6 +26,9 @@ namespace parquet::arrow {
using ::arrow::Result;
+Result<::arrow::TimeUnit::type> ArrowTimeUnitFromParquet(
+ LogicalType::TimeUnit::unit unit);
+
Result<std::shared_ptr<::arrow::DataType>> FromInt32(
const LogicalType& logical_type, const ArrowReaderProperties&
reader_properties);
Result<std::shared_ptr<::arrow::DataType>> FromInt64(
diff --git a/cpp/src/parquet/properties.h b/cpp/src/parquet/properties.h
index f8ad605dec..89f2be42bf 100644
--- a/cpp/src/parquet/properties.h
+++ b/cpp/src/parquet/properties.h
@@ -1173,7 +1173,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_(true),
+ flba_timestamp_clamp_on_overflow_(false) {}
/// \brief Set whether to use the IO thread pool to parse columns in
parallel.
///
@@ -1311,6 +1313,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<int> read_dict_indices_;
@@ -1324,6 +1349,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 d223ce0db6..e4d2fd1d59 100644
--- a/cpp/src/parquet/reader_test.cc
+++ b/cpp/src/parquet/reader_test.cc
@@ -140,6 +140,8 @@ std::string byte_stream_split_extended() {
std::string nested_lists() { return data_file("nested_lists.snappy.parquet"); }
+std::string flba12_timestamp() { return data_file("flba12_timestamp.parquet");
}
+
template <typename DType, typename ValueType = typename DType::c_type>
std::vector<ValueType> ReadColumnValues(ParquetFileReader* file_reader, int
row_group,
int column, int64_t
expected_values_read) {
@@ -1822,6 +1824,73 @@ 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<const
TimestampLogicalType&>(*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<TypedColumnReader<FLBAType>>(rg_reader->Column(c));
+ std::vector<FLBA> 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<const char*>(values[kMinRow].ptr), 12);
+ max_value.assign(reinterpret_cast<const char*>(values[kMaxRow].ptr), 12);
+
+ auto comparator = MakeComparator<FLBAType>(descr);
+ auto min_max = comparator->GetMinMax(values.data(), kNumRows);
+ ASSERT_EQ(min_value,
+ std::string(reinterpret_cast<const char*>(min_max.first.ptr),
12));
+ ASSERT_EQ(max_value,
+ std::string(reinterpret_cast<const char*>(min_max.second.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<int32_t> row_group_indices;
std::vector<int32_t> column_indices;
diff --git a/cpp/src/parquet/schema_test.cc b/cpp/src/parquet/schema_test.cc
index 6c8e6366ad..0ab957dfb9 100644
--- a/cpp/src/parquet/schema_test.cc
+++ b/cpp/src/parquet/schema_test.cc
@@ -1462,6 +1462,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 d43998ef78..71222db79f 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 {
@@ -463,6 +466,58 @@ struct RebindLogical<Float16LogicalType> {
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<Flba12TimestampType, /*is_signed=*/true> {
+ 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; }
+
+ // 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 int32_t a_hi =
+ SafeCopy<int32_t>(FromLittleEndian(SafeLoadAs<uint32_t>(a.ptr + 8)));
+ const int32_t b_hi =
+ SafeCopy<int32_t>(FromLittleEndian(SafeLoadAs<uint32_t>(b.ptr + 8)));
+ if (a_hi != b_hi) return a_hi < b_hi;
+ const uint64_t a_lo = FromLittleEndian(SafeLoadAs<uint64_t>(a.ptr));
+ const uint64_t b_lo = FromLittleEndian(SafeLoadAs<uint64_t>(b.ptr));
+ return a_lo < b_lo;
+ }
+
+ 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<Flba12TimestampType> {
+ using DType = FLBAType;
+ using c_type = DType::c_type;
+};
+
template <bool is_signed, typename DType>
class TypedComparatorImpl
: virtual public TypedComparator<typename RebindLogical<DType>::DType> {
@@ -1014,6 +1069,10 @@ std::shared_ptr<Comparator> DoMakeComparator(Type::type
physical_type,
return std::make_shared<TypedComparatorImpl<true,
Float16LogicalType>>(
type_length);
}
+ if (logical_type == LogicalType::Type::TIMESTAMP) {
+ return std::make_shared<TypedComparatorImpl<true,
Flba12TimestampType>>(
+ type_length);
+ }
return std::make_shared<TypedComparatorImpl<true,
FLBAType>>(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 bf0961c4fc..af9d96a0ad 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<FLBAType>(&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<uint8_t> large_neg_bytes = {0x00, 0x00, 0x00, 0x00, 0x00, 0x00,
+ 0x00, 0x00, 0x00, 0x00, 0x00, 0x80};
+ std::vector<uint8_t> neg256_bytes = {0x00, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF,
+ 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF};
+ std::vector<uint8_t> minus_one_bytes = {0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF,
+ 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF};
+ std::vector<uint8_t> zero_bytes = {0x00, 0x00, 0x00, 0x00, 0x00, 0x00,
+ 0x00, 0x00, 0x00, 0x00, 0x00, 0x00};
+ std::vector<uint8_t> plus_one_bytes = {0x01, 0x00, 0x00, 0x00, 0x00, 0x00,
+ 0x00, 0x00, 0x00, 0x00, 0x00, 0x00};
+ std::vector<uint8_t> large_pos_bytes = {0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF,
+ 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0x7F};
+
+ std::vector<FLBA> 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<FLBAType>(&descr);
+
+ std::vector<uint8_t> plus_one = {0x01, 0x00, 0x00, 0x00, 0x00, 0x00,
+ 0x00, 0x00, 0x00, 0x00, 0x00, 0x00};
+ std::vector<uint8_t> plus_five = {0x05, 0x00, 0x00, 0x00, 0x00, 0x00,
+ 0x00, 0x00, 0x00, 0x00, 0x00, 0x00};
+ std::vector<uint8_t> large_pos = {0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF,
+ 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0x7F};
+ std::vector<FLBA> 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 fda5e319e0..c754e8777c 100644
--- a/cpp/src/parquet/types.cc
+++ b/cpp/src/parquet/types.cc
@@ -1387,10 +1387,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;
@@ -1411,7 +1413,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),
@@ -1422,6 +1423,13 @@ 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_;
}
diff --git a/cpp/submodules/parquet-testing b/cpp/submodules/parquet-testing
index e74785d85a..56653c437c 160000
--- a/cpp/submodules/parquet-testing
+++ b/cpp/submodules/parquet-testing
@@ -1 +1 @@
-Subproject commit e74785d85a4ecee829e1e405444d6a1b24b8bc9c
+Subproject commit 56653c437c8092f704a092d0d1d4e600124cd49f