package parquet
import (
"log"
"reflect"
"runtime"
"sort"
"sync"
"sync/atomic"
"github.com/parquet-go/parquet-go/internal/debug"
)
type GenericBuffer [T any ] struct {
base Buffer
write bufferFunc [T ]
}
func NewGenericBuffer [T any ](options ...RowGroupOption ) *GenericBuffer [T ] {
config , err := NewRowGroupConfig (options ...)
if err != nil {
panic (err )
}
t := typeOf [T ]()
if config .Schema == nil && t != nil {
config .Schema = schemaOf (dereference (t ))
}
if config .Schema == nil {
panic ("generic buffer must be instantiated with schema or concrete type." )
}
buf := &GenericBuffer [T ]{
base : Buffer {config : config },
}
buf .base .configure (config .Schema )
buf .write = bufferFuncOf [T ](t , config .Schema )
return buf
}
func typeOf[T any ]() reflect .Type {
var v T
return reflect .TypeOf (v )
}
type bufferFunc[T any ] func (*GenericBuffer [T ], []T ) (int , error )
func bufferFuncOf[T any ](t reflect .Type , schema *Schema ) bufferFunc [T ] {
if t == nil {
return (*GenericBuffer [T ]).writeRows
}
switch t .Kind () {
case reflect .Interface , reflect .Map :
return (*GenericBuffer [T ]).writeRows
case reflect .Struct :
return makeBufferFunc [T ](t , schema )
case reflect .Pointer :
if e := t .Elem (); e .Kind () == reflect .Struct {
return makeBufferFunc [T ](t , schema )
}
}
panic ("cannot create buffer for values of type " + t .String ())
}
func makeBufferFunc[T any ](t reflect .Type , schema *Schema ) bufferFunc [T ] {
writeRows := writeRowsFuncOf (t , schema , nil )
return func (buf *GenericBuffer [T ], rows []T ) (n int , err error ) {
err = writeRows (buf .base .columns , makeArrayOf (rows ), columnLevels {})
if err == nil {
n = len (rows )
}
return n , err
}
}
func (buf *GenericBuffer [T ]) Size () int64 {
return buf .base .Size ()
}
func (buf *GenericBuffer [T ]) NumRows () int64 {
return buf .base .NumRows ()
}
func (buf *GenericBuffer [T ]) ColumnChunks () []ColumnChunk {
return buf .base .ColumnChunks ()
}
func (buf *GenericBuffer [T ]) ColumnBuffers () []ColumnBuffer {
return buf .base .ColumnBuffers ()
}
func (buf *GenericBuffer [T ]) SortingColumns () []SortingColumn {
return buf .base .SortingColumns ()
}
func (buf *GenericBuffer [T ]) Len () int {
return buf .base .Len ()
}
func (buf *GenericBuffer [T ]) Less (i , j int ) bool {
return buf .base .Less (i , j )
}
func (buf *GenericBuffer [T ]) Swap (i , j int ) {
buf .base .Swap (i , j )
}
func (buf *GenericBuffer [T ]) Reset () {
buf .base .Reset ()
}
func (buf *GenericBuffer [T ]) Write (rows []T ) (int , error ) {
if len (rows ) == 0 {
return 0 , nil
}
return buf .write (buf , rows )
}
func (buf *GenericBuffer [T ]) WriteRows (rows []Row ) (int , error ) {
return buf .base .WriteRows (rows )
}
func (buf *GenericBuffer [T ]) WriteRowGroup (rowGroup RowGroup ) (int64 , error ) {
return buf .base .WriteRowGroup (rowGroup )
}
func (buf *GenericBuffer [T ]) Rows () Rows {
return buf .base .Rows ()
}
func (buf *GenericBuffer [T ]) Schema () *Schema {
return buf .base .Schema ()
}
func (buf *GenericBuffer [T ]) writeRows (rows []T ) (int , error ) {
if cap (buf .base .rowbuf ) < len (rows ) {
buf .base .rowbuf = make ([]Row , len (rows ))
} else {
buf .base .rowbuf = buf .base .rowbuf [:len (rows )]
}
defer clearRows (buf .base .rowbuf )
schema := buf .base .Schema ()
for i := range rows {
buf .base .rowbuf [i ] = schema .Deconstruct (buf .base .rowbuf [i ], &rows [i ])
}
return buf .base .WriteRows (buf .base .rowbuf )
}
var (
_ RowGroup = (*GenericBuffer [any ])(nil )
_ RowGroupWriter = (*GenericBuffer [any ])(nil )
_ sort .Interface = (*GenericBuffer [any ])(nil )
_ RowGroup = (*GenericBuffer [struct {}])(nil )
_ RowGroupWriter = (*GenericBuffer [struct {}])(nil )
_ sort .Interface = (*GenericBuffer [struct {}])(nil )
_ RowGroup = (*GenericBuffer [map [struct {}]struct {}])(nil )
_ RowGroupWriter = (*GenericBuffer [map [struct {}]struct {}])(nil )
_ sort .Interface = (*GenericBuffer [map [struct {}]struct {}])(nil )
)
type Buffer struct {
config *RowGroupConfig
schema *Schema
rowbuf []Row
colbuf [][]Value
chunks []ColumnChunk
columns []ColumnBuffer
sorted []ColumnBuffer
}
func NewBuffer (options ...RowGroupOption ) *Buffer {
config , err := NewRowGroupConfig (options ...)
if err != nil {
panic (err )
}
buf := &Buffer {
config : config ,
}
if config .Schema != nil {
buf .configure (config .Schema )
}
return buf
}
func (buf *Buffer ) configure (schema *Schema ) {
if schema == nil {
return
}
sortingColumns := buf .config .Sorting .SortingColumns
buf .sorted = make ([]ColumnBuffer , len (sortingColumns ))
forEachLeafColumnOf (schema , func (leaf leafColumn ) {
nullOrdering := nullsGoLast
columnIndex := int (leaf .columnIndex )
columnType := leaf .node .Type ()
bufferCap := buf .config .ColumnBufferCapacity
dictionary := (Dictionary )(nil )
encoding := encodingOf (leaf .node )
if isDictionaryEncoding (encoding ) {
estimatedDictBufferSize := columnType .EstimateSize (bufferCap )
dictBuffer := columnType .NewValues (
make ([]byte , 0 , estimatedDictBufferSize ),
nil ,
)
dictionary = columnType .NewDictionary (columnIndex , 0 , dictBuffer )
columnType = dictionary .Type ()
}
sortingIndex := searchSortingColumn (sortingColumns , leaf .path )
if sortingIndex < len (sortingColumns ) && sortingColumns [sortingIndex ].NullsFirst () {
nullOrdering = nullsGoFirst
}
column := columnType .NewColumnBuffer (columnIndex , bufferCap )
switch {
case leaf .maxRepetitionLevel > 0 :
column = newRepeatedColumnBuffer (column , leaf .maxRepetitionLevel , leaf .maxDefinitionLevel , nullOrdering )
case leaf .maxDefinitionLevel > 0 :
column = newOptionalColumnBuffer (column , leaf .maxDefinitionLevel , nullOrdering )
}
buf .columns = append (buf .columns , column )
if sortingIndex < len (sortingColumns ) {
if sortingColumns [sortingIndex ].Descending () {
column = &reversedColumnBuffer {column }
}
buf .sorted [sortingIndex ] = column
}
})
buf .schema = schema
buf .rowbuf = make ([]Row , 0 , 1 )
buf .colbuf = make ([][]Value , len (buf .columns ))
buf .chunks = make ([]ColumnChunk , len (buf .columns ))
for i , column := range buf .columns {
buf .chunks [i ] = column
}
}
func (buf *Buffer ) Size () int64 {
size := int64 (0 )
for _ , col := range buf .columns {
size += col .Size ()
}
return size
}
func (buf *Buffer ) NumRows () int64 { return int64 (buf .Len ()) }
func (buf *Buffer ) ColumnChunks () []ColumnChunk { return buf .chunks }
func (buf *Buffer ) ColumnBuffers () []ColumnBuffer { return buf .columns }
func (buf *Buffer ) Schema () *Schema { return buf .schema }
func (buf *Buffer ) SortingColumns () []SortingColumn { return buf .config .Sorting .SortingColumns }
func (buf *Buffer ) Len () int {
if len (buf .columns ) == 0 {
return 0
} else {
return buf .columns [0 ].Len ()
}
}
func (buf *Buffer ) Less (i , j int ) bool {
for _ , col := range buf .sorted {
switch {
case col .Less (i , j ):
return true
case col .Less (j , i ):
return false
}
}
return false
}
func (buf *Buffer ) Swap (i , j int ) {
for _ , col := range buf .columns {
col .Swap (i , j )
}
}
func (buf *Buffer ) Reset () {
for _ , col := range buf .columns {
col .Reset ()
}
}
func (buf *Buffer ) Write (row interface {}) error {
if buf .schema == nil {
buf .configure (SchemaOf (row ))
}
buf .rowbuf = buf .rowbuf [:1 ]
defer clearRows (buf .rowbuf )
buf .rowbuf [0 ] = buf .schema .Deconstruct (buf .rowbuf [0 ], row )
_ , err := buf .WriteRows (buf .rowbuf )
return err
}
func (buf *Buffer ) WriteRows (rows []Row ) (int , error ) {
defer func () {
for i , colbuf := range buf .colbuf {
clearValues (colbuf )
buf .colbuf [i ] = colbuf [:0 ]
}
}()
if buf .schema == nil {
return 0 , ErrRowGroupSchemaMissing
}
for _ , row := range rows {
for _ , value := range row {
columnIndex := value .Column ()
buf .colbuf [columnIndex ] = append (buf .colbuf [columnIndex ], value )
}
}
for columnIndex , values := range buf .colbuf {
if _ , err := buf .columns [columnIndex ].WriteValues (values ); err != nil {
return 0 , err
}
}
return len (rows ), nil
}
func (buf *Buffer ) WriteRowGroup (rowGroup RowGroup ) (int64 , error ) {
rowGroupSchema := rowGroup .Schema ()
switch {
case rowGroupSchema == nil :
return 0 , ErrRowGroupSchemaMissing
case buf .schema == nil :
buf .configure (rowGroupSchema )
case !nodesAreEqual (buf .schema , rowGroupSchema ):
return 0 , ErrRowGroupSchemaMismatch
}
if !sortingColumnsHavePrefix (rowGroup .SortingColumns (), buf .SortingColumns ()) {
return 0 , ErrRowGroupSortingColumnsMismatch
}
n := buf .NumRows ()
r := rowGroup .Rows ()
defer r .Close ()
_ , err := CopyRows (bufferWriter {buf }, r )
return buf .NumRows () - n , err
}
func (buf *Buffer ) Rows () Rows { return newRowGroupRows (buf , ReadModeSync ) }
type bufferWriter struct { buf *Buffer }
func (w bufferWriter ) WriteRows (rows []Row ) (int , error ) {
return w .buf .WriteRows (rows )
}
func (w bufferWriter ) WriteValues (values []Value ) (int , error ) {
return w .buf .columns [values [0 ].Column ()].WriteValues (values )
}
func (w bufferWriter ) WritePage (page Page ) (int64 , error ) {
return CopyValues (w .buf .columns [page .Column ()], page .Values ())
}
var (
_ RowGroup = (*Buffer )(nil )
_ RowGroupWriter = (*Buffer )(nil )
_ sort .Interface = (*Buffer )(nil )
_ RowWriter = (*bufferWriter )(nil )
_ PageWriter = (*bufferWriter )(nil )
_ ValueWriter = (*bufferWriter )(nil )
)
type buffer struct {
data []byte
refc uintptr
pool *bufferPool
stack []byte
}
func (b *buffer ) refCount () int {
return int (atomic .LoadUintptr (&b .refc ))
}
func (b *buffer ) ref () {
atomic .AddUintptr (&b .refc , +1 )
}
func (b *buffer ) unref () {
if atomic .AddUintptr (&b .refc , ^uintptr (0 )) == 0 {
if b .pool != nil {
b .pool .put (b )
}
}
}
func monitorBufferRelease(b *buffer ) {
if rc := b .refCount (); rc != 0 {
log .Printf ("PARQUETGODEBUG: buffer garbage collected with non-zero reference count\n%s" , string (b .stack ))
}
}
type bufferPool struct {
buckets [bufferPoolBucketCount ]sync .Pool
}
func (p *bufferPool ) newBuffer (bufferSize , bucketSize int ) *buffer {
b := &buffer {
data : make ([]byte , bufferSize , bucketSize ),
refc : 1 ,
pool : p ,
}
if debug .TRACEBUF > 0 {
b .stack = make ([]byte , 4096 )
runtime .SetFinalizer (b , monitorBufferRelease )
}
return b
}
func (p *bufferPool ) get (bufferSize int ) *buffer {
bucketIndex , bucketSize := bufferPoolBucketIndexAndSizeOfGet (bufferSize )
b := (*buffer )(nil )
if bucketIndex >= 0 {
b , _ = p .buckets [bucketIndex ].Get ().(*buffer )
}
if b == nil {
b = p .newBuffer (bufferSize , bucketSize )
} else {
b .data = b .data [:bufferSize ]
b .ref ()
}
if debug .TRACEBUF > 0 {
b .stack = b .stack [:runtime .Stack (b .stack [:cap (b .stack )], false )]
}
return b
}
func (p *bufferPool ) put (b *buffer ) {
if b .pool != p {
panic ("BUG: buffer returned to a different pool than the one it was allocated from" )
}
if b .refCount () != 0 {
panic ("BUG: buffer returned to pool with a non-zero reference count" )
}
if bucketIndex , _ := bufferPoolBucketIndexAndSizeOfPut (cap (b .data )); bucketIndex >= 0 {
p .buckets [bucketIndex ].Put (b )
}
}
const (
bufferPoolBucketCount = 32
bufferPoolMinSize = 4096
bufferPoolLastShortBucketSize = 262144
)
func bufferPoolNextSize(size int ) int {
if size < bufferPoolLastShortBucketSize {
return size * 2
} else {
return size + (size / 2 )
}
}
func bufferPoolBucketIndexAndSizeOfGet(size int ) (int , int ) {
limit := bufferPoolMinSize
for i := 0 ; i < bufferPoolBucketCount ; i ++ {
if size <= limit {
return i , limit
}
limit = bufferPoolNextSize (limit )
}
return -1 , size
}
func bufferPoolBucketIndexAndSizeOfPut(size int ) (int , int ) {
if limit := bufferPoolMinSize ; size >= limit {
for i := 0 ; i < bufferPoolBucketCount ; i ++ {
n := bufferPoolNextSize (limit )
if size < n {
return i , limit
}
limit = n
}
}
return -1 , size
}
var (
buffers bufferPool
)
type bufferedPage struct {
Page
values *buffer
offsets *buffer
repetitionLevels *buffer
definitionLevels *buffer
}
func newBufferedPage(page Page , values , offsets , definitionLevels , repetitionLevels *buffer ) *bufferedPage {
p := &bufferedPage {
Page : page ,
values : values ,
offsets : offsets ,
definitionLevels : definitionLevels ,
repetitionLevels : repetitionLevels ,
}
bufferRef (values )
bufferRef (offsets )
bufferRef (definitionLevels )
bufferRef (repetitionLevels )
return p
}
func (p *bufferedPage ) Slice (i , j int64 ) Page {
return newBufferedPage (
p .Page .Slice (i , j ),
p .values ,
p .offsets ,
p .definitionLevels ,
p .repetitionLevels ,
)
}
func (p *bufferedPage ) Retain () {
bufferRef (p .values )
bufferRef (p .offsets )
bufferRef (p .definitionLevels )
bufferRef (p .repetitionLevels )
}
func (p *bufferedPage ) Release () {
bufferUnref (p .values )
bufferUnref (p .offsets )
bufferUnref (p .definitionLevels )
bufferUnref (p .repetitionLevels )
}
func bufferRef(buf *buffer ) {
if buf != nil {
buf .ref ()
}
}
func bufferUnref(buf *buffer ) {
if buf != nil {
buf .unref ()
}
}
func Retain (page Page ) {
if p , _ := page .(retainable ); p != nil {
p .Retain ()
}
}
func Release (page Page ) {
if p , _ := page .(releasable ); p != nil {
p .Release ()
}
}
type retainable interface {
Retain()
}
type releasable interface {
Release()
}
var (
_ retainable = (*bufferedPage )(nil )
_ releasable = (*bufferedPage )(nil )
)
The pages are generated with Golds v0.8.4 . (GOOS=linux GOARCH=amd64)
Golds is a Go 101 project developed by Tapir Liu .
PR and bug reports are welcome and can be submitted to the issue list .
Please follow @zigo_101 (reachable from the left QR code) to get the latest news of Golds .