Dataset writer layer

The dataset writer sits above the row write layer. It routes rows to multiple Parquet files.

Use it to write a partitioned dataset when rows belonging to different files are mixed together.

It can write any number of output files while keeping only a fixed number open at a time.

Each output file follows the same row-group and rollover targets as the row writer.

The examples use EventSchema and the FileParquetSource adapter below.

File sources

Implement IParquetReadWriteSource to open and reopen dataset files. This example uses local files:

// A reusable local-file adapter for dataset writing. One instance owns one open file.
sealed class FileParquetSource : IParquetReadWriteSource
{
    FileStream? _stream;

    FileStream Stream => _stream ?? throw new InvalidOperationException("The file is not open.");

    public ulong Length => checked((ulong)Stream.Length);

    public void Open(ReadOnlySpan<byte> path, FileMode mode)
    {
        if (_stream is not null)
            throw new InvalidOperationException("The previous file must be closed before this source is reused.");
        var filePath = Encoding.UTF8.GetString(path);
        Directory.CreateDirectory(Path.GetDirectoryName(Path.GetFullPath(filePath))!);
        _stream = new FileStream(filePath, mode, FileAccess.ReadWrite, FileShare.None);
    }

    public void ReadExactly(ulong offset, Span<byte> destination)
    {
        Stream.Position = checked((long)offset);
        Stream.ReadExactly(destination);
    }

    public void Write(ulong offset, ReadOnlySpan<byte> source)
    {
        Stream.Position = checked((long)offset);
        Stream.Write(source);
    }

    public void SetLength(ulong length) => Stream.SetLength(checked((long)length));

    public void Flush() => Stream.Flush();

    public void Close()
    {
        _stream?.Dispose();
        _stream = null;
    }

    public void Dispose() => Close();
}

Dispose the writer before disposing its sources.

Route rows

The route returns the UTF-8 path that should receive each row:

using var file = new FileParquetSource();
IParquetReadWriteSource[] files = [file];

using var writer = EventSchema.CreateDatasetWriter(
    static (EventSchema row, IParquetBufferPool pool, out ParquetBuffer? allocation) =>
    {
        allocation = null;
        return row.Id % 2 == 0
            ? "events/even.parquet"u8
            : "events/odd.parquet"u8;
    },
    files);

for (var id = 0; id < 6; id++)
    writer.Queue(new EventSchema { Id = id, Name = "event"u8.ToArray(), OccurredAt = DateTimeOffset.UtcNow });
// Disposing the writer flushes all remaining rows and closes the output files.

files contains the reusable read/write sources. Its length is the maximum number of files kept open.

Queue() copies the row into the writer buffers. Disposing the writer writes the remaining rows and closes every open file.

Plank appends to existing files. Use a new directory to start a fresh dataset.

Build paths at runtime

Static UTF-8 paths do not need an allocation. For paths built at runtime, use the provided buffer pool:

using var file = new FileParquetSource();
IParquetReadWriteSource[] files = [file];

using var writer = EventSchema.CreateDatasetWriter(
    static (EventSchema row, IParquetBufferPool pool, out ParquetBuffer? allocation) =>
    {
        var path = $"events/bucket={row.Id % 16}.parquet";
        var buffer = pool.Rent(checked((uint)Encoding.UTF8.GetByteCount(path)));
        var length = Encoding.UTF8.GetBytes(path, buffer.Span);
        allocation = buffer;
        return buffer.Span[..length];
    },
    files);

for (var id = 0; id < 6; id++)
    writer.Queue(new EventSchema { Id = id, OccurredAt = DateTimeOffset.UtcNow });

Plank releases the returned allocation when it no longer needs the path.