Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 42 additions & 0 deletions cpp/src/parquet/arrow/reader_writer_benchmark.cc
Original file line number Diff line number Diff line change
Expand Up @@ -273,6 +273,48 @@ BENCHMARK(BM_WriteBinaryColumn)
->Args({50, kInfiniteUniqueValues})
->Args({99, kInfiniteUniqueValues});

// Writes int32 columns with content-defined chunking into a single row group, with
// use_threads the columns are chunked by different threads in parallel
static void BM_WriteContentDefinedChunking(::benchmark::State& state) {
const auto num_columns = static_cast<int>(state.range(0));
const bool use_threads = state.range(1) != 0;
constexpr int64_t kNumRows = 1024 * 1024;

::arrow::random::RandomArrayGenerator generator(/*seed=*/500);
::arrow::FieldVector fields;
::arrow::ArrayVector columns;
for (int i = 0; i < num_columns; i++) {
fields.push_back(::arrow::field("column" + std::to_string(i), ::arrow::int32(),
/*nullable=*/false));
columns.push_back(generator.Int32(kNumRows, /*min=*/0, /*max=*/1 << 30));
}
auto schema = ::arrow::schema(fields);
auto batch = ::arrow::RecordBatch::Make(schema, kNumRows, columns);
auto properties = WriterProperties::Builder()
.enable_content_defined_chunking()
->disable_dictionary()
->build();
auto arrow_properties =
ArrowWriterProperties::Builder().set_use_threads(use_threads)->build();

for (auto _ : state) {
auto output = CreateOutputStream();
auto writer = arrow::FileWriter::Open(*schema, ::arrow::default_memory_pool(), output,
properties, arrow_properties);
EXIT_NOT_OK(writer.status());
EXIT_NOT_OK((*writer)->WriteRecordBatch(*batch));
EXIT_NOT_OK((*writer)->Close());
}
state.SetBytesProcessed(state.iterations() * num_columns * kNumRows *
static_cast<int64_t>(sizeof(int32_t)));
}

BENCHMARK(BM_WriteContentDefinedChunking)
->ArgNames({"columns", "use_threads"})
->Args({16, 0})
->Args({16, 1})
->UseRealTime();

template <typename T>
struct Examples {
static constexpr std::array<T, 2> values() { return {127, 128}; }
Expand Down
133 changes: 82 additions & 51 deletions cpp/src/parquet/chunker_internal.cc
Original file line number Diff line number Diff line change
Expand Up @@ -112,14 +112,11 @@ uint64_t CalculateMask(int64_t min_chunk_size, int64_t max_chunk_size, int norm_
}
}

} // namespace

class ContentDefinedChunker::Impl {
/// The gear hash of a column together with the state deciding the chunk boundaries
class GearHash {
public:
Impl(const LevelInfo& level_info, int64_t min_chunk_size, int64_t max_chunk_size,
int norm_level)
: level_info_(level_info),
min_chunk_size_(min_chunk_size),
GearHash(int64_t min_chunk_size, int64_t max_chunk_size, int norm_level)
: min_chunk_size_(min_chunk_size),
max_chunk_size_(max_chunk_size),
rolling_hash_mask_(CalculateMask(min_chunk_size, max_chunk_size, norm_level)) {}

Expand Down Expand Up @@ -201,6 +198,40 @@ class ContentDefinedChunker::Impl {
return false;
}

private:
// Minimum chunk size in bytes, the rolling hash will not be updated until this size is
// reached for each chunk. Note that all data sent through the hash function is counted
// towards the chunk size, including definition and repetition levels.
int64_t min_chunk_size_;
int64_t max_chunk_size_;
// The mask to match the rolling hash against to determine if a new chunk should be
// created. The mask is calculated based on min/max chunk size and the normalization
// level.
uint64_t rolling_hash_mask_;

// Whether the rolling hash has matched the mask since the last chunk creation. This
// flag is set true by the Roll() function when the mask is matched and reset to false
// by NeedNewChunk() method.
bool has_matched_ = false;
// The current run of the rolling hash, used to normalize the chunk size distribution
// by requiring multiple consecutive matches to create a new chunk.
int8_t nth_run_ = 0;
// Current chunk size in bytes, reset to 0 when a new chunk is created.
int64_t chunk_size_ = 0;
// Rolling hash state, never reset only initialized once for the entire column.
uint64_t rolling_hash_ = 0;
};

} // namespace

class ContentDefinedChunker::Impl {
public:
Impl(const LevelInfo& level_info, int64_t min_chunk_size, int64_t max_chunk_size,
int norm_level)
: level_info_(level_info), gearhash_(min_chunk_size, max_chunk_size, norm_level) {}

uint64_t GetRollingHashMask() const { return gearhash_.GetRollingHashMask(); }

void ValidateChunks(const std::vector<Chunk>& chunks, int64_t num_levels) const {
// chunks must be non-empty and monotonic increasing
ARROW_DCHECK(!chunks.empty());
Expand Down Expand Up @@ -248,6 +279,9 @@ class ContentDefinedChunker::Impl {
// requirements, we create a new chunk. See the `NeedNewChunk()` method for more
// details.
std::vector<Chunk> chunks;
// Roll a local copy of the hash so the compiler can keep it in registers instead of
// storing it to memory for every hashed byte
GearHash gearhash = gearhash_;
int64_t offset;
int64_t prev_offset = 0;
int64_t prev_value_offset = 0;
Expand All @@ -257,8 +291,8 @@ class ContentDefinedChunker::Impl {
if (!has_rep_levels && !has_def_levels) {
// fastest path for non-nested non-null data
for (offset = 0; offset < num_levels; ++offset) {
RollValue(offset);
if (NeedNewChunk()) {
RollValue(gearhash, offset);
if (gearhash.NeedNewChunk()) {
chunks.push_back({prev_offset, prev_offset, offset - prev_offset});
prev_offset = offset;
}
Expand All @@ -271,11 +305,11 @@ class ContentDefinedChunker::Impl {
for (int64_t offset = 0; offset < num_levels; ++offset) {
def_level = def_levels[offset];

Roll(&def_level);
gearhash.Roll(&def_level);
if (def_level == level_info_.def_level) {
RollValue(offset);
RollValue(gearhash, offset);
}
if (NeedNewChunk()) {
if (gearhash.NeedNewChunk()) {
chunks.push_back({prev_offset, prev_offset, offset - prev_offset});
prev_offset = offset;
}
Expand All @@ -292,13 +326,13 @@ class ContentDefinedChunker::Impl {
def_level = def_levels[offset];
rep_level = rep_levels[offset];

Roll(&def_level);
Roll(&rep_level);
gearhash.Roll(&def_level);
gearhash.Roll(&rep_level);
if (def_level == level_info_.def_level) {
RollValue(value_offset);
RollValue(gearhash, value_offset);
}

if (rep_level == 0 && NeedNewChunk()) {
if (rep_level == 0 && gearhash.NeedNewChunk()) {
// if we are at a record boundary and need a new chunk, we create a new chunk
auto levels_to_write = offset - prev_offset;
if (levels_to_write > 0) {
Expand All @@ -313,6 +347,7 @@ class ContentDefinedChunker::Impl {
}
}
}
gearhash_ = gearhash;

// add the last chunk if we have any levels left
if (prev_offset < num_levels) {
Expand All @@ -332,38 +367,43 @@ class ContentDefinedChunker::Impl {
const uint8_t* raw_values =
values.data()->GetValues<uint8_t>(/*i=*/1, /*absolute_offset=*/0) +
values.offset() * kByteWidth;
return Calculate(def_levels, rep_levels, num_levels, [&](int64_t i) {
return Roll<kByteWidth>(&raw_values[i * kByteWidth]);
});
return Calculate(def_levels, rep_levels, num_levels,
[&](GearHash& gearhash, int64_t i) {
return gearhash.Roll<kByteWidth>(&raw_values[i * kByteWidth]);
});
}

template <typename ArrayType>
std::vector<Chunk> CalculateBinaryLike(const int16_t* def_levels,
const int16_t* rep_levels, int64_t num_levels,
const ::arrow::Array& values) {
const auto& array = checked_cast<const ArrayType&>(values);
return Calculate(def_levels, rep_levels, num_levels, [&](int64_t i) {
typename ArrayType::offset_type length;
const uint8_t* value = array.GetValue(i, &length);
Roll(value, length);
});
return Calculate(def_levels, rep_levels, num_levels,
[&](GearHash& gearhash, int64_t i) {
typename ArrayType::offset_type length;
const uint8_t* value = array.GetValue(i, &length);
gearhash.Roll(value, length);
});
}

std::vector<Chunk> GetChunks(const int16_t* def_levels, const int16_t* rep_levels,
int64_t num_levels, const ::arrow::Array& values) {
auto handle_type = [&](auto&& type) -> std::vector<Chunk> {
using ArrowType = std::decay_t<decltype(type)>;
if constexpr (ArrowType::type_id == ::arrow::Type::NA) {
return Calculate(def_levels, rep_levels, num_levels, [](int64_t) {});
return Calculate(def_levels, rep_levels, num_levels, [](GearHash&, int64_t) {});
} else if constexpr (ArrowType::type_id == ::arrow::Type::BOOL) {
const auto& array = static_cast<const ::arrow::BooleanArray&>(values);
return Calculate(def_levels, rep_levels, num_levels,
[&](int64_t i) { Roll(array.Value(i)); });
return Calculate(
def_levels, rep_levels, num_levels,
[&](GearHash& gearhash, int64_t i) { gearhash.Roll(array.Value(i)); });
} else if constexpr (ArrowType::type_id == ::arrow::Type::FIXED_SIZE_BINARY) {
const auto& array = static_cast<const ::arrow::FixedSizeBinaryArray&>(values);
const auto byte_width = array.byte_width();
return Calculate(def_levels, rep_levels, num_levels,
[&](int64_t i) { Roll(array.GetValue(i), byte_width); });
[&](GearHash& gearhash, int64_t i) {
gearhash.Roll(array.GetValue(i), byte_width);
});
} else if constexpr (ArrowType::type_id == ::arrow::Type::EXTENSION) {
const auto& array = static_cast<const ::arrow::ExtensionArray&>(values);
return GetChunks(def_levels, rep_levels, num_levels, *array.storage());
Expand Down Expand Up @@ -392,38 +432,29 @@ class ContentDefinedChunker::Impl {
}

private:
// Reference to the column's level information
const internal::LevelInfo& level_info_;
// Minimum chunk size in bytes, the rolling hash will not be updated until this size is
// reached for each chunk. Note that all data sent through the hash function is counted
// towards the chunk size, including definition and repetition levels.
const int64_t min_chunk_size_;
const int64_t max_chunk_size_;
// The mask to match the rolling hash against to determine if a new chunk should be
// created. The mask is calculated based on min/max chunk size and the normalization
// level.
const uint64_t rolling_hash_mask_;

// Whether the rolling hash has matched the mask since the last chunk creation. This
// flag is set true by the Roll() function when the mask is matched and reset to false
// by NeedNewChunk() method.
bool has_matched_ = false;
// The current run of the rolling hash, used to normalize the chunk size distribution
// by requiring multiple consecutive matches to create a new chunk.
int8_t nth_run_ = 0;
// Current chunk size in bytes, reset to 0 when a new chunk is created.
int64_t chunk_size_ = 0;
// Rolling hash state, never reset only initialized once for the entire column.
uint64_t rolling_hash_ = 0;
// The column's level information
const internal::LevelInfo level_info_;
// The gear hash carried over between the calls, so the chunking continues across
// the pages and the row groups
GearHash gearhash_;
};

ContentDefinedChunker::ContentDefinedChunker(const LevelInfo& level_info,
int64_t min_chunk_size,
int64_t max_chunk_size, int norm_level)
: impl_(new Impl(level_info, min_chunk_size, max_chunk_size, norm_level)) {}

ContentDefinedChunker::ContentDefinedChunker(ContentDefinedChunker&&) noexcept = default;
ContentDefinedChunker& ContentDefinedChunker::operator=(
ContentDefinedChunker&&) noexcept = default;
ContentDefinedChunker::~ContentDefinedChunker() = default;

ContentDefinedChunker ContentDefinedChunker::Make(const LevelInfo& level_info,
const CdcOptions& options) {
return ContentDefinedChunker(level_info, options.min_chunk_size, options.max_chunk_size,
options.norm_level);
}

std::vector<Chunk> ContentDefinedChunker::GetChunks(const int16_t* def_levels,
const int16_t* rep_levels,
int64_t num_levels,
Expand Down
18 changes: 14 additions & 4 deletions cpp/src/parquet/chunker_internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@

#include "arrow/array.h"
#include "parquet/level_conversion.h"
#include "parquet/properties.h"

namespace parquet::internal {

Expand Down Expand Up @@ -72,10 +73,10 @@ struct Chunk {
/// Implementation details:
///
/// Only the parquet writer must be aware of the content defined chunking, the reader
/// doesn't need to know about it. Each parquet column writer holds a
/// ContentDefinedChunker instance depending on the writer's properties. The chunker's
/// state is maintained across the entire column without being reset between pages and row
/// groups.
/// doesn't need to know about it. The parquet file writer holds one
/// ContentDefinedChunker per leaf column depending on the writer's properties, and passes
/// it to the column writers of every row group. The chunker's state is maintained
/// across the entire column without being reset between pages and row groups.
///
/// The chunker receives the record shredded column data (def_levels, rep_levels, values)
/// and goes over the (def_level, rep_level, value) triplets one by one while adjusting
Expand Down Expand Up @@ -118,8 +119,17 @@ class PARQUET_EXPORT ContentDefinedChunker {
/// expense of fragmentation.
ContentDefinedChunker(const LevelInfo& level_info, int64_t min_chunk_size,
int64_t max_chunk_size, int norm_level = 0);
ContentDefinedChunker(ContentDefinedChunker&&) noexcept;
ContentDefinedChunker& operator=(ContentDefinedChunker&&) noexcept;
~ContentDefinedChunker();

/// Create a new ContentDefinedChunker instance using the given chunking options
///
/// @param level_info Information about definition and repetition levels
/// @param options Content defined chunking options
static ContentDefinedChunker Make(const LevelInfo& level_info,
const CdcOptions& options);

/// Get the chunk boundaries for the given column data
///
/// @param def_levels Definition levels
Expand Down
49 changes: 49 additions & 0 deletions cpp/src/parquet/chunker_internal_test.cc
Original file line number Diff line number Diff line change
Expand Up @@ -403,6 +403,20 @@ ParquetInfo GetColumnParquetInfo(const std::shared_ptr<Buffer>& data, int column
return result;
}

// The offsets where the data pages of a column end, counted in levels from the start of
// the column across its row groups, which are rows for flat columns
std::vector<int64_t> GetPageEnds(const std::shared_ptr<Buffer>& data, int column_index) {
std::vector<int64_t> page_ends;
int64_t offset = 0;
for (const auto& row_group : GetColumnParquetInfo(data, column_index)) {
for (auto page_length : row_group.page_lengths) {
offset += page_length;
page_ends.push_back(offset);
}
}
return page_ends;
}

// A git-hunk like side-by-side data structure to represent the differences between two
// vectors of uint64_t values.
using ChunkDiff = std::pair<ChunkList, ChunkList>;
Expand Down Expand Up @@ -1700,4 +1714,39 @@ TEST_F(TestCDCMultipleRowGroups, Append) {
}
}

TEST_F(TestCDCMultipleRowGroups, IndependentOfRowGroupBoundaries) {
// Splitting the data into row groups must not move the content defined page
// boundaries, only add one at the end of each row group. For example, if the pages of
// a single row group file end at rows 100, 250 and 400, then with row groups of 200
// rows the pages must end at rows 100, 200, 250 and 400.
ASSERT_OK_AND_ASSIGN(auto table, ConcatAndCombine({part1_, part2_, part3_}));
const int64_t num_rows = table->num_rows();
const int64_t row_group_length = num_rows / 6;
ASSERT_OK_AND_ASSIGN(auto single,
WriteTableToBuffer(table, kMinChunkSize, kMaxChunkSize,
/*row_group_length=*/num_rows));
ASSERT_OK_AND_ASSIGN(auto multi, WriteTableToBuffer(table, kMinChunkSize, kMaxChunkSize,
row_group_length));
ASSERT_EQ(ReadMetaData(std::make_shared<BufferReader>(single))->num_row_groups(), 1);
ASSERT_EQ(ReadMetaData(std::make_shared<BufferReader>(multi))->num_row_groups(), 6);
// compare the page ends column by column
for (int col = 0; col < table->num_columns(); col++) {
// the page ends of the single row group file, e.g. 100, 250 and 400
auto single_page_ends = GetPageEnds(single, col);
// the page ends of the multiple row group file, e.g. 100, 200, 250 and 400
auto multi_page_ends = GetPageEnds(multi, col);
// expect the page ends of the single row group file plus one at each of the 5
// boundaries between the row groups, e.g. 200, the last row group ends with the data
// where the single row group file's last page ends too
auto expected = single_page_ends;
for (int i = 1; i < 6; i++) {
expected.push_back(i * row_group_length);
}
std::sort(expected.begin(), expected.end());
EXPECT_EQ(multi_page_ends.size(), single_page_ends.size() + 5) << "column " << col;
// unlike ASSERT_EQ, ContainerEq prints the page ends that differ
ASSERT_THAT(multi_page_ends, ::testing::ContainerEq(expected)) << "column " << col;
}
}

} // namespace parquet::internal
Loading
Loading