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
Original file line number Diff line number Diff line change
@@ -0,0 +1,159 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing,
* software distributed under the License is distributed on an
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
* KIND, either express or implied. See the License for the
* specific language governing permissions and limitations
* under the License.
*/
package org.apache.parquet.benchmarks;

import java.io.IOException;
import java.util.ArrayList;
import java.util.List;
import java.util.Random;
import java.util.concurrent.TimeUnit;
import org.apache.parquet.example.data.Group;
import org.apache.parquet.example.data.simple.SimpleGroupFactory;
import org.apache.parquet.hadoop.ParquetFileWriter;
import org.apache.parquet.hadoop.ParquetWriter;
import org.apache.parquet.hadoop.example.ExampleParquetWriter;
import org.apache.parquet.hadoop.metadata.CompressionCodecName;
import org.apache.parquet.io.api.Binary;
import org.apache.parquet.schema.MessageType;
import org.apache.parquet.schema.MessageTypeParser;
import org.openjdk.jmh.annotations.Benchmark;
import org.openjdk.jmh.annotations.BenchmarkMode;
import org.openjdk.jmh.annotations.Fork;
import org.openjdk.jmh.annotations.Level;
import org.openjdk.jmh.annotations.Measurement;
import org.openjdk.jmh.annotations.Mode;
import org.openjdk.jmh.annotations.OutputTimeUnit;
import org.openjdk.jmh.annotations.Param;
import org.openjdk.jmh.annotations.Scope;
import org.openjdk.jmh.annotations.Setup;
import org.openjdk.jmh.annotations.State;
import org.openjdk.jmh.annotations.Warmup;

/**
* The write cost of content defined chunking: the same rows written with and without it, for rows of
* several shapes, since the cost is the hashing of every value's and level's bytes.
*
* <p>Both variants lift the page row count limit and turn dictionary encoding off: chunking changes
* where pages end and when a dictionary falls back, either of which would dwarf the hashing. Rows are
* built once and written to {@link BlackHoleOutputFile}, so neither data generation nor I/O is
* measured. The cost while disabled is {@link WriteBenchmarks} compared across revisions.
*/
@BenchmarkMode(Mode.AverageTime)
@Fork(1)
@Warmup(iterations = 3)
@Measurement(iterations = 5)
@OutputTimeUnit(TimeUnit.MILLISECONDS)
@State(Scope.Thread)
public class CdcWriteBenchmarks {

private static final int ROW_COUNT = 100_000;

@Param({"false", "true"})
public boolean chunking;

/**
* mixed: a long, a 128-byte binary and a list of eight ints; numbers: an int, a long and a nullable
* double; strings: short strings, one of them nullable; lists: nullable lists of zero to eight
* nullable longs.
*/
@Param({"mixed", "numbers", "strings", "lists"})
public String data;

private MessageType schema;
private List<Group> rows;

@Setup(Level.Trial)
public void setup() {
Random random = new Random(TestDataFactory.DEFAULT_SEED);
switch (data) {
case "mixed":
schema = MessageTypeParser.parseMessageType(
"message m { required int64 l; required binary b; required group g { repeated int32 i; } }");
break;
case "numbers":
schema = MessageTypeParser.parseMessageType(
"message m { required int32 i; required int64 l; optional double d; }");
break;
case "strings":
schema = MessageTypeParser.parseMessageType(
"message m { required binary s (STRING); optional binary t (STRING); }");
break;
case "lists":
schema = MessageTypeParser.parseMessageType(
"message m { optional group l (LIST) { repeated group list { optional int64 element; } } }");
break;
default:
throw new IllegalArgumentException("unknown data " + data);
}
Binary[] binaries = TestDataFactory.generateBinaryData(ROW_COUNT, 128, 0, TestDataFactory.DEFAULT_SEED);
SimpleGroupFactory factory = new SimpleGroupFactory(schema);
rows = new ArrayList<>(ROW_COUNT);
for (int i = 0; i < ROW_COUNT; i++) {
Group row = factory.newGroup();
switch (data) {
case "mixed":
row.append("l", (long) i).append("b", binaries[i]);
Group g = row.addGroup("g");
for (int j = 0; j < 8; j++) {
g.append("i", random.nextInt());
}
break;
case "numbers":
row.append("i", random.nextInt()).append("l", random.nextLong());
if (random.nextInt(10) > 0) {
row.append("d", random.nextDouble());
}
break;
case "strings":
row.append("s", "s" + random.nextInt(1_000_000));
if (random.nextInt(10) > 0) {
row.append("t", Long.toString(random.nextLong(), 36));
}
break;
default:
if (random.nextInt(10) > 0) {
Group list = row.addGroup("l");
for (int n = random.nextInt(9); n > 0; n--) {
Group element = list.addGroup("list");
if (random.nextInt(10) > 0) {
element.append("element", random.nextLong());
}
}
}
}
rows.add(row);
}
}

@Benchmark
public void write() throws IOException {
try (ParquetWriter<Group> writer = ExampleParquetWriter.builder(BlackHoleOutputFile.INSTANCE)
.withWriteMode(ParquetFileWriter.Mode.OVERWRITE)
.withType(schema)
.withCompressionCodec(CompressionCodecName.UNCOMPRESSED)
.withPageRowCountLimit(Integer.MAX_VALUE)
.withDictionaryEncoding(false)
.withContentDefinedChunkingEnabled(chunking)
.build()) {
for (Group row : rows) {
writer.write(row);
}
}
}
}
157 changes: 157 additions & 0 deletions parquet-column/src/main/java/org/apache/parquet/column/CdcOptions.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,157 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one
* or more contributor license agreements. See the NOTICE file
* distributed with this work for additional information
* regarding copyright ownership. The ASF licenses this file
* to you under the Apache License, Version 2.0 (the
* "License"); you may not use this file except in compliance
* with the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing,
* software distributed under the License is distributed on an
* "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
* KIND, either express or implied. See the License for the
* specific language governing permissions and limitations
* under the License.
*/
package org.apache.parquet.column;

import org.apache.parquet.Preconditions;
import org.apache.parquet.internal.column.chunking.RollingHashMask;

/**
* EXPERIMENTAL: The size envelope and normalization level of content defined chunking (CDC), which
* ends data pages at boundaries derived from the column's values, so that files sharing a run of
* values share byte-identical pages. Sizes are measured on values and levels before encoding.
* Dictionary ids follow the order values first appear in, so an edit that adds new values renumbers
* the later ones within its row group: columns of many distinct values deduplicate best without
* dictionary encoding.
*
* <p>The settings and defaults match {@code CdcOptions} in Arrow C++ and arrow-rs, and so do the
* chunk boundaries for the same physical values. Arrow types converted before writing, such as
* narrow integers, coerced timestamps and decimals, are hashed differently there. Chunk boundaries
* continue across row groups, as in arrow-rs; Arrow C++ restarts them in every row group, so its
* chunks match only in the first.
*
* @see ParquetProperties.Builder#withContentDefinedChunking(CdcOptions)
*/
public final class CdcOptions {

/** 256 KiB minimum, 1 MiB maximum and normalization level 0. */
public static final CdcOptions DEFAULT = builder().build();

private final long minChunkSize;
private final long maxChunkSize;
private final int normLevel;

private CdcOptions(Builder builder) {
this.minChunkSize = builder.minChunkSize;
this.maxChunkSize = builder.maxChunkSize;
this.normLevel = builder.normLevel;
// Validate now rather than when the first column writer is built.
RollingHashMask.calculate(minChunkSize, maxChunkSize, normLevel);
}

/**
* @return the minimum chunk size in bytes
*/
public long getMinChunkSize() {
return minChunkSize;
}

/**
* @return the maximum chunk size in bytes
*/
public long getMaxChunkSize() {
return maxChunkSize;
}

/**
* @return the normalization level of the rolling hash mask
*/
public int getNormLevel() {
return normLevel;
}

@Override
public String toString() {
return "CdcOptions{minChunkSize=" + minChunkSize + ", maxChunkSize=" + maxChunkSize + ", normLevel=" + normLevel
+ '}';
}

/**
* @return a builder holding the default options
*/
public static Builder builder() {
return new Builder();
}

/** EXPERIMENTAL: Builds {@link CdcOptions}. */
public static class Builder {
private long minChunkSize = 256 * 1024L;
private long maxChunkSize = 1024 * 1024L;
private int normLevel = 0;

private Builder() {}

/**
* Set the minimum chunk size in bytes, 256 KiB by default. The rolling hash is not updated
* until a chunk reaches this size, so no chunk is shorter but a file's last; pages can be, where
* a row group or a page limit ends one.
*
* @param minChunkSize the minimum chunk size in bytes
* @return this builder for method chaining
*/
public Builder withMinChunkSize(long minChunkSize) {
Preconditions.checkArgument(
minChunkSize >= 0,
"Invalid content defined chunking minimum chunk size (negative): %s",
minChunkSize);
this.minChunkSize = minChunkSize;
return this;
}

/**
* Set the maximum chunk size in bytes, 1 MiB by default. A chunk ends when it reaches this size,
* whatever the rolling hash says. {@link ParquetProperties.Builder#withPageSize(int)} separately
* limits the page size; below this it splits chunks into more pages, in the same places after
* an edit, so it does not cost deduplication.
*
* @param maxChunkSize the maximum chunk size in bytes
* @return this builder for method chaining
*/
public Builder withMaxChunkSize(long maxChunkSize) {
Preconditions.checkArgument(
maxChunkSize > 0,
"Invalid content defined chunking maximum chunk size (not positive): %s",
maxChunkSize);
this.maxChunkSize = maxChunkSize;
return this;
}

/**
* Set the normalization level of the rolling hash mask, 0 by default. Raising it makes a
* boundary more likely, which tightens the chunk size distribution and improves deduplication
* at the cost of more small pages; lowering it does the reverse. Values outside
* {@code [-3, 3]} are not useful.
*
* @param normLevel the normalization level
* @return this builder for method chaining
*/
public Builder withNormLevel(int normLevel) {
this.normLevel = normLevel;
return this;
}

/**
* @return the options
* @throws IllegalArgumentException if the maximum chunk size is not greater than the minimum,
* or the envelope is too narrow for the normalization level
*/
public CdcOptions build() {
return new CdcOptions(this);
}
}
}
Loading
Loading