From ba7b0549546da60d871f5a8477570ce92c77bfd1 Mon Sep 17 00:00:00 2001 From: Kuinox Date: Tue, 2 Dec 2025 01:45:10 +0100 Subject: [PATCH 1/6] Add benchmark. --- .../Benchmarks/TaxiCsvToParquetBenchmark.cs | 369 ++++++++++++++++++ src/Parquet.PerfRunner/Program.cs | 3 + 2 files changed, 372 insertions(+) create mode 100644 src/Parquet.PerfRunner/Benchmarks/TaxiCsvToParquetBenchmark.cs diff --git a/src/Parquet.PerfRunner/Benchmarks/TaxiCsvToParquetBenchmark.cs b/src/Parquet.PerfRunner/Benchmarks/TaxiCsvToParquetBenchmark.cs new file mode 100644 index 00000000..91bff66e --- /dev/null +++ b/src/Parquet.PerfRunner/Benchmarks/TaxiCsvToParquetBenchmark.cs @@ -0,0 +1,369 @@ +using System.Collections.Generic; +using System.Linq; +using System.Net.Http; +using BenchmarkDotNet.Attributes; +using Parquet.Data; +using Parquet.Schema; +using ParquetSharp; +using ParquetSharp.IO; +using Column = ParquetSharp.Column; +using ParquetReaderNet = Parquet.ParquetReader; +using ParquetWriterNet = Parquet.ParquetWriter; +using IOFile = System.IO.File; + +namespace Parquet.PerfRunner.Benchmarks; + +[MemoryDiagnoser] +[MarkdownExporter] +[ShortRunJob] +public class TaxiCsvToParquetBenchmark +{ + // Official NYC TLC parquet. The "tripdata" dataset is a single month; "tripdata-large" combines multiple months. + private const string BaseUrl = "https://d37ci6vzurychx.cloudfront.net/trip-data/"; + private static readonly IReadOnlyDictionary DatasetFiles = new Dictionary + { + { "tripdata", new[] { "yellow_tripdata_2024-01.parquet" } }, + { "tripdata-large", new[] { "yellow_tripdata_2024-01.parquet", "yellow_tripdata_2024-02.parquet", "yellow_tripdata_2024-03.parquet" } } + }; + + private int?[]? _vendorIds; + private int?[]? _rateCodes; + private double?[]? _passengerCounts; + private double?[]? _tripDistances; + private int?[]? _paymentTypes; + private double?[]? _fareAmounts; + + private ParquetSchema? _parquetSchema; + private DataColumn[]? _parquetNetColumns; + private Column[]? _parquetSharpColumns; + + [Params("tripdata", "tripdata-large")] + public string Dataset { get; set; } = "tripdata"; + + [GlobalSetup] + public async Task LoadDataset() + { + string[] parquetPaths = await EnsureDatasetsAsync(); + TaxiColumns columns = await LoadParquetColumnsAsync(parquetPaths); + + _vendorIds = columns.VendorIds; + _rateCodes = columns.RateCodes; + _passengerCounts = columns.PassengerCounts; + _tripDistances = columns.TripDistances; + _paymentTypes = columns.PaymentTypes; + _fareAmounts = columns.FareAmounts; + + _parquetSchema = new ParquetSchema( + new DataField("vendorid"), + new DataField("ratecodeid"), + new DataField("passenger_count"), + new DataField("trip_distance"), + new DataField("payment_type"), + new DataField("fare_amount")); + + DataField[] fields = _parquetSchema.DataFields; + _parquetNetColumns = + [ + new DataColumn(fields[0], _vendorIds), + new DataColumn(fields[1], _rateCodes), + new DataColumn(fields[2], _passengerCounts), + new DataColumn(fields[3], _tripDistances), + new DataColumn(fields[4], _paymentTypes), + new DataColumn(fields[5], _fareAmounts) + ]; + + _parquetSharpColumns = + [ + new Column("vendorid"), + new Column("ratecodeid"), + new Column("passenger_count"), + new Column("trip_distance"), + new Column("payment_type"), + new Column("fare_amount") + ]; + } + + [Benchmark(Description = "Parquet.Net -> MemoryStream")] + public async Task ParquetNet() + { + using var output = new MemoryStream(); + using ParquetWriterNet writer = await ParquetWriterNet.CreateAsync(_parquetSchema!, output); + using ParquetRowGroupWriter rowGroup = writer.CreateRowGroup(); + + foreach (DataColumn column in _parquetNetColumns!) + { + await rowGroup.WriteColumnAsync(column); + } + } + + [Benchmark(Description = "ParquetSharp -> MemoryStream")] + public void ParquetSharp() + { + using var output = new MemoryStream(); + using var managedOutput = new ManagedOutputStream(output); + using var writer = new ParquetFileWriter(managedOutput, _parquetSharpColumns!); + using RowGroupWriter rowGroup = writer.AppendRowGroup(); + + using (LogicalColumnWriter vendorWriter = rowGroup.NextColumn().LogicalWriter()) + { + vendorWriter.WriteBatch(_vendorIds!); + } + + using (LogicalColumnWriter rateCodeWriter = rowGroup.NextColumn().LogicalWriter()) + { + rateCodeWriter.WriteBatch(_rateCodes!); + } + + using (LogicalColumnWriter passengerCountWriter = rowGroup.NextColumn().LogicalWriter()) + { + passengerCountWriter.WriteBatch(_passengerCounts!); + } + + using (LogicalColumnWriter tripDistanceWriter = rowGroup.NextColumn().LogicalWriter()) + { + tripDistanceWriter.WriteBatch(_tripDistances!); + } + + using (LogicalColumnWriter paymentTypeWriter = rowGroup.NextColumn().LogicalWriter()) + { + paymentTypeWriter.WriteBatch(_paymentTypes!); + } + + using (LogicalColumnWriter fareWriter = rowGroup.NextColumn().LogicalWriter()) + { + fareWriter.WriteBatch(_fareAmounts!); + } + + writer.Close(); + } + + [Benchmark(Description = "Parquet.Net -> Disk")] + public async Task ParquetNetToDisk() + { + string path = Path.Combine(Path.GetTempPath(), $"taxi-parquetnet-{Guid.NewGuid():N}.parquet"); + try + { + await using FileStream output = IOFile.Create(path); + using ParquetWriterNet writer = await ParquetWriterNet.CreateAsync(_parquetSchema!, output); + using ParquetRowGroupWriter rowGroup = writer.CreateRowGroup(); + + foreach (DataColumn column in _parquetNetColumns!) + { + await rowGroup.WriteColumnAsync(column); + } + } + finally + { + TryDelete(path); + } + } + + [Benchmark(Description = "ParquetSharp -> Disk")] + public void ParquetSharpToDisk() + { + string path = Path.Combine(Path.GetTempPath(), $"taxi-parquetsharp-{Guid.NewGuid():N}.parquet"); + try + { + using FileStream output = IOFile.Create(path); + using var managedOutput = new ManagedOutputStream(output); + using var writer = new ParquetFileWriter(managedOutput, _parquetSharpColumns!); + using RowGroupWriter rowGroup = writer.AppendRowGroup(); + + using (LogicalColumnWriter vendorWriter = rowGroup.NextColumn().LogicalWriter()) + { + vendorWriter.WriteBatch(_vendorIds!); + } + + using (LogicalColumnWriter rateCodeWriter = rowGroup.NextColumn().LogicalWriter()) + { + rateCodeWriter.WriteBatch(_rateCodes!); + } + + using (LogicalColumnWriter passengerCountWriter = rowGroup.NextColumn().LogicalWriter()) + { + passengerCountWriter.WriteBatch(_passengerCounts!); + } + + using (LogicalColumnWriter tripDistanceWriter = rowGroup.NextColumn().LogicalWriter()) + { + tripDistanceWriter.WriteBatch(_tripDistances!); + } + + using (LogicalColumnWriter paymentTypeWriter = rowGroup.NextColumn().LogicalWriter()) + { + paymentTypeWriter.WriteBatch(_paymentTypes!); + } + + using (LogicalColumnWriter fareWriter = rowGroup.NextColumn().LogicalWriter()) + { + fareWriter.WriteBatch(_fareAmounts!); + } + + writer.Close(); + } + finally + { + TryDelete(path); + } + } + + private static string GetDataDirectory() + { + return Path.Combine(AppContext.BaseDirectory, "Data"); + } + + private static string? TryProjectDataPath(string fileName) + { + string projectPath = Path.GetFullPath(Path.Combine(AppContext.BaseDirectory, "../../../Data", fileName)); + return IOFile.Exists(projectPath) ? projectPath : null; + } + + private async Task EnsureDatasetsAsync() + { + if (!DatasetFiles.TryGetValue(Dataset, out string[]? files)) + { + throw new ArgumentOutOfRangeException(nameof(Dataset), Dataset, "Unknown dataset alias"); + } + + string dataDir = GetDataDirectory(); + Directory.CreateDirectory(dataDir); + + var paths = new List(files.Length); + using var httpClient = new HttpClient(); + + foreach (string file in files) + { + string dataPath = Path.Combine(dataDir, file); + if (IOFile.Exists(dataPath)) + { + paths.Add(dataPath); + continue; + } + + string? projectPath = TryProjectDataPath(file); + if (projectPath != null) + { + paths.Add(projectPath); + continue; + } + + string url = $"{BaseUrl}{file}"; + byte[] bytes = await httpClient.GetByteArrayAsync(url); + await IOFile.WriteAllBytesAsync(dataPath, bytes); + paths.Add(dataPath); + } + + return paths.ToArray(); + } + + private static async Task LoadParquetColumnsAsync(IEnumerable paths) + { + var vendorIds = new List(); + var rateCodes = new List(); + var passengerCounts = new List(); + var tripDistances = new List(); + var paymentTypes = new List(); + var fareAmounts = new List(); + + foreach (string path in paths) + { + using FileStream fs = IOFile.OpenRead(path); + using ParquetReaderNet reader = await ParquetReaderNet.CreateAsync(fs); + + DataField vendor = FindField(reader.Schema, "vendorid"); + DataField rateCode = FindField(reader.Schema, "ratecodeid"); + DataField passengerCount = FindField(reader.Schema, "passenger_count"); + DataField tripDistance = FindField(reader.Schema, "trip_distance"); + DataField paymentType = FindField(reader.Schema, "payment_type"); + DataField fareAmount = FindField(reader.Schema, "fare_amount"); + + for (int i = 0; i < reader.RowGroupCount; i++) + { + using ParquetRowGroupReader rg = reader.OpenRowGroupReader(i); + + vendorIds.AddRange(ToNullableInt((await rg.ReadColumnAsync(vendor)).Data)); + rateCodes.AddRange(ToNullableInt((await rg.ReadColumnAsync(rateCode)).Data)); + passengerCounts.AddRange(ToNullableDouble((await rg.ReadColumnAsync(passengerCount)).Data)); + tripDistances.AddRange(ToNullableDouble((await rg.ReadColumnAsync(tripDistance)).Data)); + paymentTypes.AddRange(ToNullableInt((await rg.ReadColumnAsync(paymentType)).Data)); + fareAmounts.AddRange(ToNullableDouble((await rg.ReadColumnAsync(fareAmount)).Data)); + } + } + + return new TaxiColumns( + vendorIds.ToArray(), + rateCodes.ToArray(), + passengerCounts.ToArray(), + tripDistances.ToArray(), + paymentTypes.ToArray(), + fareAmounts.ToArray()); + } + + private static DataField FindField(ParquetSchema schema, string name) + { + return schema.DataFields.First(f => string.Equals(f.Name, name, StringComparison.OrdinalIgnoreCase)); + } + + private static IReadOnlyList ToNullableInt(Array data) => + data switch + { + int?[] v => v, + int[] v => v.Select(i => (int?)i).ToArray(), + long?[] v => v.Select(i => i.HasValue ? (int?)checked((int)i.Value) : null).ToArray(), + long[] v => v.Select(i => (int?)checked((int)i)).ToArray(), + double?[] v => v.Select(i => i.HasValue ? (int?)checked((int)i.Value) : null).ToArray(), + double[] v => v.Select(i => (int?)checked((int)i)).ToArray(), + float?[] v => v.Select(i => i.HasValue ? (int?)checked((int)i.Value) : null).ToArray(), + float[] v => v.Select(i => (int?)checked((int)i)).ToArray(), + _ => throw new InvalidOperationException($"Unsupported numeric type: {data.GetType()}") + }; + + private static IReadOnlyList ToNullableDouble(Array data) => + data switch + { + double?[] v => v, + double[] v => v.Select(i => (double?)i).ToArray(), + float?[] v => v.Select(i => i.HasValue ? (double?)i.Value : null).ToArray(), + float[] v => v.Select(i => (double?)i).ToArray(), + int?[] v => v.Select(i => i.HasValue ? (double?)i.Value : null).ToArray(), + int[] v => v.Select(i => (double?)i).ToArray(), + long?[] v => v.Select(i => i.HasValue ? (double?)i.Value : null).ToArray(), + long[] v => v.Select(i => (double?)i).ToArray(), + _ => throw new InvalidOperationException($"Unsupported numeric type: {data.GetType()}") + }; + + private static void TryDelete(string path) + { + try + { + if (IOFile.Exists(path)) + { + IOFile.Delete(path); + } + } + catch + { + // best effort cleanup + } + } +} + +internal readonly struct TaxiColumns +{ + public TaxiColumns(int?[] vendorIds, int?[] rateCodes, double?[] passengerCounts, double?[] tripDistances, int?[] paymentTypes, double?[] fareAmounts) + { + VendorIds = vendorIds; + RateCodes = rateCodes; + PassengerCounts = passengerCounts; + TripDistances = tripDistances; + PaymentTypes = paymentTypes; + FareAmounts = fareAmounts; + } + + public int?[] VendorIds { get; } + public int?[] RateCodes { get; } + public double?[] PassengerCounts { get; } + public double?[] TripDistances { get; } + public int?[] PaymentTypes { get; } + public double?[] FareAmounts { get; } +} diff --git a/src/Parquet.PerfRunner/Program.cs b/src/Parquet.PerfRunner/Program.cs index f8ca501c..2778584d 100644 --- a/src/Parquet.PerfRunner/Program.cs +++ b/src/Parquet.PerfRunner/Program.cs @@ -12,6 +12,9 @@ case "progression": VersionedBenchmark.Run(); break; + case "taxi": + BenchmarkRunner.Run(); + break; } } else { await new DataTypes().NullableInts(); From a6741745a76d72a0d5bd1ba70bd7f50b52651713 Mon Sep 17 00:00:00 2001 From: Kuinox Date: Sun, 14 Dec 2025 03:05:10 +0100 Subject: [PATCH 2/6] Refactor benchmark. --- .../Benchmarks/DataTypes.cs | 4 +- .../ParquetSharpComparisonBenchmark.cs | 113 ++++++ .../Benchmarks/Progression.cs | 66 +++- .../Benchmarks/SelfComparisonBenchmark.cs | 42 ++ .../Benchmarks/TaxiCsvToParquetBenchmark.cs | 369 ------------------ .../Benchmarks/WriteBenchmark.cs | 6 +- src/Parquet.PerfRunner/LogicalEncoding.cs | 37 ++ .../Parquet.PerfRunner.csproj | 8 +- src/Parquet.PerfRunner/Program.cs | 35 +- .../Taxis/ParquetSharpTaxiSchema.cs | 76 ++++ src/Parquet.PerfRunner/Taxis/TaxiDataset.cs | 24 ++ .../Taxis/TaxiDatasetLoader.cs | 148 +++++++ src/Parquet.PerfRunner/Taxis/TaxiSchema.cs | 115 ++++++ .../Taxis/TaxiSchemaKind.cs | 6 + 14 files changed, 647 insertions(+), 402 deletions(-) create mode 100644 src/Parquet.PerfRunner/Benchmarks/ParquetSharpComparisonBenchmark.cs create mode 100644 src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs delete mode 100644 src/Parquet.PerfRunner/Benchmarks/TaxiCsvToParquetBenchmark.cs create mode 100644 src/Parquet.PerfRunner/LogicalEncoding.cs create mode 100644 src/Parquet.PerfRunner/Taxis/ParquetSharpTaxiSchema.cs create mode 100644 src/Parquet.PerfRunner/Taxis/TaxiDataset.cs create mode 100644 src/Parquet.PerfRunner/Taxis/TaxiDatasetLoader.cs create mode 100644 src/Parquet.PerfRunner/Taxis/TaxiSchema.cs create mode 100644 src/Parquet.PerfRunner/Taxis/TaxiSchemaKind.cs diff --git a/src/Parquet.PerfRunner/Benchmarks/DataTypes.cs b/src/Parquet.PerfRunner/Benchmarks/DataTypes.cs index 1cbe4324..7c215500 100644 --- a/src/Parquet.PerfRunner/Benchmarks/DataTypes.cs +++ b/src/Parquet.PerfRunner/Benchmarks/DataTypes.cs @@ -28,13 +28,13 @@ public static string RandomString(int length) { public DataTypes() { //_ints = new DataColumn(new DataField("c"), Enumerable.Range(0, DataSize).ToArray()); - _nullableInts = new DataColumn(_nullableIntsSchema.DataFields[0], + _nullableInts = new DataColumn(_nullableIntsSchema.GetDataFields()[0], Enumerable .Range(0, DataSize) .Select(i => i % 4 == 0 ? (int?)null : i) .ToArray()); - _nullableDecimals = new DataColumn(_nullableDecimalsSchema.DataFields[0], + _nullableDecimals = new DataColumn(_nullableDecimalsSchema.GetDataFields()[0], Enumerable .Range(0, DataSize) .Select(i => i % 4 == 0 ? (decimal?)null : (decimal)i) diff --git a/src/Parquet.PerfRunner/Benchmarks/ParquetSharpComparisonBenchmark.cs b/src/Parquet.PerfRunner/Benchmarks/ParquetSharpComparisonBenchmark.cs new file mode 100644 index 00000000..d8b1241c --- /dev/null +++ b/src/Parquet.PerfRunner/Benchmarks/ParquetSharpComparisonBenchmark.cs @@ -0,0 +1,113 @@ +using BenchmarkDotNet.Attributes; +using Parquet.Data; +using Parquet.Meta; +using Parquet.PerfRunner.Taxis; +using ParquetSharp; +using ParquetSharp.IO; +using IOFile = System.IO.File; +using ParquetSharpEncoding = ParquetSharp.Encoding; +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(); + _parquetSharpOptions = _parquetSharpSchema.CreateParquetSharpWriterProperties(LogicalEncoding); + } + + [Benchmark(Description = "Parquet.Net -> MemoryStream")] + public async Task ParquetNetAsync() { + _memoryStream.Position = 0; + using ParquetWriterNet writer = await ParquetWriterNet.CreateAsync(_parquetNetSchema.Schema, _memoryStream, _parquetNetOptions); + writer.CompressionMethod = CompressionMethod.Snappy; + using ParquetRowGroupWriter rowGroup = writer.CreateRowGroup(); + + foreach(DataColumn column in _parquetNetSchema.Columns) { + await rowGroup.WriteColumnAsync(column); + } + } + + [Benchmark(Description = "Parquet.Net -> Disk")] + public async Task ParquetNetToDiskAsync() { + string path = Path.Combine(Path.GetTempPath(), GetFileName("parquetnet")); + if(Path.Exists(path)) + IOFile.Delete(path); + + try { + await using FileStream output = IOFile.Create(path); + using ParquetWriterNet writer = await ParquetWriterNet.CreateAsync(_parquetNetSchema.Schema, output, _parquetNetOptions); + writer.CompressionMethod = CompressionMethod.Snappy; + using ParquetRowGroupWriter rowGroup = writer.CreateRowGroup(); + + foreach(DataColumn column in _parquetNetSchema.Columns) { + await rowGroup.WriteColumnAsync(column); + } + } 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(); + + _parquetNetSchema.WriteParquetSharp(rowGroup, _dataset); + writer.Close(); + } + + [Benchmark(Description = "ParquetSharp -> Disk")] + public void ParquetSharpToDisk() { + string path = Path.Combine(Path.GetTempPath(), GetFileName("parquetsharp")); + if(Path.Exists(path)) + IOFile.Delete(path); + 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(); + + _parquetNetSchema.WriteParquetSharp(rowGroup, _dataset); + + writer.Close(); + } finally { + //IOFile.Delete(path); + } + } +} diff --git a/src/Parquet.PerfRunner/Benchmarks/Progression.cs b/src/Parquet.PerfRunner/Benchmarks/Progression.cs index e6bf9f82..8a8b7940 100644 --- a/src/Parquet.PerfRunner/Benchmarks/Progression.cs +++ b/src/Parquet.PerfRunner/Benchmarks/Progression.cs @@ -7,6 +7,7 @@ using BenchmarkDotNet.Configs; using BenchmarkDotNet.Jobs; using BenchmarkDotNet.Running; +using Microsoft.Diagnostics.Tracing.Parsers.Kernel; using Parquet.Data; using Parquet.Schema; @@ -18,22 +19,55 @@ namespace Parquet.PerfRunner.Benchmarks { [MemoryDiagnoser] [RPlotExporter] public class VersionedBenchmark { - - public static void Run() { - BenchmarkRunner.Run(); - } - public class NuConfig : ManualConfig { public NuConfig() { Job baseJob = Job.ShortRun; - //AddJob(baseJob.WithNuGet("Parquet.Net", "4.2.3")); - //AddJob(baseJob.WithNuGet("Parquet.Net", "4.3.0")); - //AddJob(baseJob.WithNuGet("Parquet.Net", "4.3.2")); - //AddJob(baseJob.WithNuGet("Parquet.Net", "4.4.1")); - AddJob(baseJob.WithNuGet("Parquet.Net", "4.5.0")); - AddJob(baseJob.WithNuGet("Parquet.Net", "4.9.1")); - AddJob(baseJob.WithNuGet("Parquet.Net", "4.12.0")); + //AddJob(CreatePackageJob(baseJob, "4.2.3")); + //AddJob(CreatePackageJob(baseJob, "4.3.0")); + //AddJob(CreatePackageJob(baseJob, "4.3.2")); + //AddJob(CreatePackageJob(baseJob, "4.4.1")); + AddJob(CreatePackageJob(baseJob, "4.5.0")); + AddJob(CreatePackageJob(baseJob, "4.9.1")); + AddJob(CreatePackageJob(baseJob, "4.12.0")); + AddJob(CreatePackageJob(baseJob, "4.13.0")); + AddJob(CreatePackageJob(baseJob, "4.14.0")); + AddJob(CreatePackageJob(baseJob, "4.15.0")); + AddJob(CreatePackageJob(baseJob, "4.16.0")); + AddJob(CreatePackageJob(baseJob, "4.16.1")); + AddJob(CreatePackageJob(baseJob, "4.16.2")); + AddJob(CreatePackageJob(baseJob, "4.16.3")); + AddJob(CreatePackageJob(baseJob, "4.16.4")); + AddJob(CreatePackageJob(baseJob, "4.17.0")); + AddJob(CreatePackageJob(baseJob, "4.18.0")); + AddJob(CreatePackageJob(baseJob, "4.18.1")); + AddJob(CreatePackageJob(baseJob, "4.19.0")); + AddJob(CreatePackageJob(baseJob, "4.20.0")); + AddJob(CreatePackageJob(baseJob, "4.20.1")); + AddJob(CreatePackageJob(baseJob, "4.22.0")); + AddJob(CreatePackageJob(baseJob, "4.22.1")); + AddJob(CreatePackageJob(baseJob, "4.23.0")); + AddJob(CreatePackageJob(baseJob, "4.23.1")); + AddJob(CreatePackageJob(baseJob, "4.23.2")); + AddJob(CreatePackageJob(baseJob, "4.23.3")); + AddJob(CreatePackageJob(baseJob, "4.23.4")); + AddJob(CreatePackageJob(baseJob, "4.23.5")); + AddJob(CreatePackageJob(baseJob, "4.24.0")); + AddJob(CreatePackageJob(baseJob, "4.25.0")); + AddJob(CreatePackageJob(baseJob, "5.0.0")); + AddJob(CreatePackageJob(baseJob, "5.0.1")); + AddJob(CreatePackageJob(baseJob, "5.0.2")); + AddJob(CreatePackageJob(baseJob, "5.1.0")); + AddJob(CreatePackageJob(baseJob, "5.1.1")); + AddJob(CreatePackageJob(baseJob, "5.2.0")); + AddJob(CreatePackageJob(baseJob, "5.3.0")); + AddJob(CreatePackageJob(baseJob, "5.4.0")); + } + + private static Job CreatePackageJob(Job baseJob, string version) { + return baseJob + .WithId($"Parquet.Net {version}") + .WithMsBuildArguments($"/p:ParquetNuGetVersion={version}"); } } @@ -53,8 +87,14 @@ public static string RandomString(int length) { #endregion + [Params(LogicalEncoding.Plain, LogicalEncoding.RleDictionary, LogicalEncoding.DeltaBinaryPacked)] + public LogicalEncoding LogicalEncoding { get; set; } + + ParquetOptions _options = null!; + [GlobalSetup] public async Task Setup() { + _options = LogicalEncoding.CreateOptions(); _ints = new DataColumn(_intsSchema.GetDataFields()[0], Enumerable.Range(0, DataSize).Select(i => i % 4 == 0 ? (int?)null : i).ToArray(), null); @@ -94,7 +134,7 @@ public async Task Setup() { private async Task MakeFile(ParquetSchema schema, DataColumn c) { var ms = new MemoryStream(); - using(ParquetWriter writer = await ParquetWriter.CreateAsync(schema, ms)) { + using(ParquetWriter writer = await ParquetWriter.CreateAsync(schema, ms, _options)) { writer.CompressionMethod = CompressionMethod.None; // create a new row group in the file using(ParquetRowGroupWriter groupWriter = writer.CreateRowGroup()) { diff --git a/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs b/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs new file mode 100644 index 00000000..213dd62e --- /dev/null +++ b/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs @@ -0,0 +1,42 @@ +using BenchmarkDotNet.Attributes; +using Parquet.Data; +using Parquet.Meta; +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 SelfComparisonBenchmark { + [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(); + } + + [Benchmark(Description = "Parquet.Net source")] + public async Task ParquetNetAsync() { + using var output = new MemoryStream(); + using ParquetWriterNet writer = await ParquetWriterNet.CreateAsync(_schema.Schema, output, _options); + using ParquetRowGroupWriter rowGroup = writer.CreateRowGroup(); + + foreach(DataColumn column in _schema.Columns) { + await rowGroup.WriteColumnAsync(column); + } + } +} diff --git a/src/Parquet.PerfRunner/Benchmarks/TaxiCsvToParquetBenchmark.cs b/src/Parquet.PerfRunner/Benchmarks/TaxiCsvToParquetBenchmark.cs deleted file mode 100644 index 91bff66e..00000000 --- a/src/Parquet.PerfRunner/Benchmarks/TaxiCsvToParquetBenchmark.cs +++ /dev/null @@ -1,369 +0,0 @@ -using System.Collections.Generic; -using System.Linq; -using System.Net.Http; -using BenchmarkDotNet.Attributes; -using Parquet.Data; -using Parquet.Schema; -using ParquetSharp; -using ParquetSharp.IO; -using Column = ParquetSharp.Column; -using ParquetReaderNet = Parquet.ParquetReader; -using ParquetWriterNet = Parquet.ParquetWriter; -using IOFile = System.IO.File; - -namespace Parquet.PerfRunner.Benchmarks; - -[MemoryDiagnoser] -[MarkdownExporter] -[ShortRunJob] -public class TaxiCsvToParquetBenchmark -{ - // Official NYC TLC parquet. The "tripdata" dataset is a single month; "tripdata-large" combines multiple months. - private const string BaseUrl = "https://d37ci6vzurychx.cloudfront.net/trip-data/"; - private static readonly IReadOnlyDictionary DatasetFiles = new Dictionary - { - { "tripdata", new[] { "yellow_tripdata_2024-01.parquet" } }, - { "tripdata-large", new[] { "yellow_tripdata_2024-01.parquet", "yellow_tripdata_2024-02.parquet", "yellow_tripdata_2024-03.parquet" } } - }; - - private int?[]? _vendorIds; - private int?[]? _rateCodes; - private double?[]? _passengerCounts; - private double?[]? _tripDistances; - private int?[]? _paymentTypes; - private double?[]? _fareAmounts; - - private ParquetSchema? _parquetSchema; - private DataColumn[]? _parquetNetColumns; - private Column[]? _parquetSharpColumns; - - [Params("tripdata", "tripdata-large")] - public string Dataset { get; set; } = "tripdata"; - - [GlobalSetup] - public async Task LoadDataset() - { - string[] parquetPaths = await EnsureDatasetsAsync(); - TaxiColumns columns = await LoadParquetColumnsAsync(parquetPaths); - - _vendorIds = columns.VendorIds; - _rateCodes = columns.RateCodes; - _passengerCounts = columns.PassengerCounts; - _tripDistances = columns.TripDistances; - _paymentTypes = columns.PaymentTypes; - _fareAmounts = columns.FareAmounts; - - _parquetSchema = new ParquetSchema( - new DataField("vendorid"), - new DataField("ratecodeid"), - new DataField("passenger_count"), - new DataField("trip_distance"), - new DataField("payment_type"), - new DataField("fare_amount")); - - DataField[] fields = _parquetSchema.DataFields; - _parquetNetColumns = - [ - new DataColumn(fields[0], _vendorIds), - new DataColumn(fields[1], _rateCodes), - new DataColumn(fields[2], _passengerCounts), - new DataColumn(fields[3], _tripDistances), - new DataColumn(fields[4], _paymentTypes), - new DataColumn(fields[5], _fareAmounts) - ]; - - _parquetSharpColumns = - [ - new Column("vendorid"), - new Column("ratecodeid"), - new Column("passenger_count"), - new Column("trip_distance"), - new Column("payment_type"), - new Column("fare_amount") - ]; - } - - [Benchmark(Description = "Parquet.Net -> MemoryStream")] - public async Task ParquetNet() - { - using var output = new MemoryStream(); - using ParquetWriterNet writer = await ParquetWriterNet.CreateAsync(_parquetSchema!, output); - using ParquetRowGroupWriter rowGroup = writer.CreateRowGroup(); - - foreach (DataColumn column in _parquetNetColumns!) - { - await rowGroup.WriteColumnAsync(column); - } - } - - [Benchmark(Description = "ParquetSharp -> MemoryStream")] - public void ParquetSharp() - { - using var output = new MemoryStream(); - using var managedOutput = new ManagedOutputStream(output); - using var writer = new ParquetFileWriter(managedOutput, _parquetSharpColumns!); - using RowGroupWriter rowGroup = writer.AppendRowGroup(); - - using (LogicalColumnWriter vendorWriter = rowGroup.NextColumn().LogicalWriter()) - { - vendorWriter.WriteBatch(_vendorIds!); - } - - using (LogicalColumnWriter rateCodeWriter = rowGroup.NextColumn().LogicalWriter()) - { - rateCodeWriter.WriteBatch(_rateCodes!); - } - - using (LogicalColumnWriter passengerCountWriter = rowGroup.NextColumn().LogicalWriter()) - { - passengerCountWriter.WriteBatch(_passengerCounts!); - } - - using (LogicalColumnWriter tripDistanceWriter = rowGroup.NextColumn().LogicalWriter()) - { - tripDistanceWriter.WriteBatch(_tripDistances!); - } - - using (LogicalColumnWriter paymentTypeWriter = rowGroup.NextColumn().LogicalWriter()) - { - paymentTypeWriter.WriteBatch(_paymentTypes!); - } - - using (LogicalColumnWriter fareWriter = rowGroup.NextColumn().LogicalWriter()) - { - fareWriter.WriteBatch(_fareAmounts!); - } - - writer.Close(); - } - - [Benchmark(Description = "Parquet.Net -> Disk")] - public async Task ParquetNetToDisk() - { - string path = Path.Combine(Path.GetTempPath(), $"taxi-parquetnet-{Guid.NewGuid():N}.parquet"); - try - { - await using FileStream output = IOFile.Create(path); - using ParquetWriterNet writer = await ParquetWriterNet.CreateAsync(_parquetSchema!, output); - using ParquetRowGroupWriter rowGroup = writer.CreateRowGroup(); - - foreach (DataColumn column in _parquetNetColumns!) - { - await rowGroup.WriteColumnAsync(column); - } - } - finally - { - TryDelete(path); - } - } - - [Benchmark(Description = "ParquetSharp -> Disk")] - public void ParquetSharpToDisk() - { - string path = Path.Combine(Path.GetTempPath(), $"taxi-parquetsharp-{Guid.NewGuid():N}.parquet"); - try - { - using FileStream output = IOFile.Create(path); - using var managedOutput = new ManagedOutputStream(output); - using var writer = new ParquetFileWriter(managedOutput, _parquetSharpColumns!); - using RowGroupWriter rowGroup = writer.AppendRowGroup(); - - using (LogicalColumnWriter vendorWriter = rowGroup.NextColumn().LogicalWriter()) - { - vendorWriter.WriteBatch(_vendorIds!); - } - - using (LogicalColumnWriter rateCodeWriter = rowGroup.NextColumn().LogicalWriter()) - { - rateCodeWriter.WriteBatch(_rateCodes!); - } - - using (LogicalColumnWriter passengerCountWriter = rowGroup.NextColumn().LogicalWriter()) - { - passengerCountWriter.WriteBatch(_passengerCounts!); - } - - using (LogicalColumnWriter tripDistanceWriter = rowGroup.NextColumn().LogicalWriter()) - { - tripDistanceWriter.WriteBatch(_tripDistances!); - } - - using (LogicalColumnWriter paymentTypeWriter = rowGroup.NextColumn().LogicalWriter()) - { - paymentTypeWriter.WriteBatch(_paymentTypes!); - } - - using (LogicalColumnWriter fareWriter = rowGroup.NextColumn().LogicalWriter()) - { - fareWriter.WriteBatch(_fareAmounts!); - } - - writer.Close(); - } - finally - { - TryDelete(path); - } - } - - private static string GetDataDirectory() - { - return Path.Combine(AppContext.BaseDirectory, "Data"); - } - - private static string? TryProjectDataPath(string fileName) - { - string projectPath = Path.GetFullPath(Path.Combine(AppContext.BaseDirectory, "../../../Data", fileName)); - return IOFile.Exists(projectPath) ? projectPath : null; - } - - private async Task EnsureDatasetsAsync() - { - if (!DatasetFiles.TryGetValue(Dataset, out string[]? files)) - { - throw new ArgumentOutOfRangeException(nameof(Dataset), Dataset, "Unknown dataset alias"); - } - - string dataDir = GetDataDirectory(); - Directory.CreateDirectory(dataDir); - - var paths = new List(files.Length); - using var httpClient = new HttpClient(); - - foreach (string file in files) - { - string dataPath = Path.Combine(dataDir, file); - if (IOFile.Exists(dataPath)) - { - paths.Add(dataPath); - continue; - } - - string? projectPath = TryProjectDataPath(file); - if (projectPath != null) - { - paths.Add(projectPath); - continue; - } - - string url = $"{BaseUrl}{file}"; - byte[] bytes = await httpClient.GetByteArrayAsync(url); - await IOFile.WriteAllBytesAsync(dataPath, bytes); - paths.Add(dataPath); - } - - return paths.ToArray(); - } - - private static async Task LoadParquetColumnsAsync(IEnumerable paths) - { - var vendorIds = new List(); - var rateCodes = new List(); - var passengerCounts = new List(); - var tripDistances = new List(); - var paymentTypes = new List(); - var fareAmounts = new List(); - - foreach (string path in paths) - { - using FileStream fs = IOFile.OpenRead(path); - using ParquetReaderNet reader = await ParquetReaderNet.CreateAsync(fs); - - DataField vendor = FindField(reader.Schema, "vendorid"); - DataField rateCode = FindField(reader.Schema, "ratecodeid"); - DataField passengerCount = FindField(reader.Schema, "passenger_count"); - DataField tripDistance = FindField(reader.Schema, "trip_distance"); - DataField paymentType = FindField(reader.Schema, "payment_type"); - DataField fareAmount = FindField(reader.Schema, "fare_amount"); - - for (int i = 0; i < reader.RowGroupCount; i++) - { - using ParquetRowGroupReader rg = reader.OpenRowGroupReader(i); - - vendorIds.AddRange(ToNullableInt((await rg.ReadColumnAsync(vendor)).Data)); - rateCodes.AddRange(ToNullableInt((await rg.ReadColumnAsync(rateCode)).Data)); - passengerCounts.AddRange(ToNullableDouble((await rg.ReadColumnAsync(passengerCount)).Data)); - tripDistances.AddRange(ToNullableDouble((await rg.ReadColumnAsync(tripDistance)).Data)); - paymentTypes.AddRange(ToNullableInt((await rg.ReadColumnAsync(paymentType)).Data)); - fareAmounts.AddRange(ToNullableDouble((await rg.ReadColumnAsync(fareAmount)).Data)); - } - } - - return new TaxiColumns( - vendorIds.ToArray(), - rateCodes.ToArray(), - passengerCounts.ToArray(), - tripDistances.ToArray(), - paymentTypes.ToArray(), - fareAmounts.ToArray()); - } - - private static DataField FindField(ParquetSchema schema, string name) - { - return schema.DataFields.First(f => string.Equals(f.Name, name, StringComparison.OrdinalIgnoreCase)); - } - - private static IReadOnlyList ToNullableInt(Array data) => - data switch - { - int?[] v => v, - int[] v => v.Select(i => (int?)i).ToArray(), - long?[] v => v.Select(i => i.HasValue ? (int?)checked((int)i.Value) : null).ToArray(), - long[] v => v.Select(i => (int?)checked((int)i)).ToArray(), - double?[] v => v.Select(i => i.HasValue ? (int?)checked((int)i.Value) : null).ToArray(), - double[] v => v.Select(i => (int?)checked((int)i)).ToArray(), - float?[] v => v.Select(i => i.HasValue ? (int?)checked((int)i.Value) : null).ToArray(), - float[] v => v.Select(i => (int?)checked((int)i)).ToArray(), - _ => throw new InvalidOperationException($"Unsupported numeric type: {data.GetType()}") - }; - - private static IReadOnlyList ToNullableDouble(Array data) => - data switch - { - double?[] v => v, - double[] v => v.Select(i => (double?)i).ToArray(), - float?[] v => v.Select(i => i.HasValue ? (double?)i.Value : null).ToArray(), - float[] v => v.Select(i => (double?)i).ToArray(), - int?[] v => v.Select(i => i.HasValue ? (double?)i.Value : null).ToArray(), - int[] v => v.Select(i => (double?)i).ToArray(), - long?[] v => v.Select(i => i.HasValue ? (double?)i.Value : null).ToArray(), - long[] v => v.Select(i => (double?)i).ToArray(), - _ => throw new InvalidOperationException($"Unsupported numeric type: {data.GetType()}") - }; - - private static void TryDelete(string path) - { - try - { - if (IOFile.Exists(path)) - { - IOFile.Delete(path); - } - } - catch - { - // best effort cleanup - } - } -} - -internal readonly struct TaxiColumns -{ - public TaxiColumns(int?[] vendorIds, int?[] rateCodes, double?[] passengerCounts, double?[] tripDistances, int?[] paymentTypes, double?[] fareAmounts) - { - VendorIds = vendorIds; - RateCodes = rateCodes; - PassengerCounts = passengerCounts; - TripDistances = tripDistances; - PaymentTypes = paymentTypes; - FareAmounts = fareAmounts; - } - - public int?[] VendorIds { get; } - public int?[] RateCodes { get; } - public double?[] PassengerCounts { get; } - public double?[] TripDistances { get; } - public int?[] PaymentTypes { get; } - public double?[] FareAmounts { get; } -} diff --git a/src/Parquet.PerfRunner/Benchmarks/WriteBenchmark.cs b/src/Parquet.PerfRunner/Benchmarks/WriteBenchmark.cs index 28b9b19f..ec07fe39 100644 --- a/src/Parquet.PerfRunner/Benchmarks/WriteBenchmark.cs +++ b/src/Parquet.PerfRunner/Benchmarks/WriteBenchmark.cs @@ -26,14 +26,12 @@ public class WriteBenchmark : BenchmarkBase { private Array? _ar; [GlobalSetup] - public Task SetupAsync() { + public void Setup() { _schema = new ParquetSchema(new DataField("test", DataType!)); _ar = CreateTestData(DataType); - _c = new DataColumn(_schema.DataFields[0], _ar); + _c = new DataColumn(_schema.GetDataFields()[0], _ar); _psc = new Column(DataType!, "test"); - - return Task.CompletedTask; } [Benchmark] diff --git a/src/Parquet.PerfRunner/LogicalEncoding.cs b/src/Parquet.PerfRunner/LogicalEncoding.cs new file mode 100644 index 00000000..b62571b8 --- /dev/null +++ b/src/Parquet.PerfRunner/LogicalEncoding.cs @@ -0,0 +1,37 @@ +using System; +using System.Collections.Generic; +using System.Linq; +using System.Text; +using System.Threading.Tasks; + +namespace Parquet.PerfRunner; + +/// +/// Allow to easily control the logical encoding used in benchmarks. +/// +public enum LogicalEncoding { + Plain, + RleDictionary, + DeltaBinaryPacked +} + +public static class LogicalEncodingExtensions { + public static ParquetOptions CreateOptions(this LogicalEncoding encoding) => + encoding switch { + LogicalEncoding.RleDictionary => new ParquetOptions { + UseDictionaryEncoding = true, + // Force dictionary extraction for every supported column to match the ParquetSharp 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") + }; +} diff --git a/src/Parquet.PerfRunner/Parquet.PerfRunner.csproj b/src/Parquet.PerfRunner/Parquet.PerfRunner.csproj index 5a7ebd3c..a7aca8ef 100644 --- a/src/Parquet.PerfRunner/Parquet.PerfRunner.csproj +++ b/src/Parquet.PerfRunner/Parquet.PerfRunner.csproj @@ -8,6 +8,8 @@ latest true ../fake.snk + + @@ -15,7 +17,11 @@ - + + + + + diff --git a/src/Parquet.PerfRunner/Program.cs b/src/Parquet.PerfRunner/Program.cs index 2778584d..a7cf9d63 100644 --- a/src/Parquet.PerfRunner/Program.cs +++ b/src/Parquet.PerfRunner/Program.cs @@ -1,21 +1,30 @@ // for performance tests only +using BenchmarkDotNet.Configs; using BenchmarkDotNet.Running; using Parquet; using Parquet.PerfRunner.Benchmarks; -if(args.Length == 1) { - switch(args[0]) { - case "write": - BenchmarkRunner.Run(); - break; - case "progression": - VersionedBenchmark.Run(); - break; - case "taxi": - BenchmarkRunner.Run(); - break; - } +bool haveArgs = args.Length > 0; +string benchName; +if(!haveArgs) { + Console.WriteLine("Enter the benchmark name to run:"); + benchName = Console.ReadLine()!; } else { - await new DataTypes().NullableInts(); + benchName = args[0]; +} + +switch(benchName) { + case "write": + BenchmarkRunner.Run(); + break; + case "progression": + BenchmarkRunner.Run(); + break; + case "sharp-compare-taxi": + BenchmarkRunner.Run(); + break; + case "self-compare-taxi": + BenchmarkRunner.Run(); + break; } diff --git a/src/Parquet.PerfRunner/Taxis/ParquetSharpTaxiSchema.cs b/src/Parquet.PerfRunner/Taxis/ParquetSharpTaxiSchema.cs new file mode 100644 index 00000000..13e1fdb5 --- /dev/null +++ b/src/Parquet.PerfRunner/Taxis/ParquetSharpTaxiSchema.cs @@ -0,0 +1,76 @@ +using ParquetSharp; +using Column = ParquetSharp.Column; + +namespace Parquet.PerfRunner.Taxis; + +sealed record ParquetSharpTaxiSchema(Column[] Columns) { + public static ParquetSharpTaxiSchema Full() { + Column[] columns = [ + new Column("VendorID"), + new Column("tpep_pickup_datetime"), + new Column("tpep_dropoff_datetime"), + new Column("passenger_count"), + new Column("trip_distance"), + new Column("RatecodeID"), + new Column("store_and_fwd_flag"), + new Column("PULocationID"), + new Column("DOLocationID"), + new Column("payment_type"), + new Column("fare_amount"), + new Column("extra"), + new Column("mta_tax"), + new Column("tip_amount"), + new Column("tolls_amount"), + new Column("improvement_surcharge"), + new Column("total_amount"), + new Column("congestion_surcharge"), + new Column("Airport_fee") + ]; + + return new ParquetSharpTaxiSchema(columns); + } + + public static ParquetSharpTaxiSchema Small() { + Column[] columns = [ + new Column("VendorID"), + new Column("passenger_count"), + new Column("trip_distance"), + new Column("RatecodeID"), + new Column("payment_type"), + new Column("fare_amount") + ]; + return new ParquetSharpTaxiSchema(columns); + } + + public WriterProperties CreateParquetSharpWriterProperties(LogicalEncoding encoding) { + WriterPropertiesBuilder builder = new WriterPropertiesBuilder() + .Compression(Compression.Snappy); + + switch(encoding) { + case LogicalEncoding.RleDictionary: + builder.EnableDictionary(); + builder.Encoding(Encoding.Plain); + break; + case LogicalEncoding.DeltaBinaryPacked: + builder.DisableDictionary(); + builder.Encoding(Encoding.Plain); + foreach(Column column in Columns) { + Type type = Nullable.GetUnderlyingType(column.LogicalSystemType) + ?? column.LogicalSystemType; + + if(type == typeof(int) || type == typeof(long)) { + builder.Encoding(column.Name, Encoding.DeltaBinaryPacked); + } + } + + break; + case LogicalEncoding.Plain: + builder.DisableDictionary(); + builder.Encoding(Encoding.Plain); + break; + default: + throw new ArgumentOutOfRangeException(nameof(encoding), encoding, "Unknown logical encoding"); + } + return builder.Build(); + } +} diff --git a/src/Parquet.PerfRunner/Taxis/TaxiDataset.cs b/src/Parquet.PerfRunner/Taxis/TaxiDataset.cs new file mode 100644 index 00000000..4b7c7ed6 --- /dev/null +++ b/src/Parquet.PerfRunner/Taxis/TaxiDataset.cs @@ -0,0 +1,24 @@ +namespace Parquet.PerfRunner.Taxis; + +record struct TaxiDataset( + int?[] VendorID, + DateTime?[] tpep_pickup_datetime, + DateTime?[] tpep_dropoff_datetime, + long?[] passenger_count, + double?[] trip_distance, + long?[] RatecodeID, + string?[] store_and_fwd_flag, + int?[] PULocationID, + int?[] DOLocationID, + long?[] payment_type, + double?[] fare_amount, + double?[] extra, + double?[] mta_tax, + double?[] tip_amount, + double?[] tolls_amount, + double?[] improvement_surcharge, + double?[] total_amount, + double?[] congestion_surcharge, + double?[] Airport_fee + ) { +} diff --git a/src/Parquet.PerfRunner/Taxis/TaxiDatasetLoader.cs b/src/Parquet.PerfRunner/Taxis/TaxiDatasetLoader.cs new file mode 100644 index 00000000..8f8dadce --- /dev/null +++ b/src/Parquet.PerfRunner/Taxis/TaxiDatasetLoader.cs @@ -0,0 +1,148 @@ +using Parquet.Schema; +using IOFile = System.IO.File; +using ParquetReaderNet = Parquet.ParquetReader; + +namespace Parquet.PerfRunner.Taxis; + +class TaxiDatasetLoader { + + public static TaxiDatasetLoader Instance { get; } = new TaxiDatasetLoader(); + TaxiDatasetLoader() {} + + + // Official NYC TLC parquet. The "tripdata" dataset is a single month; "tripdata-large" combines multiple months. + private const string _baseUrl = "https://d37ci6vzurychx.cloudfront.net/trip-data/"; + private static readonly HttpClient _httpClient = new(); + + private static readonly IReadOnlyDictionary DatasetFiles = new Dictionary + { + { "tripdata", new[] { "yellow_tripdata_2024-01.parquet" } }, + { "tripdata-large", new[] { "yellow_tripdata_2024-01.parquet", "yellow_tripdata_2024-02.parquet", "yellow_tripdata_2024-03.parquet" } } + }; + + static readonly string _dataFolder = Path.Combine(Path.GetTempPath(), "nyc_tlc_data"); + + public async Task LoadAsync(string datasetName) { + Directory.CreateDirectory(_dataFolder); + + if(!DatasetFiles.TryGetValue(datasetName, out string[]? files)) { + throw new ArgumentException($"Unknown dataset: {datasetName}", nameof(datasetName)); + } + + await PreloadFiles(files); + TaxiDataset dataset = await LoadParquetColumnsAsync(files); + return dataset; + } + + private static async Task LoadParquetColumnsAsync(IEnumerable paths) { + List vendorIds = []; + List pickupTimes = []; + List dropoffTimes = []; + List passengerCounts = []; + List tripDistances = []; + List rateCodes = []; + List storeAndFwdFlags = []; + List puLocationIds = []; + List doLocationIds = []; + List paymentTypes = []; + List fareAmounts = []; + List extras = []; + List mtaTaxes = []; + List tipAmounts = []; + List tollsAmounts = []; + List improvementSurcharges = []; + List totalAmounts = []; + List congestionSurcharges = []; + List airportFees = []; + + foreach(string path in paths) { + using FileStream fs = IOFile.OpenRead(Path.Combine(_dataFolder, path)); + using ParquetReaderNet reader = await ParquetReaderNet.CreateAsync(fs); + + DataField vendor = FindDataField(reader.Schema, "VendorID"); + DataField pickup = FindDataField(reader.Schema, "tpep_pickup_datetime"); + DataField dropoff = FindDataField(reader.Schema, "tpep_dropoff_datetime"); + DataField passengerCount = FindDataField(reader.Schema, "passenger_count"); + DataField tripDistance = FindDataField(reader.Schema, "trip_distance"); + DataField rateCode = FindDataField(reader.Schema, "RatecodeID"); + DataField storeAndFwdFlag = FindDataField(reader.Schema, "store_and_fwd_flag"); + DataField puLocationId = FindDataField(reader.Schema, "PULocationID"); + DataField doLocationId = FindDataField(reader.Schema, "DOLocationID"); + DataField paymentType = FindDataField(reader.Schema, "payment_type"); + DataField fareAmount = FindDataField(reader.Schema, "fare_amount"); + DataField extra = FindDataField(reader.Schema, "extra"); + DataField mtaTax = FindDataField(reader.Schema, "mta_tax"); + DataField tipAmount = FindDataField(reader.Schema, "tip_amount"); + DataField tollsAmount = FindDataField(reader.Schema, "tolls_amount"); + DataField improvementSurcharge = FindDataField(reader.Schema, "improvement_surcharge"); + DataField totalAmount = FindDataField(reader.Schema, "total_amount"); + DataField congestionSurcharge = FindDataField(reader.Schema, "congestion_surcharge"); + DataField airportFee = FindDataField(reader.Schema, "Airport_fee"); + + for(int i = 0; i < reader.RowGroupCount; i++) { + using ParquetRowGroupReader rg = reader.OpenRowGroupReader(i); + + vendorIds.AddRange((int?[])(await rg.ReadColumnAsync(vendor)).Data); + pickupTimes.AddRange((DateTime?[])(await rg.ReadColumnAsync(pickup)).Data); + dropoffTimes.AddRange((DateTime?[])(await rg.ReadColumnAsync(dropoff)).Data); + passengerCounts.AddRange((long?[])(await rg.ReadColumnAsync(passengerCount)).Data); + tripDistances.AddRange((double?[])(await rg.ReadColumnAsync(tripDistance)).Data); + rateCodes.AddRange((long?[])(await rg.ReadColumnAsync(rateCode)).Data); + storeAndFwdFlags.AddRange((string?[])(await rg.ReadColumnAsync(storeAndFwdFlag)).Data); + puLocationIds.AddRange((int?[])(await rg.ReadColumnAsync(puLocationId)).Data); + doLocationIds.AddRange((int?[])(await rg.ReadColumnAsync(doLocationId)).Data); + paymentTypes.AddRange((long?[])(await rg.ReadColumnAsync(paymentType)).Data); + fareAmounts.AddRange((double?[])(await rg.ReadColumnAsync(fareAmount)).Data); + extras.AddRange((double?[])(await rg.ReadColumnAsync(extra)).Data); + mtaTaxes.AddRange((double?[])(await rg.ReadColumnAsync(mtaTax)).Data); + tipAmounts.AddRange((double?[])(await rg.ReadColumnAsync(tipAmount)).Data); + tollsAmounts.AddRange((double?[])(await rg.ReadColumnAsync(tollsAmount)).Data); + improvementSurcharges.AddRange((double?[])(await rg.ReadColumnAsync(improvementSurcharge)).Data); + totalAmounts.AddRange((double?[])(await rg.ReadColumnAsync(totalAmount)).Data); + congestionSurcharges.AddRange((double?[])(await rg.ReadColumnAsync(congestionSurcharge)).Data); + airportFees.AddRange((double?[])(await rg.ReadColumnAsync(airportFee)).Data); + + } + } + + return new TaxiDataset( + vendorIds.ToArray(), + pickupTimes.ToArray(), + dropoffTimes.ToArray(), + passengerCounts.ToArray(), + tripDistances.ToArray(), + rateCodes.ToArray(), + storeAndFwdFlags.ToArray(), + puLocationIds.ToArray(), + doLocationIds.ToArray(), + paymentTypes.ToArray(), + fareAmounts.ToArray(), + extras.ToArray(), + mtaTaxes.ToArray(), + tipAmounts.ToArray(), + tollsAmounts.ToArray(), + improvementSurcharges.ToArray(), + totalAmounts.ToArray(), + congestionSurcharges.ToArray(), + airportFees.ToArray() + ); + } + + private static DataField FindDataField(ParquetSchema schema, string v) => schema.GetDataFields().First(f => f.Name == v); + + private static async Task PreloadFiles(IEnumerable files) + => await Task.WhenAll( + files.Select(DownloadFileAsync) + ); + + private static async Task DownloadFileAsync(string fileName) { + string dataPath = Path.Combine(_dataFolder, fileName); + if(IOFile.Exists(dataPath)) { + return; + } + string url = $"{_baseUrl}{fileName}"; + using Stream stream = await _httpClient.GetStreamAsync(url); + using FileStream fileStream = IOFile.Create(dataPath); + await stream.CopyToAsync(fileStream); + } +} diff --git a/src/Parquet.PerfRunner/Taxis/TaxiSchema.cs b/src/Parquet.PerfRunner/Taxis/TaxiSchema.cs new file mode 100644 index 00000000..c5784fe0 --- /dev/null +++ b/src/Parquet.PerfRunner/Taxis/TaxiSchema.cs @@ -0,0 +1,115 @@ +using Parquet.Data; +using Parquet.Schema; +using ParquetSharp; +using RowGroupWriter = ParquetSharp.RowGroupWriter; +using ParquetWriterNet = Parquet.ParquetWriter; + +namespace Parquet.PerfRunner.Taxis; + +sealed record TaxiSchema(TaxiSchemaKind Kind, ParquetSchema Schema, DataColumn[] Columns) { + public static TaxiSchema Full(TaxiDataset dataset) { + ParquetSchema schema = new( + new DataField("VendorID"), + new DateTimeDataField("tpep_pickup_datetime", DateTimeFormat.Timestamp, isNullable: true, unit: DateTimeTimeUnit.Micros), + new DateTimeDataField("tpep_dropoff_datetime", DateTimeFormat.Timestamp, isNullable: true, unit: DateTimeTimeUnit.Micros), + new DataField("passenger_count"), + new DataField("trip_distance"), + new DataField("RatecodeID"), + new DataField("store_and_fwd_flag"), + new DataField("PULocationID"), + new DataField("DOLocationID"), + new DataField("payment_type"), + new DataField("fare_amount"), + new DataField("extra"), + new DataField("mta_tax"), + new DataField("tip_amount"), + new DataField("tolls_amount"), + new DataField("improvement_surcharge"), + new DataField("total_amount"), + new DataField("congestion_surcharge"), + new DataField("Airport_fee") + ); + DataField[] dataFields = schema.GetDataFields(); // better API not available in older version we compare against. + DataColumn[] columns = [ + new DataColumn(dataFields[0], dataset.VendorID), + new DataColumn(dataFields[1], dataset.tpep_pickup_datetime), + new DataColumn(dataFields[2], dataset.tpep_dropoff_datetime), + new DataColumn(dataFields[3], dataset.passenger_count), + new DataColumn(dataFields[4], dataset.trip_distance), + new DataColumn(dataFields[5], dataset.RatecodeID), + new DataColumn(dataFields[6], dataset.store_and_fwd_flag), + new DataColumn(dataFields[7], dataset.PULocationID), + new DataColumn(dataFields[8], dataset.DOLocationID), + new DataColumn(dataFields[9], dataset.payment_type), + new DataColumn(dataFields[10], dataset.fare_amount), + new DataColumn(dataFields[11], dataset.extra), + new DataColumn(dataFields[12], dataset.mta_tax), + new DataColumn(dataFields[13], dataset.tip_amount), + new DataColumn(dataFields[14], dataset.tolls_amount), + new DataColumn(dataFields[15], dataset.improvement_surcharge), + new DataColumn(dataFields[16], dataset.total_amount), + new DataColumn(dataFields[17], dataset.congestion_surcharge), + new DataColumn(dataFields[18], dataset.Airport_fee) + ]; + + return new TaxiSchema(TaxiSchemaKind.Full, schema, columns); + } + + public static TaxiSchema Small(TaxiDataset dataset) { + ParquetSchema schema = new( + new DataField("VendorID"), + new DataField("passenger_count"), + new DataField("trip_distance"), + new DataField("RatecodeID"), + new DataField("payment_type"), + new DataField("fare_amount") + ); + DataField[] dataFields = schema.GetDataFields(); + DataColumn[] columns = [ + new DataColumn(dataFields[0], dataset.VendorID), + new DataColumn(dataFields[1], dataset.passenger_count), + new DataColumn(dataFields[2], dataset.trip_distance), + new DataColumn(dataFields[3], dataset.RatecodeID), + new DataColumn(dataFields[4], dataset.payment_type), + new DataColumn(dataFields[5], dataset.fare_amount) + ]; + + return new TaxiSchema(TaxiSchemaKind.Small, schema, columns); + } + + public void WriteParquetSharp(RowGroupWriter rowGroup, TaxiDataset dataset) { + static void WriteColumn(RowGroupWriter groupWriter, T[] data) { + using LogicalColumnWriter columnWriter = groupWriter.NextColumn().LogicalWriter(); + columnWriter.WriteBatch(data); + } + + if(Kind == TaxiSchemaKind.Small) { + WriteColumn(rowGroup, dataset.VendorID); + WriteColumn(rowGroup, dataset.passenger_count); + WriteColumn(rowGroup, dataset.trip_distance); + WriteColumn(rowGroup, dataset.RatecodeID); + WriteColumn(rowGroup, dataset.payment_type); + WriteColumn(rowGroup, dataset.fare_amount); + } else { + WriteColumn(rowGroup, dataset.VendorID); + WriteColumn(rowGroup, dataset.tpep_pickup_datetime); + WriteColumn(rowGroup, dataset.tpep_dropoff_datetime); + WriteColumn(rowGroup, dataset.passenger_count); + WriteColumn(rowGroup, dataset.trip_distance); + WriteColumn(rowGroup, dataset.RatecodeID); + WriteColumn(rowGroup, dataset.store_and_fwd_flag); + WriteColumn(rowGroup, dataset.PULocationID); + WriteColumn(rowGroup, dataset.DOLocationID); + WriteColumn(rowGroup, dataset.payment_type); + WriteColumn(rowGroup, dataset.fare_amount); + WriteColumn(rowGroup, dataset.extra); + WriteColumn(rowGroup, dataset.mta_tax); + WriteColumn(rowGroup, dataset.tip_amount); + WriteColumn(rowGroup, dataset.tolls_amount); + WriteColumn(rowGroup, dataset.improvement_surcharge); + WriteColumn(rowGroup, dataset.total_amount); + WriteColumn(rowGroup, dataset.congestion_surcharge); + WriteColumn(rowGroup, dataset.Airport_fee); + } + } +} diff --git a/src/Parquet.PerfRunner/Taxis/TaxiSchemaKind.cs b/src/Parquet.PerfRunner/Taxis/TaxiSchemaKind.cs new file mode 100644 index 00000000..31d8e384 --- /dev/null +++ b/src/Parquet.PerfRunner/Taxis/TaxiSchemaKind.cs @@ -0,0 +1,6 @@ +namespace Parquet.PerfRunner.Taxis; + +enum TaxiSchemaKind { + Small, + Full +} From 23aa83ec05e68e9a4bf91324c854a5e3e000727c Mon Sep 17 00:00:00 2001 From: Kuinox Date: Sun, 14 Dec 2025 03:29:43 +0100 Subject: [PATCH 3/6] And now the self compare bench. --- .../Benchmarks/SelfComparisonBenchmark.cs | 17 +++++++++++++---- 1 file changed, 13 insertions(+), 4 deletions(-) diff --git a/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs b/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs index 213dd62e..a0b00dea 100644 --- a/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs +++ b/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs @@ -1,18 +1,27 @@ using BenchmarkDotNet.Attributes; +using BenchmarkDotNet.Configs; +using BenchmarkDotNet.Jobs; using Parquet.Data; using Parquet.Meta; using Parquet.PerfRunner.Taxis; -using ParquetSharp; -using ParquetSharp.IO; -using IOFile = System.IO.File; using ParquetWriterNet = Parquet.ParquetWriter; namespace Parquet.PerfRunner.Benchmarks; +[Config(typeof(NuConfig))] [MemoryDiagnoser] [MarkdownExporter] -[ShortRunJob] 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.4.0") + .WithMsBuildArguments("/p:ParquetNuGetVersion=5.4.0")); + } + } + [Params("tripdata", "tripdata-large")] public string Dataset { get; set; } = "tripdata"; From 5f5805d87910108c4e2a492480b550f2c670d940 Mon Sep 17 00:00:00 2001 From: Kuinox Date: Sun, 14 Dec 2025 03:54:52 +0100 Subject: [PATCH 4/6] And some cleanup --- .../ParquetSharpComparisonBenchmark.cs | 15 ++--- .../Benchmarks/Progression.cs | 9 +-- .../Benchmarks/SelfComparisonBenchmark.cs | 1 - src/Parquet.PerfRunner/LogicalEncoding.cs | 8 +-- src/Parquet.PerfRunner/Program.cs | 2 - .../Taxis/ParquetSharpTaxiSchema.cs | 55 ++++++++++++++++++- src/Parquet.PerfRunner/Taxis/TaxiSchema.cs | 54 ++++-------------- .../Taxis/TaxiSchemaKind.cs | 6 -- 8 files changed, 71 insertions(+), 79 deletions(-) delete mode 100644 src/Parquet.PerfRunner/Taxis/TaxiSchemaKind.cs diff --git a/src/Parquet.PerfRunner/Benchmarks/ParquetSharpComparisonBenchmark.cs b/src/Parquet.PerfRunner/Benchmarks/ParquetSharpComparisonBenchmark.cs index d8b1241c..6de1c851 100644 --- a/src/Parquet.PerfRunner/Benchmarks/ParquetSharpComparisonBenchmark.cs +++ b/src/Parquet.PerfRunner/Benchmarks/ParquetSharpComparisonBenchmark.cs @@ -1,11 +1,9 @@ using BenchmarkDotNet.Attributes; using Parquet.Data; -using Parquet.Meta; using Parquet.PerfRunner.Taxis; using ParquetSharp; using ParquetSharp.IO; using IOFile = System.IO.File; -using ParquetSharpEncoding = ParquetSharp.Encoding; using ParquetWriterNet = Parquet.ParquetWriter; namespace Parquet.PerfRunner.Benchmarks; @@ -64,8 +62,6 @@ public async Task ParquetNetAsync() { [Benchmark(Description = "Parquet.Net -> Disk")] public async Task ParquetNetToDiskAsync() { string path = Path.Combine(Path.GetTempPath(), GetFileName("parquetnet")); - if(Path.Exists(path)) - IOFile.Delete(path); try { await using FileStream output = IOFile.Create(path); @@ -77,7 +73,7 @@ public async Task ParquetNetToDiskAsync() { await rowGroup.WriteColumnAsync(column); } } finally { - //IOFile.Delete(path); + IOFile.Delete(path); } } @@ -88,26 +84,25 @@ public void ParquetSharp() { using var writer = new ParquetFileWriter(managedOutput, _parquetSharpSchema.Columns, _parquetSharpOptions); using RowGroupWriter rowGroup = writer.AppendRowGroup(); - _parquetNetSchema.WriteParquetSharp(rowGroup, _dataset); + _parquetSharpSchema.Write(rowGroup, _dataset); writer.Close(); } [Benchmark(Description = "ParquetSharp -> Disk")] public void ParquetSharpToDisk() { string path = Path.Combine(Path.GetTempPath(), GetFileName("parquetsharp")); - if(Path.Exists(path)) - IOFile.Delete(path); + 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(); - _parquetNetSchema.WriteParquetSharp(rowGroup, _dataset); + _parquetSharpSchema.Write(rowGroup, _dataset); writer.Close(); } finally { - //IOFile.Delete(path); + IOFile.Delete(path); } } } diff --git a/src/Parquet.PerfRunner/Benchmarks/Progression.cs b/src/Parquet.PerfRunner/Benchmarks/Progression.cs index 8a8b7940..b78d4465 100644 --- a/src/Parquet.PerfRunner/Benchmarks/Progression.cs +++ b/src/Parquet.PerfRunner/Benchmarks/Progression.cs @@ -1,13 +1,6 @@ -using System; -using System.Collections.Generic; -using System.Linq; -using System.Text; -using System.Threading.Tasks; -using BenchmarkDotNet.Attributes; +using BenchmarkDotNet.Attributes; using BenchmarkDotNet.Configs; using BenchmarkDotNet.Jobs; -using BenchmarkDotNet.Running; -using Microsoft.Diagnostics.Tracing.Parsers.Kernel; using Parquet.Data; using Parquet.Schema; diff --git a/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs b/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs index a0b00dea..071569ed 100644 --- a/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs +++ b/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs @@ -2,7 +2,6 @@ using BenchmarkDotNet.Configs; using BenchmarkDotNet.Jobs; using Parquet.Data; -using Parquet.Meta; using Parquet.PerfRunner.Taxis; using ParquetWriterNet = Parquet.ParquetWriter; diff --git a/src/Parquet.PerfRunner/LogicalEncoding.cs b/src/Parquet.PerfRunner/LogicalEncoding.cs index b62571b8..9c49fa0a 100644 --- a/src/Parquet.PerfRunner/LogicalEncoding.cs +++ b/src/Parquet.PerfRunner/LogicalEncoding.cs @@ -1,10 +1,4 @@ -using System; -using System.Collections.Generic; -using System.Linq; -using System.Text; -using System.Threading.Tasks; - -namespace Parquet.PerfRunner; +namespace Parquet.PerfRunner; /// /// Allow to easily control the logical encoding used in benchmarks. diff --git a/src/Parquet.PerfRunner/Program.cs b/src/Parquet.PerfRunner/Program.cs index a7cf9d63..9ba5287b 100644 --- a/src/Parquet.PerfRunner/Program.cs +++ b/src/Parquet.PerfRunner/Program.cs @@ -1,8 +1,6 @@ // for performance tests only -using BenchmarkDotNet.Configs; using BenchmarkDotNet.Running; -using Parquet; using Parquet.PerfRunner.Benchmarks; bool haveArgs = args.Length > 0; diff --git a/src/Parquet.PerfRunner/Taxis/ParquetSharpTaxiSchema.cs b/src/Parquet.PerfRunner/Taxis/ParquetSharpTaxiSchema.cs index 13e1fdb5..2c290364 100644 --- a/src/Parquet.PerfRunner/Taxis/ParquetSharpTaxiSchema.cs +++ b/src/Parquet.PerfRunner/Taxis/ParquetSharpTaxiSchema.cs @@ -3,7 +3,20 @@ namespace Parquet.PerfRunner.Taxis; -sealed record ParquetSharpTaxiSchema(Column[] Columns) { +sealed class ParquetSharpTaxiSchema { + enum Kind { + Small, + Full + } + + public Column[] Columns { get; } + private readonly Kind _kind; + + private ParquetSharpTaxiSchema(Column[] columns, Kind kind) { + Columns = columns; + _kind = kind; + } + public static ParquetSharpTaxiSchema Full() { Column[] columns = [ new Column("VendorID"), @@ -27,7 +40,7 @@ public static ParquetSharpTaxiSchema Full() { new Column("Airport_fee") ]; - return new ParquetSharpTaxiSchema(columns); + return new ParquetSharpTaxiSchema(columns, Kind.Full); } public static ParquetSharpTaxiSchema Small() { @@ -39,7 +52,7 @@ public static ParquetSharpTaxiSchema Small() { new Column("payment_type"), new Column("fare_amount") ]; - return new ParquetSharpTaxiSchema(columns); + return new ParquetSharpTaxiSchema(columns, Kind.Small); } public WriterProperties CreateParquetSharpWriterProperties(LogicalEncoding encoding) { @@ -73,4 +86,40 @@ public WriterProperties CreateParquetSharpWriterProperties(LogicalEncoding encod } return builder.Build(); } + + public void Write(RowGroupWriter rowGroup, TaxiDataset dataset) { + static void WriteColumn(RowGroupWriter groupWriter, T[] data) { + using LogicalColumnWriter columnWriter = groupWriter.NextColumn().LogicalWriter(); + columnWriter.WriteBatch(data); + } + + if(_kind == Kind.Small) { + WriteColumn(rowGroup, dataset.VendorID); + WriteColumn(rowGroup, dataset.passenger_count); + WriteColumn(rowGroup, dataset.trip_distance); + WriteColumn(rowGroup, dataset.RatecodeID); + WriteColumn(rowGroup, dataset.payment_type); + WriteColumn(rowGroup, dataset.fare_amount); + } else { + WriteColumn(rowGroup, dataset.VendorID); + WriteColumn(rowGroup, dataset.tpep_pickup_datetime); + WriteColumn(rowGroup, dataset.tpep_dropoff_datetime); + WriteColumn(rowGroup, dataset.passenger_count); + WriteColumn(rowGroup, dataset.trip_distance); + WriteColumn(rowGroup, dataset.RatecodeID); + WriteColumn(rowGroup, dataset.store_and_fwd_flag); + WriteColumn(rowGroup, dataset.PULocationID); + WriteColumn(rowGroup, dataset.DOLocationID); + WriteColumn(rowGroup, dataset.payment_type); + WriteColumn(rowGroup, dataset.fare_amount); + WriteColumn(rowGroup, dataset.extra); + WriteColumn(rowGroup, dataset.mta_tax); + WriteColumn(rowGroup, dataset.tip_amount); + WriteColumn(rowGroup, dataset.tolls_amount); + WriteColumn(rowGroup, dataset.improvement_surcharge); + WriteColumn(rowGroup, dataset.total_amount); + WriteColumn(rowGroup, dataset.congestion_surcharge); + WriteColumn(rowGroup, dataset.Airport_fee); + } + } } diff --git a/src/Parquet.PerfRunner/Taxis/TaxiSchema.cs b/src/Parquet.PerfRunner/Taxis/TaxiSchema.cs index c5784fe0..2f0a9b76 100644 --- a/src/Parquet.PerfRunner/Taxis/TaxiSchema.cs +++ b/src/Parquet.PerfRunner/Taxis/TaxiSchema.cs @@ -1,12 +1,18 @@ using Parquet.Data; using Parquet.Schema; -using ParquetSharp; -using RowGroupWriter = ParquetSharp.RowGroupWriter; -using ParquetWriterNet = Parquet.ParquetWriter; namespace Parquet.PerfRunner.Taxis; -sealed record TaxiSchema(TaxiSchemaKind Kind, ParquetSchema Schema, DataColumn[] Columns) { +sealed class TaxiSchema { + + public ParquetSchema Schema { get; } + public DataColumn[] Columns { get; } + + private TaxiSchema(ParquetSchema schema, DataColumn[] columns) { + Schema = schema; + Columns = columns; + } + public static TaxiSchema Full(TaxiDataset dataset) { ParquetSchema schema = new( new DataField("VendorID"), @@ -52,7 +58,7 @@ public static TaxiSchema Full(TaxiDataset dataset) { new DataColumn(dataFields[18], dataset.Airport_fee) ]; - return new TaxiSchema(TaxiSchemaKind.Full, schema, columns); + return new TaxiSchema(schema, columns); } public static TaxiSchema Small(TaxiDataset dataset) { @@ -74,42 +80,6 @@ public static TaxiSchema Small(TaxiDataset dataset) { new DataColumn(dataFields[5], dataset.fare_amount) ]; - return new TaxiSchema(TaxiSchemaKind.Small, schema, columns); - } - - public void WriteParquetSharp(RowGroupWriter rowGroup, TaxiDataset dataset) { - static void WriteColumn(RowGroupWriter groupWriter, T[] data) { - using LogicalColumnWriter columnWriter = groupWriter.NextColumn().LogicalWriter(); - columnWriter.WriteBatch(data); - } - - if(Kind == TaxiSchemaKind.Small) { - WriteColumn(rowGroup, dataset.VendorID); - WriteColumn(rowGroup, dataset.passenger_count); - WriteColumn(rowGroup, dataset.trip_distance); - WriteColumn(rowGroup, dataset.RatecodeID); - WriteColumn(rowGroup, dataset.payment_type); - WriteColumn(rowGroup, dataset.fare_amount); - } else { - WriteColumn(rowGroup, dataset.VendorID); - WriteColumn(rowGroup, dataset.tpep_pickup_datetime); - WriteColumn(rowGroup, dataset.tpep_dropoff_datetime); - WriteColumn(rowGroup, dataset.passenger_count); - WriteColumn(rowGroup, dataset.trip_distance); - WriteColumn(rowGroup, dataset.RatecodeID); - WriteColumn(rowGroup, dataset.store_and_fwd_flag); - WriteColumn(rowGroup, dataset.PULocationID); - WriteColumn(rowGroup, dataset.DOLocationID); - WriteColumn(rowGroup, dataset.payment_type); - WriteColumn(rowGroup, dataset.fare_amount); - WriteColumn(rowGroup, dataset.extra); - WriteColumn(rowGroup, dataset.mta_tax); - WriteColumn(rowGroup, dataset.tip_amount); - WriteColumn(rowGroup, dataset.tolls_amount); - WriteColumn(rowGroup, dataset.improvement_surcharge); - WriteColumn(rowGroup, dataset.total_amount); - WriteColumn(rowGroup, dataset.congestion_surcharge); - WriteColumn(rowGroup, dataset.Airport_fee); - } + return new TaxiSchema(schema, columns); } } diff --git a/src/Parquet.PerfRunner/Taxis/TaxiSchemaKind.cs b/src/Parquet.PerfRunner/Taxis/TaxiSchemaKind.cs deleted file mode 100644 index 31d8e384..00000000 --- a/src/Parquet.PerfRunner/Taxis/TaxiSchemaKind.cs +++ /dev/null @@ -1,6 +0,0 @@ -namespace Parquet.PerfRunner.Taxis; - -enum TaxiSchemaKind { - Small, - Full -} From cfd714319be44d527a036763264bba3903845815 Mon Sep 17 00:00:00 2001 From: Kuinox Date: Sun, 14 Dec 2025 15:49:21 +0100 Subject: [PATCH 5/6] Fix self compare msbuild arg. --- src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs b/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs index 071569ed..f780707c 100644 --- a/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs +++ b/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs @@ -17,7 +17,7 @@ public NuConfig() { AddJob(baseJob); AddJob(baseJob .WithId("Parquet.Net 5.4.0") - .WithMsBuildArguments("/p:ParquetNuGetVersion=5.4.0")); + .WithMsBuildArguments("/p:ParquetVersion=5.4.0")); } } From 1f08ad5a08b1b69d3d2625a70b137da014dd9cd0 Mon Sep 17 00:00:00 2001 From: Kuinox Date: Tue, 5 May 2026 15:03:05 +0200 Subject: [PATCH 6/6] Compare self benchmark against previous release --- .../ParquetSharpComparisonBenchmark.cs | 18 +++-- .../Benchmarks/SelfComparisonBenchmark.cs | 13 ++-- src/Parquet.PerfRunner/LogicalEncoding.cs | 27 +++++++- .../Parquet.PerfRunner.csproj | 14 ++++ src/Parquet.PerfRunner/Program.cs | 4 ++ .../Taxis/TaxiDatasetLoader.cs | 30 ++++++++ src/Parquet.PerfRunner/Taxis/TaxiSchema.cs | 68 ++++++++++++++++++- 7 files changed, 159 insertions(+), 15 deletions(-) diff --git a/src/Parquet.PerfRunner/Benchmarks/ParquetSharpComparisonBenchmark.cs b/src/Parquet.PerfRunner/Benchmarks/ParquetSharpComparisonBenchmark.cs index 86f9ba6d..3f0a53b1 100644 --- a/src/Parquet.PerfRunner/Benchmarks/ParquetSharpComparisonBenchmark.cs +++ b/src/Parquet.PerfRunner/Benchmarks/ParquetSharpComparisonBenchmark.cs @@ -49,12 +49,15 @@ public async Task LoadDatasetAsync() { [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(); - foreach(TaxiColumn column in _parquetNetSchema.Columns) { - await column.WriteAsync(rowGroup); - } + await _parquetNetSchema.WriteAsync(rowGroup); } [Benchmark(Description = "Parquet.Net -> Disk")] @@ -63,12 +66,15 @@ public async Task ParquetNetToDiskAsync() { 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(); - foreach(TaxiColumn column in _parquetNetSchema.Columns) { - await column.WriteAsync(rowGroup); - } + await _parquetNetSchema.WriteAsync(rowGroup); } finally { IOFile.Delete(path); } diff --git a/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs b/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs index 32b452f7..92298e04 100644 --- a/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs +++ b/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs @@ -15,8 +15,8 @@ public NuConfig() { Job baseJob = Job.ShortRun.WithId("Parquet.Net (local)"); AddJob(baseJob); AddJob(baseJob - .WithId("Parquet.Net 6.0.0") - .WithMsBuildArguments("/p:ParquetVersion=6.0.0")); + .WithId("Parquet.Net 5.6.0") + .WithMsBuildArguments("/p:ParquetVersion=5.6.0")); } } @@ -39,11 +39,14 @@ public async Task LoadDatasetAsync() { [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(); - foreach(TaxiColumn column in _schema.Columns) { - await column.WriteAsync(rowGroup); - } + await _schema.WriteAsync(rowGroup); } } diff --git a/src/Parquet.PerfRunner/LogicalEncoding.cs b/src/Parquet.PerfRunner/LogicalEncoding.cs index b5e01d89..c7c21b4f 100644 --- a/src/Parquet.PerfRunner/LogicalEncoding.cs +++ b/src/Parquet.PerfRunner/LogicalEncoding.cs @@ -12,14 +12,38 @@ public enum LogicalEncoding { } 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(), + 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") @@ -40,4 +64,5 @@ public static ParquetOptions CreateOptions(this LogicalEncoding encoding, Parque return options; } +#endif } diff --git a/src/Parquet.PerfRunner/Parquet.PerfRunner.csproj b/src/Parquet.PerfRunner/Parquet.PerfRunner.csproj index e077a2be..3bf1d9df 100644 --- a/src/Parquet.PerfRunner/Parquet.PerfRunner.csproj +++ b/src/Parquet.PerfRunner/Parquet.PerfRunner.csproj @@ -27,4 +27,18 @@ + + $(DefineConstants);PARQUET_PACKAGE + + + + + + + + + + + + diff --git a/src/Parquet.PerfRunner/Program.cs b/src/Parquet.PerfRunner/Program.cs index 8248b6dd..5854b953 100644 --- a/src/Parquet.PerfRunner/Program.cs +++ b/src/Parquet.PerfRunner/Program.cs @@ -5,6 +5,7 @@ if(args.Length == 1) { switch(args[0]) { +#if !PARQUET_PACKAGE case "highLevel": HighLevel.Run(); break; @@ -22,6 +23,7 @@ case "compression": BenchmarkRunner.Run(); break; +#endif case "sharp-compare-taxi": BenchmarkRunner.Run(); break; @@ -30,6 +32,8 @@ break; } } else { +#if !PARQUET_PACKAGE await new DataTypes().RandomStrings(); +#endif //await SampleGenerator.GenerateFiles(); } diff --git a/src/Parquet.PerfRunner/Taxis/TaxiDatasetLoader.cs b/src/Parquet.PerfRunner/Taxis/TaxiDatasetLoader.cs index c986aab3..8ef3e7ff 100644 --- a/src/Parquet.PerfRunner/Taxis/TaxiDatasetLoader.cs +++ b/src/Parquet.PerfRunner/Taxis/TaxiDatasetLoader.cs @@ -1,6 +1,8 @@ using Parquet.Schema; using IOFile = System.IO.File; +#if !PARQUET_PACKAGE using Parquet.Data; +#endif using ParquetReaderNet = Parquet.ParquetReader; namespace Parquet.PerfRunner.Taxis; @@ -58,7 +60,11 @@ private static async Task LoadParquetColumnsAsync(IEnumerable LoadParquetColumnsAsync(IEnumerable(rg, vendor, vendorIds); await AddNullableValuesAsync(rg, pickup, pickupTimes); await AddNullableValuesAsync(rg, dropoff, dropoffTimes); @@ -102,6 +129,7 @@ private static async Task LoadParquetColumnsAsync(IEnumerable(rg, totalAmount, totalAmounts); await AddNullableValuesAsync(rg, congestionSurcharge, congestionSurcharges); await AddNullableValuesAsync(rg, airportFee, airportFees); +#endif } } @@ -131,6 +159,7 @@ private static async Task LoadParquetColumnsAsync(IEnumerable schema.GetDataFields().First(f => f.Name == v); +#if !PARQUET_PACKAGE private static async Task AddNullableValuesAsync( ParquetRowGroupReader rowGroup, DataField field, @@ -170,6 +199,7 @@ private static async Task AddNullableStringsAsync( destination.Add(definitionLevel == field.MaxDefinitionLevel ? new string(data.Values[valueIndex++].Span) : null); } } +#endif private static async Task PreloadFiles(IEnumerable files) => await Task.WhenAll( diff --git a/src/Parquet.PerfRunner/Taxis/TaxiSchema.cs b/src/Parquet.PerfRunner/Taxis/TaxiSchema.cs index ceb32979..9a897f96 100644 --- a/src/Parquet.PerfRunner/Taxis/TaxiSchema.cs +++ b/src/Parquet.PerfRunner/Taxis/TaxiSchema.cs @@ -1,3 +1,6 @@ +#if PARQUET_PACKAGE +using Parquet.Data; +#endif using Parquet.Schema; namespace Parquet.PerfRunner.Taxis; @@ -5,11 +8,33 @@ namespace Parquet.PerfRunner.Taxis; sealed class TaxiSchema { public ParquetSchema Schema { get; } - public TaxiColumn[] Columns { get; } + +#if PARQUET_PACKAGE + private readonly DataColumn[] _columns; + + private TaxiSchema(ParquetSchema schema, DataColumn[] columns) { + Schema = schema; + _columns = columns; + } +#else + private readonly TaxiColumn[] _columns; private TaxiSchema(ParquetSchema schema, TaxiColumn[] columns) { Schema = schema; - Columns = columns; + _columns = columns; + } +#endif + + public async Task WriteAsync(ParquetRowGroupWriter rowGroup) { +#if PARQUET_PACKAGE + foreach(DataColumn column in _columns) { + await rowGroup.WriteColumnAsync(column); + } +#else + foreach(TaxiColumn column in _columns) { + await column.WriteAsync(rowGroup); + } +#endif } public static TaxiSchema Full(TaxiDataset dataset) { @@ -35,6 +60,29 @@ public static TaxiSchema Full(TaxiDataset dataset) { new DataField("Airport_fee") ); DataField[] dataFields = schema.GetDataFields(); +#if PARQUET_PACKAGE + DataColumn[] columns = [ + new DataColumn(dataFields[0], dataset.VendorID), + new DataColumn(dataFields[1], dataset.tpep_pickup_datetime), + new DataColumn(dataFields[2], dataset.tpep_dropoff_datetime), + new DataColumn(dataFields[3], dataset.passenger_count), + new DataColumn(dataFields[4], dataset.trip_distance), + new DataColumn(dataFields[5], dataset.RatecodeID), + new DataColumn(dataFields[6], dataset.store_and_fwd_flag), + new DataColumn(dataFields[7], dataset.PULocationID), + new DataColumn(dataFields[8], dataset.DOLocationID), + new DataColumn(dataFields[9], dataset.payment_type), + new DataColumn(dataFields[10], dataset.fare_amount), + new DataColumn(dataFields[11], dataset.extra), + new DataColumn(dataFields[12], dataset.mta_tax), + new DataColumn(dataFields[13], dataset.tip_amount), + new DataColumn(dataFields[14], dataset.tolls_amount), + new DataColumn(dataFields[15], dataset.improvement_surcharge), + new DataColumn(dataFields[16], dataset.total_amount), + new DataColumn(dataFields[17], dataset.congestion_surcharge), + new DataColumn(dataFields[18], dataset.Airport_fee) + ]; +#else TaxiColumn[] columns = [ new(dataFields[0], rowGroup => rowGroup.WriteAsync(dataFields[0], dataset.VendorID)), new(dataFields[1], rowGroup => rowGroup.WriteAsync(dataFields[1], dataset.tpep_pickup_datetime)), @@ -56,6 +104,7 @@ public static TaxiSchema Full(TaxiDataset dataset) { new(dataFields[17], rowGroup => rowGroup.WriteAsync(dataFields[17], dataset.congestion_surcharge)), new(dataFields[18], rowGroup => rowGroup.WriteAsync(dataFields[18], dataset.Airport_fee)) ]; +#endif return new TaxiSchema(schema, columns); } @@ -70,6 +119,16 @@ public static TaxiSchema Small(TaxiDataset dataset) { new DataField("fare_amount") ); DataField[] dataFields = schema.GetDataFields(); +#if PARQUET_PACKAGE + DataColumn[] columns = [ + new DataColumn(dataFields[0], dataset.VendorID), + new DataColumn(dataFields[1], dataset.passenger_count), + new DataColumn(dataFields[2], dataset.trip_distance), + new DataColumn(dataFields[3], dataset.RatecodeID), + new DataColumn(dataFields[4], dataset.payment_type), + new DataColumn(dataFields[5], dataset.fare_amount) + ]; +#else TaxiColumn[] columns = [ new(dataFields[0], rowGroup => rowGroup.WriteAsync(dataFields[0], dataset.VendorID)), new(dataFields[1], rowGroup => rowGroup.WriteAsync(dataFields[1], dataset.passenger_count)), @@ -78,9 +137,12 @@ public static TaxiSchema Small(TaxiDataset dataset) { new(dataFields[4], rowGroup => rowGroup.WriteAsync(dataFields[4], dataset.payment_type)), new(dataFields[5], rowGroup => rowGroup.WriteAsync(dataFields[5], dataset.fare_amount)) ]; +#endif return new TaxiSchema(schema, columns); } } -public sealed record TaxiColumn(DataField Field, Func WriteAsync); +#if !PARQUET_PACKAGE +sealed record TaxiColumn(DataField Field, Func WriteAsync); +#endif