Skip to content
Open
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
3 changes: 1 addition & 2 deletions src/Parquet.PerfRunner/Benchmarks/DataTypes.cs
Original file line number Diff line number Diff line change
@@ -1,4 +1,3 @@
using Parquet.Data;
using Parquet.Schema;

namespace Parquet.PerfRunner.Benchmarks;
Expand Down Expand Up @@ -94,7 +93,7 @@ public async Task NullableInts() {
public async Task RandomStrings() {
using var ms = new MemoryStream();
await using ParquetWriter w = await ParquetWriter.CreateAsync(_nullableStringSchema, ms, new ParquetOptions { CompressionMethod = CompressionMethod.None });
using ParquetRowGroupWriter rgw = w.CreateRowGroup();
using ParquetRowGroupWriter rgw = w.CreateRowGroup();
await rgw.WriteAsync(_nullableStringSchema.DataFields[0], _nullableStrings);
}
}
111 changes: 111 additions & 0 deletions src/Parquet.PerfRunner/Benchmarks/ParquetSharpComparisonBenchmark.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,111 @@
using BenchmarkDotNet.Attributes;
using Parquet.PerfRunner.Taxis;
using ParquetSharp;
using ParquetSharp.IO;
using IOFile = System.IO.File;
using ParquetWriterNet = Parquet.ParquetWriter;

namespace Parquet.PerfRunner.Benchmarks;

[MemoryDiagnoser]
[MarkdownExporter]
[ShortRunJob]
public class ParquetSharpComparisonBenchmark {

[Params("tripdata", "tripdata-large")]
public string Dataset { get; set; } = "tripdata";

[Params("small", "full")]
public string Schema { get; set; } = "small";

[Params(LogicalEncoding.Plain, LogicalEncoding.RleDictionary, LogicalEncoding.DeltaBinaryPacked)]
public LogicalEncoding LogicalEncoding { get; set; }


TaxiSchema _parquetNetSchema = null!;
ParquetSharpTaxiSchema _parquetSharpSchema = null!;
ParquetOptions _parquetNetOptions = null!;
WriterProperties _parquetSharpOptions = null!;
readonly MemoryStream _memoryStream = new(2_000_000_000); // 2GB capacity so it wont cause allocations
TaxiDataset _dataset;

string GetFileName(string libName) =>
$"taxi-{Dataset}-schema_{Schema}-{libName}-{LogicalEncoding.ToString().ToLowerInvariant()}.parquet";

[GlobalSetup]
public async Task LoadDatasetAsync() {
_dataset = await TaxiDatasetLoader.Instance.LoadAsync(Dataset);
if(Schema == "small") {
_parquetNetSchema = TaxiSchema.Small(_dataset);
_parquetSharpSchema = ParquetSharpTaxiSchema.Small();
} else {
_parquetNetSchema = TaxiSchema.Full(_dataset);
_parquetSharpSchema = ParquetSharpTaxiSchema.Full();
}
_parquetNetOptions = LogicalEncoding.CreateOptions(_parquetNetSchema.Schema);
_parquetSharpOptions = _parquetSharpSchema.CreateParquetSharpWriterProperties(LogicalEncoding);
}

[Benchmark(Description = "Parquet.Net -> MemoryStream")]
public async Task ParquetNetAsync() {
_memoryStream.Position = 0;
#if PARQUET_PACKAGE
using ParquetWriterNet writer = await ParquetWriterNet.CreateAsync(_parquetNetSchema.Schema, _memoryStream, _parquetNetOptions);
writer.CompressionMethod = CompressionMethod.None;
#else
await using ParquetWriterNet writer = await ParquetWriterNet.CreateAsync(_parquetNetSchema.Schema, _memoryStream, _parquetNetOptions);
#endif
using ParquetRowGroupWriter rowGroup = writer.CreateRowGroup();

await _parquetNetSchema.WriteAsync(rowGroup);
}

[Benchmark(Description = "Parquet.Net -> Disk")]
public async Task ParquetNetToDiskAsync() {
string path = Path.Combine(Path.GetTempPath(), GetFileName("parquetnet"));

try {
await using FileStream output = IOFile.Create(path);
#if PARQUET_PACKAGE
using ParquetWriterNet writer = await ParquetWriterNet.CreateAsync(_parquetNetSchema.Schema, output, _parquetNetOptions);
writer.CompressionMethod = CompressionMethod.None;
#else
await using ParquetWriterNet writer = await ParquetWriterNet.CreateAsync(_parquetNetSchema.Schema, output, _parquetNetOptions);
#endif
using ParquetRowGroupWriter rowGroup = writer.CreateRowGroup();

await _parquetNetSchema.WriteAsync(rowGroup);
} finally {
IOFile.Delete(path);
}
}

[Benchmark(Description = "ParquetSharp -> MemoryStream")]
public void ParquetSharp() {
_memoryStream.Position = 0;
using var managedOutput = new ManagedOutputStream(_memoryStream, leaveOpen: true);
using var writer = new ParquetFileWriter(managedOutput, _parquetSharpSchema.Columns, _parquetSharpOptions);
using RowGroupWriter rowGroup = writer.AppendRowGroup();

_parquetSharpSchema.Write(rowGroup, _dataset);
writer.Close();
}

[Benchmark(Description = "ParquetSharp -> Disk")]
public void ParquetSharpToDisk() {
string path = Path.Combine(Path.GetTempPath(), GetFileName("parquetsharp"));

try {
using FileStream output = IOFile.Create(path);
using var managedOutput = new ManagedOutputStream(output);
using var writer = new ParquetFileWriter(managedOutput, _parquetSharpSchema.Columns, _parquetSharpOptions);
using RowGroupWriter rowGroup = writer.AppendRowGroup();

_parquetSharpSchema.Write(rowGroup, _dataset);

writer.Close();
} finally {
IOFile.Delete(path);
}
}
}
52 changes: 52 additions & 0 deletions src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
using BenchmarkDotNet.Attributes;
using BenchmarkDotNet.Configs;
using BenchmarkDotNet.Jobs;
using Parquet.PerfRunner.Taxis;
using ParquetWriterNet = Parquet.ParquetWriter;

namespace Parquet.PerfRunner.Benchmarks;

[Config(typeof(NuConfig))]
[MemoryDiagnoser]
[MarkdownExporter]
public class SelfComparisonBenchmark {
public class NuConfig : ManualConfig {
public NuConfig() {
Job baseJob = Job.ShortRun.WithId("Parquet.Net (local)");
AddJob(baseJob);
AddJob(baseJob
.WithId("Parquet.Net 5.6.0")
.WithMsBuildArguments("/p:ParquetVersion=5.6.0"));
}
}

[Params("tripdata", "tripdata-large")]
public string Dataset { get; set; } = "tripdata";

[Params(LogicalEncoding.Plain, LogicalEncoding.RleDictionary, LogicalEncoding.DeltaBinaryPacked)]
public LogicalEncoding LogicalEncoding { get; set; }

TaxiSchema _schema = null!;
TaxiDataset _dataset;
ParquetOptions _options = null!;
[GlobalSetup]
public async Task LoadDatasetAsync() {
_dataset = await TaxiDatasetLoader.Instance.LoadAsync(Dataset);
_schema = TaxiSchema.Full(_dataset);
_options = LogicalEncoding.CreateOptions(_schema.Schema);
}

[Benchmark(Description = "Parquet.Net source")]
public async Task ParquetNetAsync() {
using var output = new MemoryStream();
#if PARQUET_PACKAGE
using ParquetWriterNet writer = await ParquetWriterNet.CreateAsync(_schema.Schema, output, _options);
writer.CompressionMethod = CompressionMethod.None;
#else
await using ParquetWriterNet writer = await ParquetWriterNet.CreateAsync(_schema.Schema, output, _options);
#endif
using ParquetRowGroupWriter rowGroup = writer.CreateRowGroup();

await _schema.WriteAsync(rowGroup);
}
}
2 changes: 1 addition & 1 deletion src/Parquet.PerfRunner/Benchmarks/WriteBenchmark.cs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
using BenchmarkDotNet.Attributes;
using BenchmarkDotNet.Attributes;
using Parquet.Schema;
using ParquetSharp;
using ParquetSharp.IO;
Expand Down
68 changes: 68 additions & 0 deletions src/Parquet.PerfRunner/LogicalEncoding.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,68 @@
using Parquet.Schema;

namespace Parquet.PerfRunner;

/// <summary>
/// Allow to easily control the logical encoding used in benchmarks.
/// </summary>
public enum LogicalEncoding {
Plain,
RleDictionary,
DeltaBinaryPacked
}

public static class LogicalEncodingExtensions {
#if PARQUET_PACKAGE
public static ParquetOptions CreateOptions(this LogicalEncoding encoding, ParquetSchema? schema = null) =>
encoding switch {
LogicalEncoding.RleDictionary => new ParquetOptions {
UseDictionaryEncoding = true,
// Force dictionary extraction for every supported column to match the local benchmark setup.
DictionaryEncodingThreshold = double.MaxValue,
UseDeltaBinaryPackedEncoding = false
},
LogicalEncoding.DeltaBinaryPacked => new ParquetOptions {
UseDictionaryEncoding = false,
UseDeltaBinaryPackedEncoding = true
},
LogicalEncoding.Plain => new ParquetOptions {
UseDictionaryEncoding = false,
UseDeltaBinaryPackedEncoding = false
},
_ => throw new ArgumentOutOfRangeException(nameof(encoding), encoding, "Unknown logical encoding")
};
#else
public static ParquetOptions CreateOptions(this LogicalEncoding encoding, ParquetSchema? schema = null) {
ParquetOptions options = encoding switch {
LogicalEncoding.RleDictionary => new ParquetOptions {
CompressionMethod = CompressionMethod.None,
// Force dictionary extraction for every supported column to match the ParquetSharp benchmark setup.
DictionaryEncodingThreshold = double.MaxValue
},
LogicalEncoding.DeltaBinaryPacked => new ParquetOptions {
CompressionMethod = CompressionMethod.None
},
LogicalEncoding.Plain => new ParquetOptions {
CompressionMethod = CompressionMethod.None,
DictionaryEncodingThreshold = -1
},
_ => throw new ArgumentOutOfRangeException(nameof(encoding), encoding, "Unknown logical encoding")
};

if(schema != null) {
EncodingHint hint = encoding switch {
LogicalEncoding.RleDictionary => EncodingHint.Dictionary,
LogicalEncoding.DeltaBinaryPacked => EncodingHint.DeltaBinaryPacked,
LogicalEncoding.Plain => EncodingHint.Default,
_ => throw new ArgumentOutOfRangeException(nameof(encoding), encoding, "Unknown logical encoding")
};

foreach(DataField field in schema.GetDataFields()) {
options.ColumnEncodingHints[field.Path.ToString()] = hint;
}
}

return options;
}
#endif
}
16 changes: 15 additions & 1 deletion src/Parquet.PerfRunner/Parquet.PerfRunner.csproj
Original file line number Diff line number Diff line change
Expand Up @@ -27,4 +27,18 @@
<PackageReference Include="Parquet.Net" Version="$(ParquetVersion)" />
</ItemGroup>

</Project>
<PropertyGroup Condition="'$(ParquetVersion)' != ''">
<DefineConstants>$(DefineConstants);PARQUET_PACKAGE</DefineConstants>
</PropertyGroup>

<ItemGroup Condition="'$(ParquetVersion)' != ''">
<Compile Remove="SampleGenerator.cs" />
<Compile Remove="Benchmarks\CompressionBenchmarks.cs" />
<Compile Remove="Benchmarks\Curiosities.cs" />
<Compile Remove="Benchmarks\DataTypes.cs" />
<Compile Remove="Benchmarks\EncodingBenchmarks.cs" />
<Compile Remove="Benchmarks\HighLevel.cs" />
<Compile Remove="Benchmarks\WriteBenchmark.cs" />
</ItemGroup>

</Project>
14 changes: 11 additions & 3 deletions src/Parquet.PerfRunner/Program.cs
Original file line number Diff line number Diff line change
@@ -1,12 +1,11 @@
// for performance tests only
// for performance tests only

using BenchmarkDotNet.Running;
using Parquet;
using Parquet.PerfRunner;
using Parquet.PerfRunner.Benchmarks;

if(args.Length == 1) {
switch(args[0]) {
#if !PARQUET_PACKAGE
case "highLevel":
HighLevel.Run();
break;
Expand All @@ -24,8 +23,17 @@
case "compression":
BenchmarkRunner.Run<CompressionBenchmarks>();
break;
#endif
case "sharp-compare-taxi":
BenchmarkRunner.Run<ParquetSharpComparisonBenchmark>();
break;
case "self-compare-taxi":
BenchmarkRunner.Run<SelfComparisonBenchmark>();
break;
}
} else {
#if !PARQUET_PACKAGE
await new DataTypes().RandomStrings();
#endif
//await SampleGenerator.GenerateFiles();
}
Loading
Loading