Skip to content
Merged
117 changes: 117 additions & 0 deletions cpp/src/parquet/arrow/arrow_reader_writer_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -2166,6 +2166,123 @@ TEST(TestArrowReadWrite, CoerceTimestampsLosePrecision) {
allow_truncation_to_micros));
}
Comment thread
divjotarora marked this conversation as resolved.

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());
Comment thread
divjotarora marked this conversation as resolved.
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;
Expand Down
32 changes: 32 additions & 0 deletions cpp/src/parquet/arrow/arrow_schema_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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()},
Expand Down Expand Up @@ -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;
Expand Down
77 changes: 76 additions & 1 deletion cpp/src/parquet/arrow/reader_internal.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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>
Expand All @@ -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"
Expand Down Expand Up @@ -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) \
Expand Down Expand Up @@ -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 =
Comment thread
divjotarora marked this conversation as resolved.
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:
Expand Down
37 changes: 23 additions & 14 deletions cpp/src/parquet/arrow/schema_internal.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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);
Comment thread
divjotarora marked this conversation as resolved.
default:
return Status::NotImplemented("Unhandled logical_type ", logical_type.ToString(),
Expand Down
4 changes: 4 additions & 0 deletions cpp/src/parquet/arrow/schema_internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -18,13 +18,17 @@
#pragma once

#include "arrow/result.h"
#include "arrow/type.h"
#include "arrow/type_fwd.h"
#include "parquet/schema.h"

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(
Expand Down
Loading
Loading