diff --git a/src/Parquet.PerfRunner/Benchmarks/DataTypes.cs b/src/Parquet.PerfRunner/Benchmarks/DataTypes.cs index 00bd22fd..4f13eb63 100644 --- a/src/Parquet.PerfRunner/Benchmarks/DataTypes.cs +++ b/src/Parquet.PerfRunner/Benchmarks/DataTypes.cs @@ -1,4 +1,3 @@ -using Parquet.Data; using Parquet.Schema; namespace Parquet.PerfRunner.Benchmarks; @@ -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); } } diff --git a/src/Parquet.PerfRunner/Benchmarks/ParquetSharpComparisonBenchmark.cs b/src/Parquet.PerfRunner/Benchmarks/ParquetSharpComparisonBenchmark.cs new file mode 100644 index 00000000..3f0a53b1 --- /dev/null +++ b/src/Parquet.PerfRunner/Benchmarks/ParquetSharpComparisonBenchmark.cs @@ -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); + } + } +} diff --git a/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs b/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs new file mode 100644 index 00000000..92298e04 --- /dev/null +++ b/src/Parquet.PerfRunner/Benchmarks/SelfComparisonBenchmark.cs @@ -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); + } +} diff --git a/src/Parquet.PerfRunner/Benchmarks/WriteBenchmark.cs b/src/Parquet.PerfRunner/Benchmarks/WriteBenchmark.cs index 1b1f26d8..9384d197 100644 --- a/src/Parquet.PerfRunner/Benchmarks/WriteBenchmark.cs +++ b/src/Parquet.PerfRunner/Benchmarks/WriteBenchmark.cs @@ -1,4 +1,4 @@ -using BenchmarkDotNet.Attributes; +using BenchmarkDotNet.Attributes; using Parquet.Schema; using ParquetSharp; using ParquetSharp.IO; diff --git a/src/Parquet.PerfRunner/LogicalEncoding.cs b/src/Parquet.PerfRunner/LogicalEncoding.cs new file mode 100644 index 00000000..c7c21b4f --- /dev/null +++ b/src/Parquet.PerfRunner/LogicalEncoding.cs @@ -0,0 +1,68 @@ +using Parquet.Schema; + +namespace Parquet.PerfRunner; + +/// +/// Allow to easily control the logical encoding used in benchmarks. +/// +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 +} diff --git a/src/Parquet.PerfRunner/Parquet.PerfRunner.csproj b/src/Parquet.PerfRunner/Parquet.PerfRunner.csproj index deb0a62b..3bf1d9df 100644 --- a/src/Parquet.PerfRunner/Parquet.PerfRunner.csproj +++ b/src/Parquet.PerfRunner/Parquet.PerfRunner.csproj @@ -27,4 +27,18 @@ - \ No newline at end of file + + $(DefineConstants);PARQUET_PACKAGE + + + + + + + + + + + + + diff --git a/src/Parquet.PerfRunner/Program.cs b/src/Parquet.PerfRunner/Program.cs index a80969ff..5854b953 100644 --- a/src/Parquet.PerfRunner/Program.cs +++ b/src/Parquet.PerfRunner/Program.cs @@ -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; @@ -24,8 +23,17 @@ case "compression": BenchmarkRunner.Run(); break; +#endif + case "sharp-compare-taxi": + BenchmarkRunner.Run(); + break; + case "self-compare-taxi": + BenchmarkRunner.Run(); + break; } } else { +#if !PARQUET_PACKAGE await new DataTypes().RandomStrings(); +#endif //await SampleGenerator.GenerateFiles(); } diff --git a/src/Parquet.PerfRunner/Taxis/ParquetSharpTaxiSchema.cs b/src/Parquet.PerfRunner/Taxis/ParquetSharpTaxiSchema.cs new file mode 100644 index 00000000..2c290364 --- /dev/null +++ b/src/Parquet.PerfRunner/Taxis/ParquetSharpTaxiSchema.cs @@ -0,0 +1,125 @@ +using ParquetSharp; +using Column = ParquetSharp.Column; + +namespace Parquet.PerfRunner.Taxis; + +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"), + 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, Kind.Full); + } + + 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, Kind.Small); + } + + 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(); + } + + 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/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..8ef3e7ff --- /dev/null +++ b/src/Parquet.PerfRunner/Taxis/TaxiDatasetLoader.cs @@ -0,0 +1,219 @@ +using Parquet.Schema; +using IOFile = System.IO.File; +#if !PARQUET_PACKAGE +using Parquet.Data; +#endif +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)); +#if PARQUET_PACKAGE + using ParquetReaderNet reader = await ParquetReaderNet.CreateAsync(fs); +#else + await using ParquetReaderNet reader = await ParquetReaderNet.CreateAsync(fs); +#endif + + 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); + +#if PARQUET_PACKAGE + 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); +#else + await AddNullableValuesAsync(rg, vendor, vendorIds); + await AddNullableValuesAsync(rg, pickup, pickupTimes); + await AddNullableValuesAsync(rg, dropoff, dropoffTimes); + await AddNullableValuesAsync(rg, passengerCount, passengerCounts); + await AddNullableValuesAsync(rg, tripDistance, tripDistances); + await AddNullableValuesAsync(rg, rateCode, rateCodes); + await AddNullableStringsAsync(rg, storeAndFwdFlag, storeAndFwdFlags); + await AddNullableValuesAsync(rg, puLocationId, puLocationIds); + await AddNullableValuesAsync(rg, doLocationId, doLocationIds); + await AddNullableValuesAsync(rg, paymentType, paymentTypes); + await AddNullableValuesAsync(rg, fareAmount, fareAmounts); + await AddNullableValuesAsync(rg, extra, extras); + await AddNullableValuesAsync(rg, mtaTax, mtaTaxes); + await AddNullableValuesAsync(rg, tipAmount, tipAmounts); + await AddNullableValuesAsync(rg, tollsAmount, tollsAmounts); + await AddNullableValuesAsync(rg, improvementSurcharge, improvementSurcharges); + await AddNullableValuesAsync(rg, totalAmount, totalAmounts); + await AddNullableValuesAsync(rg, congestionSurcharge, congestionSurcharges); + await AddNullableValuesAsync(rg, airportFee, airportFees); +#endif + + } + } + + 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); + +#if !PARQUET_PACKAGE + private static async Task AddNullableValuesAsync( + ParquetRowGroupReader rowGroup, + DataField field, + ICollection destination) where T : struct { + using RawColumnData rawData = await rowGroup.ReadRawColumnDataBaseAsync(field); + RawColumnData data = (RawColumnData)rawData; + + if(field.MaxDefinitionLevel == 0) { + foreach(T value in data.Values) { + destination.Add(value); + } + return; + } + + int valueIndex = 0; + foreach(int definitionLevel in data.DefinitionLevels) { + destination.Add(definitionLevel == field.MaxDefinitionLevel ? data.Values[valueIndex++] : null); + } + } + + private static async Task AddNullableStringsAsync( + ParquetRowGroupReader rowGroup, + DataField field, + ICollection destination) { + using RawColumnData rawData = await rowGroup.ReadRawColumnDataBaseAsync(field); + RawColumnData> data = (RawColumnData>)rawData; + + if(field.MaxDefinitionLevel == 0) { + foreach(ReadOnlyMemory value in data.Values) { + destination.Add(new string(value.Span)); + } + return; + } + + int valueIndex = 0; + foreach(int definitionLevel in data.DefinitionLevels) { + destination.Add(definitionLevel == field.MaxDefinitionLevel ? new string(data.Values[valueIndex++].Span) : null); + } + } +#endif + + 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..9a897f96 --- /dev/null +++ b/src/Parquet.PerfRunner/Taxis/TaxiSchema.cs @@ -0,0 +1,148 @@ +#if PARQUET_PACKAGE +using Parquet.Data; +#endif +using Parquet.Schema; + +namespace Parquet.PerfRunner.Taxis; + +sealed class TaxiSchema { + + public ParquetSchema Schema { 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; + } +#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) { + 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(); +#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)), + new(dataFields[2], rowGroup => rowGroup.WriteAsync(dataFields[2], dataset.tpep_dropoff_datetime)), + new(dataFields[3], rowGroup => rowGroup.WriteAsync(dataFields[3], dataset.passenger_count)), + new(dataFields[4], rowGroup => rowGroup.WriteAsync(dataFields[4], dataset.trip_distance)), + new(dataFields[5], rowGroup => rowGroup.WriteAsync(dataFields[5], dataset.RatecodeID)), + new(dataFields[6], rowGroup => rowGroup.WriteAsync(dataFields[6], dataset.store_and_fwd_flag)), + new(dataFields[7], rowGroup => rowGroup.WriteAsync(dataFields[7], dataset.PULocationID)), + new(dataFields[8], rowGroup => rowGroup.WriteAsync(dataFields[8], dataset.DOLocationID)), + new(dataFields[9], rowGroup => rowGroup.WriteAsync(dataFields[9], dataset.payment_type)), + new(dataFields[10], rowGroup => rowGroup.WriteAsync(dataFields[10], dataset.fare_amount)), + new(dataFields[11], rowGroup => rowGroup.WriteAsync(dataFields[11], dataset.extra)), + new(dataFields[12], rowGroup => rowGroup.WriteAsync(dataFields[12], dataset.mta_tax)), + new(dataFields[13], rowGroup => rowGroup.WriteAsync(dataFields[13], dataset.tip_amount)), + new(dataFields[14], rowGroup => rowGroup.WriteAsync(dataFields[14], dataset.tolls_amount)), + new(dataFields[15], rowGroup => rowGroup.WriteAsync(dataFields[15], dataset.improvement_surcharge)), + new(dataFields[16], rowGroup => rowGroup.WriteAsync(dataFields[16], dataset.total_amount)), + 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); + } + + 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(); +#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)), + new(dataFields[2], rowGroup => rowGroup.WriteAsync(dataFields[2], dataset.trip_distance)), + new(dataFields[3], rowGroup => rowGroup.WriteAsync(dataFields[3], dataset.RatecodeID)), + 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); + } +} + +#if !PARQUET_PACKAGE +sealed record TaxiColumn(DataField Field, Func WriteAsync); +#endif