Class PipelineRowWriterBase<TSlot>
Coordinates buffered row batches for a generated pipeline row writer.
public abstract class PipelineRowWriterBase<TSlot> : RowWriterBase<TSlot>, IDisposable where TSlot : RowBufferSlot
Type Parameters
TSlotThe generated buffer-slot type.
- Inheritance
-
RowWriterBase<TSlot>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
fileIParquetWriteSourceThe reusable destination used for each produced file.
filePathParquetFilePathSelects the path of each produced file.
schemaParquetSchemaThe generated Parquet schema.
maxParallelismuintThe maximum number of serialization workers.
onFlushAction<int>An optional callback invoked with each flushed row count.
optionsParquetWriterOptionsThe Parquet writer options.
rowBatchSizeintThe initial row capacity of each generated buffer slot.
workerThreadNamePrefixstringThe 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
streamStreamThe destination stream.
schemaParquetSchemaThe generated Parquet schema.
maxParallelismuintThe maximum number of serialization workers.
onFlushAction<int>An optional callback invoked with each flushed row count.
optionsParquetWriterOptionsThe Parquet writer options.
rowBatchSizeintThe number of rows in each generated buffer slot.
workerThreadNamePrefixstringThe worker-thread name prefix.
Properties
BufferGeneration
Identifies the current writable buffers for generated cursor invalidation.
protected long BufferGeneration { get; }
Property Value
RowBatchSize
Gets the generated writer's row-batch size.
protected int RowBatchSize { get; }
Property Value
WorkerThreadNamePrefix
Gets the name prefix used for serialization worker threads.
protected override string WorkerThreadNamePrefix { get; }
Property Value
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
slotTSlotThe current active slot.
rowsPerGroupintThe 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
slotTSlotThe current active slot.
rowSizeBytesulongThe 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
fixedRowSizeBytesulongThe 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
slotTSlotThe slot that was written.
ResetSlotForReuse(TSlot)
Resets a generated slot before it is reused.
protected override void ResetSlotForReuse(TSlot slot)
Parameters
slotTSlotThe slot to reset.
ResetWriter(Stream)
Resets a completed generated writer to a new destination stream.
protected void ResetWriter(Stream stream)
Parameters
streamStreamThe new destination stream.
SerializeSlot(TSlot)
Serializes a generated row-buffer slot.
protected override void SerializeSlot(TSlot slot)
Parameters
slotTSlotThe slot to serialize.
WriteSerializedSlot(TSlot, RowGroupWriter)
Writes a serialized slot to a row group.
protected override void WriteSerializedSlot(TSlot slot, RowGroupWriter rowGroupWriter)
Parameters
slotTSlotThe serialized slot.
rowGroupWriterRowGroupWriterThe destination row-group writer.