如何将uint64_t以DECIMAL(30,0)逻辑类型写入Parquet文件?
问题
我需要将uint64_t值以逻辑类型DECIMAL(30, 0)、物理类型FIXED_LEN_BYTE_ARRAY写入Parquet文件,以下是我的尝试及遇到的问题:
由于parquet::StreamWriter要求FIXED_LEN_BYTE_ARRAY列的逻辑类型必须为LogicalType::None,无法使用StreamWriter中定义的>>运算符。因此我自定义了MyStreamWriter类,接收uint64_t值并将其转换为FixedLenByteArray,再通过FixedLenByteArrayWriter写入。
但当我写入uint64_t val = 2时,Parquet文件中存储的值变为2658455991569831745807614120560689152,显然不正确。该错误值的二进制表示为100000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000000,我推测写入逻辑将正确比特放在最高有效位,但其余位置被补零。
请问我哪里出错了?有没有更直接的方法将uint64_t作为小数写入Parquet文件?
#include "arrow/util/decimal.h" #include "arrow/io/file.h" #include "arrow/type.h" #include "parquet/schema.h" #include "parquet/stream_writer.h" class MyStreamWriter : public parquet::StreamWriter { using parquet::StreamWriter::StreamWriter; public: MyStreamWriter& writeDecimal(uint64_t v) { arrow::Decimal128 decimal(arrow::BasicDecimal128(0, v)); std::array<uint8_t, 16> decimalBytes = decimal.ToBytes(); const uint8_t* decimalBytesPtr = reinterpret_cast<const uint8_t*>(decimalBytes.data()); parquet::FixedLenByteArray flba(decimalBytesPtr); Write<parquet::FixedLenByteArrayWriter>(flba); return *this; } }; // Define our schema. parquet::schema::NodeVector fields; int32_t precision = 30; int32_t scale = 0; fields.push_back(parquet::schema::PrimitiveNode::Make( "my_decimal_col", parquet::Repetition::REQUIRED, parquet::LogicalType::Decimal(precision, scale), parquet::Type::FIXED_LEN_BYTE_ARRAY, arrow::Decimal128Type(precision, scale).byte_width())); auto schema = static_pointer_cast<parquet::schema::GroupNode>( parquet::schema::GroupNode::Make("schema", parquet::Repetition::REPEATED, fields)); // Open the writer. const shared_ptr<arrow::io::OutputStream> os = arrow::io::FileOutputStream::Open(_filepath).ValueOrDie(); writer = make_unique<MyStreamWriter>( parquet::ParquetFileWriter::Open(os, schema)); // Try writing an arbitrary value. uint64_t val = 100; writer->WriteDecimal(val);
解答
错误原因分析
- Decimal128构造与字节序问题:直接用
arrow::BasicDecimal128(0, v)构造Decimal128时,未考虑Decimal128的有符号属性,若uint64_t值超过INT64_MAX会被解析为负数,导致二进制补码完全错误;同时创建FixedLenByteArray时未明确指定长度,可能引发字节数组解析异常。 - Schema根节点类型错误:将根GroupNode设为
REPETITION::REPEATED,会导致表结构变为数组而非普通行结构,引发解析逻辑混乱。 - FixedLenByteArray构造不规范:仅传入指针的构造函数无法确保字节数组长度符合DECIMAL(30,0)要求的16字节标准。
修复方案
方案一:直接构造大端字节数组(无需Decimal128)
uint64_t最大值远小于DECIMAL(30,0)的上限,可直接将uint64_t转为大端字节序的16字节数组(前8字节补0):
class MyStreamWriter : public parquet::StreamWriter { using parquet::StreamWriter::StreamWriter; public: MyStreamWriter& writeDecimal(uint64_t v) { std::array<uint8_t, 16> bytes = {}; // 将uint64_t转为大端格式,存入数组后8位 for (int i = 0; i < 8; ++i) { bytes[8 + 7 - i] = static_cast<uint8_t>((v >> (i * 8)) & 0xFF); } // 明确指定字节数组长度为16 parquet::FixedLenByteArray flba(bytes.data(), 16); Write<parquet::FixedLenByteArrayWriter>(flba); return *this; } };
方案二:正确使用Arrow Decimal128转换
利用Arrow提供的Decimal128::FromUint64方法,安全处理无符号转有符号Decimal128的逻辑:
class MyStreamWriter : public parquet::StreamWriter { using parquet::StreamWriter::StreamWriter; public: MyStreamWriter& writeDecimal(uint64_t v) { arrow::Decimal128 decimal; auto status = arrow::Decimal128::FromUint64(v, &decimal); if (!status.ok()) { throw std::runtime_error("Failed to convert uint64_t to Decimal128: " + status.message()); } std::array<uint8_t, 16> decimalBytes = decimal.ToBytes(); parquet::FixedLenByteArray flba(decimalBytes.data(), 16); Write<parquet::FixedLenByteArrayWriter>(flba); return *this; } };
修正Schema根节点类型
将根GroupNode的重复类型改为REQUIRED,符合Parquet标准表结构:
auto schema = static_pointer_cast<parquet::schema::GroupNode>( parquet::schema::GroupNode::Make("schema", parquet::Repetition::REQUIRED, fields));
更直接的写入方式
跳过自定义StreamWriter,直接使用FixedLenByteArrayWriter写入:
// 获取对应列的writer auto col_writer = writer->column(0)->mutable_writer()->As<parquet::FixedLenByteArrayWriter>(); // 转换uint64_t为大端16字节数组 std::array<uint8_t, 16> bytes = {}; for (int i = 0; i < 8; ++i) { bytes[8 + 7 - i] = static_cast<uint8_t>((val >> (i * 8)) & 0xFF); } parquet::FixedLenByteArray flba(bytes.data(), 16); // 写入单条数据 col_writer->WriteBatch(1, nullptr, nullptr, &flba);
内容的提问来源于stack exchange,提问作者1step1leap

