Skip to content

Commit a13b20f

Browse files
Deepak Majetijulienledem
authored andcommitted
PARQUET-499: Complete PlainEncoder implementation for all primitive types and test end to end
Includes tests for end to end plain encoding and decoding of all data types. Author: Deepak Majeti <deepak.majeti@hp.com> Closes apache#52 from majetideepak/PARQUET-499 and squashes the following commits: 897859b [Deepak Majeti] minor edits 2067ef5 [Deepak Majeti] renamed a test dfb19f8 [Deepak Majeti] templated all types 059967a [Deepak Majeti] templated int and real tests da86d4d [Deepak Majeti] minor fix 4976bec [Deepak Majeti] include pruning d0f8ab9 [Deepak Majeti] addressed comments 07257c0 [Deepak Majeti] minor format edits 6ca0b30 [Deepak Majeti] fixed formatting and casting issues 9815062 [Deepak Majeti] PARQUET-499 Change-Id: I45b2811e9abc8cad1277a533280d7fc3727d13e7
1 parent af04814 commit a13b20f

3 files changed

Lines changed: 265 additions & 19 deletions

File tree

cpp/src/parquet/encodings/plain-encoding-test.cc

Lines changed: 156 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@
2424

2525
#include "parquet/encodings/plain-encoding.h"
2626
#include "parquet/types.h"
27+
#include "parquet/schema/types.h"
2728
#include "parquet/util/bit-util.h"
2829
#include "parquet/util/output.h"
2930
#include "parquet/util/test-common.h"
@@ -35,11 +36,10 @@ namespace parquet_cpp {
3536

3637
namespace test {
3738

38-
TEST(BooleanTest, TestEncodeDecode) {
39+
TEST(VectorBooleanTest, TestEncodeDecode) {
3940
// PARQUET-454
40-
41-
size_t nvalues = 100;
42-
size_t nbytes = BitUtil::RoundUp(nvalues, 8) / 8;
41+
size_t nvalues = 10000;
42+
size_t nbytes = BitUtil::Ceil(nvalues, 8);
4343

4444
// seed the prng so failure is deterministic
4545
vector<bool> draws = flip_coins_seed(nvalues, 0.5, 0);
@@ -50,23 +50,171 @@ TEST(BooleanTest, TestEncodeDecode) {
5050
InMemoryOutputStream dst;
5151
encoder.Encode(draws, nvalues, &dst);
5252

53-
std::vector<uint8_t> encode_buffer;
53+
vector<uint8_t> encode_buffer;
5454
dst.Transfer(&encode_buffer);
5555

5656
ASSERT_EQ(nbytes, encode_buffer.size());
5757

58-
std::vector<uint8_t> decode_buffer(nbytes);
58+
vector<uint8_t> decode_buffer(nbytes);
5959
const uint8_t* decode_data = &decode_buffer[0];
6060

6161
decoder.SetData(nvalues, &encode_buffer[0], encode_buffer.size());
6262
size_t values_decoded = decoder.Decode(&decode_buffer[0], nvalues);
6363
ASSERT_EQ(nvalues, values_decoded);
6464

6565
for (size_t i = 0; i < nvalues; ++i) {
66-
ASSERT_EQ(BitUtil::GetArrayBit(decode_data, i), draws[i]) << i;
66+
ASSERT_EQ(draws[i], BitUtil::GetArrayBit(decode_data, i)) << i;
6767
}
6868
}
6969

70-
} // namespace test
70+
template<typename T, int TYPE>
71+
class EncodeDecode{
72+
public:
73+
void init_data(int nvalues) {
74+
num_values_ = nvalues;
75+
input_bytes_.resize(num_values_ * sizeof(T));
76+
output_bytes_.resize(num_values_ * sizeof(T));
77+
draws_ = reinterpret_cast<T*>(input_bytes_.data());
78+
decode_buf_ = reinterpret_cast<T*>(output_bytes_.data());
79+
}
80+
81+
void generate_data() {
82+
// seed the prng so failure is deterministic
83+
random_numbers(num_values_, 0.5, draws_);
84+
}
85+
86+
void encode_decode(ColumnDescriptor *d) {
87+
PlainEncoder<TYPE> encoder(d);
88+
PlainDecoder<TYPE> decoder(d);
89+
90+
InMemoryOutputStream dst;
91+
encoder.Encode(draws_, num_values_, &dst);
92+
93+
dst.Transfer(&encode_buffer_);
94+
95+
decoder.SetData(num_values_, &encode_buffer_[0], encode_buffer_.size());
96+
size_t values_decoded = decoder.Decode(decode_buf_, num_values_);
97+
ASSERT_EQ(num_values_, values_decoded);
98+
}
99+
100+
void verify_results() {
101+
for (size_t i = 0; i < num_values_; ++i) {
102+
ASSERT_EQ(draws_[i], decode_buf_[i]) << i;
103+
}
104+
}
105+
106+
void execute(int nvalues, ColumnDescriptor *d) {
107+
init_data(nvalues);
108+
generate_data();
109+
encode_decode(d);
110+
verify_results();
111+
}
112+
113+
private:
114+
int num_values_;
115+
T* draws_;
116+
T* decode_buf_;
117+
vector<uint8_t> input_bytes_;
118+
vector<uint8_t> output_bytes_;
119+
vector<uint8_t> data_buffer_;
120+
vector<uint8_t> encode_buffer_;
121+
};
122+
123+
template<>
124+
void EncodeDecode<bool, Type::BOOLEAN>::generate_data() {
125+
// seed the prng so failure is deterministic
126+
random_bools(num_values_, 0.5, 0, draws_);
127+
}
128+
129+
template<>
130+
void EncodeDecode<Int96, Type::INT96>::verify_results() {
131+
for (size_t i = 0; i < num_values_; ++i) {
132+
ASSERT_EQ(draws_[i].value[0], decode_buf_[i].value[0]) << i;
133+
ASSERT_EQ(draws_[i].value[1], decode_buf_[i].value[1]) << i;
134+
ASSERT_EQ(draws_[i].value[2], decode_buf_[i].value[2]) << i;
135+
}
136+
}
137+
138+
template<>
139+
void EncodeDecode<ByteArray, Type::BYTE_ARRAY>::generate_data() {
140+
// seed the prng so failure is deterministic
141+
int max_byte_array_len = 12 + sizeof(uint32_t);
142+
size_t nbytes = num_values_ * max_byte_array_len;
143+
data_buffer_.resize(nbytes);
144+
random_byte_array(num_values_, 0.5, data_buffer_.data(), draws_,
145+
max_byte_array_len);
146+
}
147+
148+
template<>
149+
void EncodeDecode<ByteArray, Type::BYTE_ARRAY>::verify_results() {
150+
for (size_t i = 0; i < num_values_; ++i) {
151+
ASSERT_EQ(draws_[i].len, decode_buf_[i].len) << i;
152+
ASSERT_EQ(0, memcmp(draws_[i].ptr, decode_buf_[i].ptr, draws_[i].len)) << i;
153+
}
154+
}
155+
156+
static int flba_length = 8;
157+
template<>
158+
void EncodeDecode<FLBA, Type::FIXED_LEN_BYTE_ARRAY>::generate_data() {
159+
// seed the prng so failure is deterministic
160+
size_t nbytes = num_values_ * flba_length;
161+
data_buffer_.resize(nbytes);
162+
ASSERT_EQ(nbytes, data_buffer_.size());
163+
random_fixed_byte_array(num_values_, 0.5, data_buffer_.data(), flba_length, draws_);
164+
}
165+
166+
template<>
167+
void EncodeDecode<FLBA, Type::FIXED_LEN_BYTE_ARRAY>::verify_results() {
168+
for (size_t i = 0; i < 1000; ++i) {
169+
ASSERT_EQ(0, memcmp(draws_[i].ptr, decode_buf_[i].ptr, flba_length)) << i;
170+
}
171+
}
71172

173+
int num_values = 10000;
174+
175+
TEST(BoolEncodeDecode, TestEncodeDecode) {
176+
EncodeDecode<bool, Type::BOOLEAN> obj;
177+
obj.execute(num_values, nullptr);
178+
}
179+
180+
TEST(Int32EncodeDecode, TestEncodeDecode) {
181+
EncodeDecode<int32_t, Type::INT32> obj;
182+
obj.execute(num_values, nullptr);
183+
}
184+
185+
TEST(Int64EncodeDecode, TestEncodeDecode) {
186+
EncodeDecode<int64_t, Type::INT64> obj;
187+
obj.execute(num_values, nullptr);
188+
}
189+
190+
TEST(FloatEncodeDecode, TestEncodeDecode) {
191+
EncodeDecode<float, Type::FLOAT> obj;
192+
obj.execute(num_values, nullptr);
193+
}
194+
195+
TEST(DoubleEncodeDecode, TestEncodeDecode) {
196+
EncodeDecode<double, Type::DOUBLE> obj;
197+
obj.execute(num_values, nullptr);
198+
}
199+
200+
TEST(Int96EncodeDecode, TestEncodeDecode) {
201+
EncodeDecode<Int96, Type::INT96> obj;
202+
obj.execute(num_values, nullptr);
203+
}
204+
205+
TEST(BAEncodeDecode, TestEncodeDecode) {
206+
EncodeDecode<ByteArray, Type::BYTE_ARRAY> obj;
207+
obj.execute(num_values, nullptr);
208+
}
209+
210+
TEST(FLBAEncodeDecode, TestEncodeDecode) {
211+
schema::NodePtr node;
212+
node = schema::PrimitiveNode::MakeFLBA("name", Repetition::OPTIONAL,
213+
Type::FIXED_LEN_BYTE_ARRAY, flba_length, LogicalType::UTF8);
214+
ColumnDescriptor d(node, 0, 0);
215+
EncodeDecode<FixedLenByteArray, Type::FIXED_LEN_BYTE_ARRAY> obj;
216+
obj.execute(num_values, &d);
217+
}
218+
219+
} // namespace test
72220
} // namespace parquet_cpp

cpp/src/parquet/encodings/plain-encoding.h

Lines changed: 29 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -67,22 +67,26 @@ inline int PlainDecoder<TYPE>::Decode(T* buffer, int max_values) {
6767
}
6868

6969
// Template specialization for BYTE_ARRAY
70+
// BA does not currently own its data
71+
// the lifetime is tied to the input stream
7072
template <>
7173
inline int PlainDecoder<Type::BYTE_ARRAY>::Decode(ByteArray* buffer,
7274
int max_values) {
7375
max_values = std::min(max_values, num_values_);
7476
for (int i = 0; i < max_values; ++i) {
75-
buffer[i].len = *reinterpret_cast<const uint32_t*>(data_);
76-
if (len_ < sizeof(uint32_t) + buffer[i].len) ParquetException::EofException();
77+
uint32_t len = buffer[i].len = *reinterpret_cast<const uint32_t*>(data_);
78+
if (len_ < sizeof(uint32_t) + len) ParquetException::EofException();
7779
buffer[i].ptr = data_ + sizeof(uint32_t);
78-
data_ += sizeof(uint32_t) + buffer[i].len;
79-
len_ -= sizeof(uint32_t) + buffer[i].len;
80+
data_ += sizeof(uint32_t) + len;
81+
len_ -= sizeof(uint32_t) + len;
8082
}
8183
num_values_ -= max_values;
8284
return max_values;
8385
}
8486

8587
// Template specialization for FIXED_LEN_BYTE_ARRAY
88+
// FLBA does not currently own its data
89+
// the lifetime is tied to the input stream
8690
template <>
8791
inline int PlainDecoder<Type::FIXED_LEN_BYTE_ARRAY>::Decode(
8892
FixedLenByteArray* buffer, int max_values) {
@@ -161,11 +165,21 @@ class PlainEncoder<Type::BOOLEAN> : public Encoder<Type::BOOLEAN> {
161165
Encoder<Type::BOOLEAN>(descr, Encoding::PLAIN) {}
162166

163167
virtual void Encode(const bool* src, int num_values, OutputStream* dst) {
164-
throw ParquetException("this API for encoding bools not implemented");
168+
size_t bytes_required = BitUtil::Ceil(num_values, 8);
169+
std::vector<uint8_t> tmp_buffer(bytes_required);
170+
171+
BitWriter bit_writer(&tmp_buffer[0], bytes_required);
172+
for (size_t i = 0; i < num_values; ++i) {
173+
bit_writer.PutValue(src[i], 1);
174+
}
175+
bit_writer.Flush();
176+
177+
// Write the result to the output stream
178+
dst->Write(bit_writer.buffer(), bit_writer.bytes_written());
165179
}
166180

167181
void Encode(const std::vector<bool>& src, int num_values, OutputStream* dst) {
168-
size_t bytes_required = BitUtil::RoundUp(num_values, 8) / 8;
182+
size_t bytes_required = BitUtil::Ceil(num_values, 8);
169183

170184
// TODO(wesm)
171185
// Use a temporary buffer for now and copy, because the BitWriter is not
@@ -193,15 +207,21 @@ inline void PlainEncoder<TYPE>::Encode(const T* buffer, int num_values,
193207
template <>
194208
inline void PlainEncoder<Type::BYTE_ARRAY>::Encode(const ByteArray* src,
195209
int num_values, OutputStream* dst) {
196-
ParquetException::NYI("byte array encoding");
210+
for (size_t i = 0; i < num_values; ++i) {
211+
// Write the result to the output stream
212+
dst->Write(reinterpret_cast<const uint8_t*>(&src[i].len), sizeof(uint32_t));
213+
dst->Write(reinterpret_cast<const uint8_t*>(src[i].ptr), src[i].len);
214+
}
197215
}
198216

199217
template <>
200218
inline void PlainEncoder<Type::FIXED_LEN_BYTE_ARRAY>::Encode(
201219
const FixedLenByteArray* src, int num_values, OutputStream* dst) {
202-
ParquetException::NYI("FLBA encoding");
220+
for (size_t i = 0; i < num_values; ++i) {
221+
// Write the result to the output stream
222+
dst->Write(reinterpret_cast<const uint8_t*>(src[i].ptr), descr_->type_length());
223+
}
203224
}
204-
205225
} // namespace parquet_cpp
206226

207227
#endif

cpp/src/parquet/util/test-common.h

Lines changed: 80 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -19,9 +19,12 @@
1919
#define PARQUET_UTIL_TEST_COMMON_H
2020

2121
#include <iostream>
22+
#include <limits>
2223
#include <random>
2324
#include <vector>
2425

26+
#include "parquet/types.h"
27+
2528
using std::vector;
2629

2730
namespace parquet_cpp {
@@ -81,7 +84,6 @@ static inline vector<bool> flip_coins_seed(size_t n, double p, uint32_t seed) {
8184
return draws;
8285
}
8386

84-
8587
static inline vector<bool> flip_coins(size_t n, double p) {
8688
std::random_device rd;
8789
std::mt19937 gen(rd());
@@ -104,8 +106,84 @@ void random_bytes(int n, uint32_t seed, std::vector<uint8_t>* out) {
104106
}
105107
}
106108

107-
} // namespace test
109+
template <typename T>
110+
void random_numbers(int n, uint32_t seed, T* out) {
111+
std::mt19937 gen(seed);
112+
std::uniform_real_distribution<T> d(std::numeric_limits<T>::lowest(),
113+
std::numeric_limits<T>::max());
114+
for (int i = 0; i < n; ++i) {
115+
out[i] = d(gen);
116+
}
117+
}
118+
119+
void random_bools(int n, double p, uint32_t seed, bool* out) {
120+
std::mt19937 gen(seed);
121+
std::bernoulli_distribution d(p);
122+
for (int i = 0; i < n; ++i) {
123+
out[i] = d(gen);
124+
}
125+
}
126+
127+
template <>
128+
void random_numbers(int n, uint32_t seed, int32_t* out) {
129+
std::mt19937 gen(seed);
130+
std::uniform_int_distribution<int32_t> d(std::numeric_limits<int32_t>::lowest(),
131+
std::numeric_limits<int32_t>::max());
132+
for (int i = 0; i < n; ++i) {
133+
out[i] = d(gen);
134+
}
135+
}
136+
137+
template <>
138+
void random_numbers(int n, uint32_t seed, int64_t* out) {
139+
std::mt19937 gen(seed);
140+
std::uniform_int_distribution<int64_t> d(std::numeric_limits<int64_t>::lowest(),
141+
std::numeric_limits<int64_t>::max());
142+
for (int i = 0; i < n; ++i) {
143+
out[i] = d(gen);
144+
}
145+
}
108146

147+
template <>
148+
void random_numbers(int n, uint32_t seed, Int96* out) {
149+
std::mt19937 gen(seed);
150+
std::uniform_int_distribution<uint32_t> d(std::numeric_limits<uint32_t>::lowest(),
151+
std::numeric_limits<uint32_t>::max());
152+
for (int i = 0; i < n; ++i) {
153+
out[i].value[0] = d(gen);
154+
out[i].value[1] = d(gen);
155+
out[i].value[2] = d(gen);
156+
}
157+
}
158+
159+
void random_fixed_byte_array(int n, uint32_t seed, uint8_t *buf, int len,
160+
FLBA* out) {
161+
std::mt19937 gen(seed);
162+
std::uniform_int_distribution<int> d(0, 255);
163+
for (int i = 0; i < n; ++i) {
164+
out[i].ptr = buf;
165+
for (int j = 0; j < len; ++j) {
166+
buf[j] = d(gen) & 0xFF;
167+
}
168+
buf += len;
169+
}
170+
}
171+
172+
void random_byte_array(int n, uint32_t seed, uint8_t *buf,
173+
ByteArray* out, int max_size) {
174+
std::mt19937 gen(seed);
175+
std::uniform_int_distribution<int> d1(0, max_size);
176+
std::uniform_int_distribution<int> d2(0, 255);
177+
for (int i = 0; i < n; ++i) {
178+
out[i].len = d1(gen);
179+
out[i].ptr = buf;
180+
for (int j = 0; j < out[i].len; ++j) {
181+
buf[j] = d2(gen) & 0xFF;
182+
}
183+
buf += out[i].len;
184+
}
185+
}
186+
} // namespace test
109187
} // namespace parquet_cpp
110188

111189
#endif // PARQUET_UTIL_TEST_COMMON_H

0 commit comments

Comments
 (0)