Table of Contents

Class PipelineRowWriterBase<TSlot>

Namespace
Plank.RowApi
Assembly
Plank.dll

Coordinates buffered row batches for a generated pipeline row writer.

public abstract class PipelineRowWriterBase<TSlot> : RowWriterBase<TSlot>, IDisposable where TSlot : RowBufferSlot

Type Parameters

TSlot

The generated buffer-slot type.

Inheritance
PipelineRowWriterBase<TSlot>
Implements
Inherited Members

Remarks

This unstable API supports Plank-generated code and is not intended for direct use by applications.

Constructors

PipelineRowWriterBase(IParquetWriteSource, ParquetFilePath, ParquetSchema, uint, Action<int>?, ParquetWriterOptions, int, string)

Initializes a rolling generated pipeline row writer.

protected PipelineRowWriterBase(IParquetWriteSource file, ParquetFilePath filePath, ParquetSchema schema, uint maxParallelism, Action<int>? onFlush, ParquetWriterOptions options, int rowBatchSize, string workerThreadNamePrefix)

Parameters

file IParquetWriteSource

The reusable destination used for each produced file.

filePath ParquetFilePath

Selects the path of each produced file.

schema ParquetSchema

The generated Parquet schema.

maxParallelism uint

The maximum number of serialization workers.

onFlush Action<int>

An optional callback invoked with each flushed row count.

options ParquetWriterOptions

The Parquet writer options.

rowBatchSize int

The initial row capacity of each generated buffer slot.

workerThreadNamePrefix string

The worker-thread name prefix.

PipelineRowWriterBase(Stream, ParquetSchema, uint, Action<int>?, ParquetWriterOptions, int, string)

Initializes the infrastructure for a generated pipeline row writer.

protected PipelineRowWriterBase(Stream stream, ParquetSchema schema, uint maxParallelism, Action<int>? onFlush, ParquetWriterOptions options, int rowBatchSize, string workerThreadNamePrefix)

Parameters

stream Stream

The destination stream.

schema ParquetSchema

The generated Parquet schema.

maxParallelism uint

The maximum number of serialization workers.

onFlush Action<int>

An optional callback invoked with each flushed row count.

options ParquetWriterOptions

The Parquet writer options.

rowBatchSize int

The number of rows in each generated buffer slot.

workerThreadNamePrefix string

The worker-thread name prefix.

Properties

BufferGeneration

Identifies the current writable buffers for generated cursor invalidation.

protected long BufferGeneration { get; }

Property Value

long

RowBatchSize

Gets the generated writer's row-batch size.

protected int RowBatchSize { get; }

Property Value

int

WorkerThreadNamePrefix

Gets the name prefix used for serialization worker threads.

protected override string WorkerThreadNamePrefix { get; }

Property Value

string

Methods

CommitFixedRow(TSlot, int)

Commits an already validated fixed-width row and returns the slot prepared for the next row.

protected TSlot CommitFixedRow(TSlot slot, int rowsPerGroup)

Parameters

slot TSlot

The current active slot.

rowsPerGroup int

The generated row-count cutoff for one row group.

Returns

TSlot

The slot prepared for the next row.

CommitVariableRow(TSlot, ulong)

Commits an already validated variable-width row and returns the slot prepared for the next row.

protected TSlot CommitVariableRow(TSlot slot, ulong rowSizeBytes)

Parameters

slot TSlot

The current active slot.

rowSizeBytes ulong

The generated estimate of the current row's buffered size.

Returns

TSlot

The slot prepared for the next row.

CompleteWriter()

Flushes pending rows and completes the generated writer.

protected void CompleteWriter()

GetFixedRowsPerGroup(ulong)

Calculates the generated row-count cutoff for a fixed-width schema.

protected int GetFixedRowsPerGroup(ulong fixedRowSizeBytes)

Parameters

fixedRowSizeBytes ulong

The generated fixed size of one row.

Returns

int

The number of rows to buffer before flushing.

GetSlotForRow()

Gets the current buffer slot for generated row assignment.

protected TSlot GetSlotForRow()

Returns

TSlot

The current writable slot.

OnSlotWritten(TSlot)

Handles successful writing of a generated slot.

protected override void OnSlotWritten(TSlot slot)

Parameters

slot TSlot

The slot that was written.

ResetSlotForReuse(TSlot)

Resets a generated slot before it is reused.

protected override void ResetSlotForReuse(TSlot slot)

Parameters

slot TSlot

The slot to reset.

ResetWriter(Stream)

Resets a completed generated writer to a new destination stream.

protected void ResetWriter(Stream stream)

Parameters

stream Stream

The new destination stream.

SerializeSlot(TSlot)

Serializes a generated row-buffer slot.

protected override void SerializeSlot(TSlot slot)

Parameters

slot TSlot

The slot to serialize.

WriteSerializedSlot(TSlot, RowGroupWriter)

Writes a serialized slot to a row group.

protected override void WriteSerializedSlot(TSlot slot, RowGroupWriter rowGroupWriter)

Parameters

slot TSlot

The serialized slot.

rowGroupWriter RowGroupWriter

The destination row-group writer.