parquet

package module
v0.32.1 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 2, 2026 License: Apache-2.0 Imports: 75 Imported by: 0

README


parquet-go/parquet-go

High-performance Go library to manipulate parquet files, initially developed at Twilio Segment.

Motivation

Parquet has been established as a powerful solution to represent columnar data on persistent storage mediums, achieving levels of compression and query performance that enable managing data sets at scales that reach the petabytes. In addition, having intensive data applications sharing a common format creates opportunities for interoperation in our tool kits, providing greater leverage and value to engineers maintaining and operating those systems.

The creation and evolution of large scale data management systems, combined with realtime expectations come with challenging maintenance and performance requirements, that existing solutions to use parquet with Go were not addressing.

The parquet-go/parquet-go package was designed and developed to respond to those challenges, offering high level APIs to read and write parquet files, while keeping a low compute and memory footprint in order to be used in environments where data volumes and cost constraints require software to achieve high levels of efficiency.

Specification

Columnar storage allows Parquet to store data more efficiently than, say, using JSON or Protobuf. For more information, refer to the Parquet Format Specification.

Installation

The package is distributed as a standard Go module that programs can take a dependency on and install with the following command:

go get github.com/parquet-go/parquet-go

Go 1.22 or later is required to use the package.

Compatibility Guarantees

The package is currently released as a pre-v1 version, which gives maintainers the freedom to break backward compatibility to help improve the APIs as we learn which initial design decisions would need to be revisited to better support the use cases that the library solves for. These occurrences are expected to be rare in frequency and documentation will be produce to guide users on how to adapt their programs to breaking changes.

Usage

The following sections describe how to use APIs exposed by the library, highlighting the use cases with code examples to demonstrate how they are used in practice.

Struct Tags

When using Go structs to define the schema of parquet files, struct fields may include a parquet tag to configure properties of the parquet column such as its name, compression, encoding, and logical type. The first value in the tag sets the column name, and additional comma-separated values set options.

type Record struct {
    ID        int64     `parquet:"id,delta"`
    Name      string    `parquet:"name,dict,zstd"`
    Timestamp int64     `parquet:"timestamp,timestamp(microsecond)"`
    Score     float64   `parquet:"score,split"`
    Tags      []string  `parquet:"tags,list"`
    Optional  *string   `parquet:"optional,optional"`
}

Map keys and values can be configured with the parquet-key and parquet-value tags, and list elements with the parquet-element tag.

For the full reference of supported tags, type constraints, and examples, see the SchemaOf documentation.

Variant Types

The library supports the Parquet VARIANT logical type for storing semi-structured data (JSON-like values) in a columnar format.

The simplest way to use variants is with the variant struct tag, which automatically marshals and unmarshals Go values to/from the variant binary format:

type Event struct {
    ID   int64 `parquet:"id"`
    Data any   `parquet:"data,variant"`
}

writer := parquet.NewGenericWriter[Event](output)
writer.Write([]Event{
    {ID: 1, Data: "hello"},
    {ID: 2, Data: int32(42)},
    {ID: 3, Data: map[string]any{"key": "value"}},
})
writer.Close()

For shredded variants (which store typed columns alongside the raw value for faster queries), build a schema with parquet.ShreddedVariant():

shreddedType, err := parquet.ShreddedVariant(parquet.String())
schema := parquet.NewSchema("Record", parquet.Group{
    "data": shreddedType,
})
writer := parquet.NewGenericWriter[Record](output, schema)

For low-level access to raw variant bytes, use a struct with Metadata and Value []byte fields:

type VariantData struct {
    Metadata []byte `parquet:"metadata"`
    Value    []byte `parquet:"value"`
}

For more examples, see the ExampleVariant and ExampleShreddedVariant functions in example_variant_test.go.

Writing Parquet Files: parquet.GenericWriter[T]

A parquet file is a collection of rows sharing the same schema, arranged in columns to support faster scan operations on subsets of the data set.

For simple use cases, the parquet.WriteFile[T] function allows the creation of parquet files on the file system from a slice of Go values representing the rows to write to the file.

type RowType struct { FirstName, LastName string }

if err := parquet.WriteFile("file.parquet", []RowType{
    {FirstName: "Bob"},
    {FirstName: "Alice"},
}); err != nil {
    ...
}

The parquet.GenericWriter[T] type denormalizes rows into columns, then encodes the columns into a parquet file, generating row groups, column chunks, and pages based on configurable heuristics.

type RowType struct { FirstName, LastName string }

writer := parquet.NewGenericWriter[RowType](output)

_, err := writer.Write([]RowType{
    ...
})
if err != nil {
    ...
}

// Closing the writer is necessary to flush buffers and write the file footer.
if err := writer.Close(); err != nil {
    ...
}

Explicit declaration of the parquet schema on a writer is useful when the application needs to ensure that data written to a file adheres to a predefined schema, which may differ from the schema derived from the writer's type parameter. The parquet.Schema type is a in-memory representation of the schema of parquet rows, translated from the type of Go values, and can be used for this purpose.

schema := parquet.SchemaOf(new(RowType))
writer := parquet.NewGenericWriter[any](output, schema)
...

Reading Parquet Files: parquet.GenericReader[T]

For simple use cases where the data set fits in memory and the program will read most rows of the file, the parquet.ReadFile[T] function returns a slice of Go values representing the rows read from the file.

type RowType struct { FirstName, LastName string }

rows, err := parquet.ReadFile[RowType]("file.parquet")
if err != nil {
    ...
}

for _, c := range rows {
    fmt.Printf("%+v\n", c)
}

The expected schema of rows can be explicitly declared when the reader is constructed, which is useful to ensure that the program receives rows matching an specific format; for example, when dealing with files from remote storage sources that applications cannot trust to have used an expected schema.

Configuring the schema of a reader is done by passing a parquet.Schema instance as argument when constructing a reader. When the schema is declared, conversion rules implemented by the package are applied to ensure that rows read by the application match the desired format (see Evolving Parquet Schemas).

schema := parquet.SchemaOf(new(RowType))
reader := parquet.NewReader(file, schema)
...

Inspecting Parquet Files: parquet.File

Sometimes, lower-level APIs can be useful to leverage the columnar layout of parquet files. The parquet.File type is intended to provide such features to Go applications, by exposing APIs to iterate over the various parts of a parquet file.

f, err := parquet.OpenFile(file, size)
if err != nil {
    ...
}

for _, rowGroup := range f.RowGroups() {
    for _, columnChunk := range rowGroup.ColumnChunks() {
        ...
    }
}

Evolving Parquet Schemas: parquet.Convert

Parquet files embed all the metadata necessary to interpret their content, including a description of the schema of the tables represented by the rows and columns they contain.

Parquet files are also immutable; once written, there is not mechanism for updating a file. If their contents need to be changed, rows must be read, modified, and written to a new file.

Because applications evolve, the schema written to parquet files also tend to evolve over time. Those requirements creating challenges when applications need to operate on parquet files with heterogenous schemas: algorithms that expect new columns to exist may have issues dealing with rows that come from files with mismatching schema versions.

To help build applications that can handle evolving schemas, parquet-go/parquet-go implements conversion rules that create views of row groups to translate between schema versions.

The parquet.Convert function is the low-level routine constructing conversion rules from a source to a target schema. The function is used to build converted views of parquet.RowReader or parquet.RowGroup, for example:

type RowTypeV1 struct { ID int64; FirstName string }
type RowTypeV2 struct { ID int64; FirstName, LastName string }

source := parquet.SchemaOf(RowTypeV1{})
target := parquet.SchemaOf(RowTypeV2{})

conversion, err := parquet.Convert(target, source)
if err != nil {
    ...
}

targetRowGroup := parquet.ConvertRowGroup(sourceRowGroup, conversion)
...

Conversion rules are automatically applied by the parquet.CopyRows function when the reader and writers passed to the function also implement the parquet.RowReaderWithSchema and parquet.RowWriterWithSchema interfaces. The copy determines whether the reader and writer schemas can be converted from one to the other, and automatically applies the conversion rules to facilitate the translation between schemas.

At this time, conversion rules only supports adding or removing columns from the schemas, there are no type conversions performed, nor ways to rename columns, etc... More advanced conversion rules may be added in the future.

Sorting Row Groups: parquet.GenericBuffer[T]

The parquet.GenericWriter[T] type is optimized for minimal memory usage, keeping the order of rows unchanged and flushing pages as soon as they are filled.

Parquet supports expressing columns by which rows are sorted through the declaration of sorting columns on row groups. Sorting row groups requires buffering all rows before ordering and writing them to a parquet file.

To help with those use cases, the parquet-go/parquet-go package exposes the parquet.GenericBuffer[T] type which acts as a buffer of rows and implements sort.Interface to allow applications to sort rows prior to writing them to a file.

The columns that rows are ordered by are configured when creating parquet.GenericBuffer[T] instances using the parquet.SortingColumns function to construct row group options configuring the buffer. The type of parquet columns defines how values are compared, see Parquet Logical Types for details.

When written to a file, the buffer is materialized into a single row group with the declared sorting columns. After being written, buffers can be reused by calling their Reset method.

The following example shows how to use a parquet.GenericBuffer[T] to order rows written to a parquet file:

type RowType struct { FirstName, LastName string }

buffer := parquet.NewGenericBuffer[RowType](
    parquet.SortingRowGroupConfig(
        parquet.SortingColumns(
            parquet.Ascending("LastName"),
            parquet.Ascending("FistName"),
        ),
    ),
)

buffer.Write([]RowType{
    {FirstName: "Luke", LastName: "Skywalker"},
    {FirstName: "Han", LastName: "Solo"},
    {FirstName: "Anakin", LastName: "Skywalker"},
})

sort.Sort(buffer)

writer := parquet.NewGenericWriter[RowType](output)
_, err := parquet.CopyRows(writer, buffer.Rows())
if err != nil {
    ...
}
if err := writer.Close(); err != nil {
    ...
}

Merging Row Groups: parquet.MergeRowGroups

Parquet files are often used as part of the underlying engine for data processing or storage layers, in which cases merging multiple row groups into one that contains more rows can be a useful operation to improve query performance; for example, bloom filters in parquet files are stored for each row group, the larger the row group, the fewer filters need to be stored and the more effective they become.

The parquet-go/parquet-go package supports creating merged views of row groups, where the view contains all the rows of the merged groups, maintaining the order defined by the sorting columns of the groups.

There are a few constraints when merging row groups:

  • The sorting columns of all the row groups must be the same, or the merge operation must be explicitly configured a set of sorting columns which are a prefix of the sorting columns of all merged row groups.

  • The schemas of row groups must all be equal, or the merge operation must be explicitly configured with a schema that all row groups can be converted to, in which case the limitations of schema conversions apply.

Once a merged view is created, it may be written to a new parquet file or buffer in order to create a larger row group:

merge, err := parquet.MergeRowGroups(rowGroups)
if err != nil {
    ...
}

writer := parquet.NewGenericWriter[RowType](output)
_, err := parquet.CopyRows(writer, merge.Rows())
if err != nil {
    ...
}
if err := writer.Close(); err != nil {
    ...
}

Using Bloom Filters: parquet.BloomFilter

Parquet files can embed bloom filters to help improve the performance of point lookups in the files. The format of parquet bloom filters is documented in the parquet specification: Parquet Bloom Filter

By default, no bloom filters are created in parquet files, but applications can configure the list of columns to create filters for using the parquet.BloomFilters option when instantiating writers; for example:

type RowType struct {
    FirstName string `parquet:"first_name"`
    LastName  string `parquet:"last_name"`
}

const filterBitsPerValue = 10
writer := parquet.NewGenericWriter[RowType](output,
    parquet.BloomFilters(
        // Configures the write to generate split-block bloom filters for the
        // "first_name" and "last_name" columns of the parquet schema of rows
        // witten by the application.
        parquet.SplitBlockFilter(filterBitsPerValue, "first_name"),
        parquet.SplitBlockFilter(filterBitsPerValue, "last_name"),
    ),
)
...

Generating bloom filters requires to know how many values exist in a column chunk in order to properly size the filter, which requires buffering all the values written to the column in memory. Because of it, the memory footprint of parquet.GenericWriter[T] increases linearly with the number of columns that the writer needs to generate filters for. This extra cost is optimized away when rows are copied from a parquet.GenericBuffer[T] to a writer, since in this case the number of values per column in known since the buffer already holds all the values in memory.

When reading parquet files, column chunks expose the generated bloom filters with the parquet.ColumnChunk.BloomFilter method, returning a parquet.BloomFilter instance if a filter was available, or nil when there were no filters.

Using bloom filters in parquet files is useful when performing point-lookups in parquet files; searching for column rows matching a given value. Programs can quickly eliminate column chunks that they know does not contain the value they search for by checking the filter first, which is often multiple orders of magnitude faster than scanning the column.

The following code snippet hilights how filters are typically used:

var candidateChunks []parquet.ColumnChunk

for _, rowGroup := range file.RowGroups() {
    columnChunk := rowGroup.ColumnChunks()[columnIndex]
    bloomFilter := columnChunk.BloomFilter()

    if bloomFilter != nil {
        if ok, err := bloomFilter.Check(value); err != nil {
            ...
        } else if !ok {
            // Bloom filters may return false positives, but never return false
            // negatives, we know this column chunk does not contain the value.
            continue
        }
    }

    candidateChunks = append(candidateChunks, columnChunk)
}

Optimizations

The following sections describe common optimization techniques supported by the library.

Optimizing Reads

Lower level APIs used to read parquet files offer more efficient ways to access column values. Consecutive sequences of values are grouped into pages which are represented by the parquet.Page interface.

A column chunk may contain multiple pages, each holding a section of the column values. Applications can retrieve the column values either by reading them into buffers of parquet.Value, or type asserting the pages to read arrays of primitive Go values. The following example demonstrates how to use both mechanisms to read column values:

pages := column.Pages()
defer func() {
    checkErr(pages.Close())
}()

for {
    p, err := pages.ReadPage()
    if err != nil {
        ... // io.EOF when there are no more pages
    }

    switch page := p.Values().(type) {
    case parquet.Int32Reader:
        values := make([]int32, page.NumValues())
        _, err := page.ReadInt32s(values)
        ...
    case parquet.Int64Reader:
        values := make([]int64, page.NumValues())
        _, err := page.ReadInt64s(values)
        ...
    default:
        values := make([]parquet.Value, page.NumValues())
        _, err := page.ReadValues(values)
        ...
    }
}

Reading arrays of typed values is often preferable when performing aggregations on the values as this model offers a more compact representation of the values in memory, and pairs well with the use of optimizations like SIMD vectorization.

Optimizing Writes

Applications that deal with columnar storage are sometimes designed to work with columnar data throughout the abstraction layers; it then becomes possible to write columns of values directly instead of reconstructing rows from the column values. The package offers two main mechanisms to satisfy those use cases:

A. Writing Columns of Typed Arrays

The first solution assumes that the program works with in-memory arrays of typed values, for example slices of primitive Go types like []float32; this would be the case if the application is built on top of a framework like Apache Arrow.

parquet.GenericBuffer[T] is an implementation of the parquet.RowGroup interface which maintains in-memory buffers of column values. Rows can be written by either boxing primitive values into arrays of parquet.Value, or type asserting the columns to a access specialized versions of the write methods accepting arrays of Go primitive types.

When using either of these models, the application is responsible for ensuring that the same number of rows are written to each column or the resulting parquet file will be malformed.

The following examples demonstrate how to use these two models to write columns of Go values:

type RowType struct { FirstName, LastName string }

func writeColumns(buffer *parquet.GenericBuffer[RowType], firstNames []string) error {
    values := make([]parquet.Value, len(firstNames))
    for i := range firstNames {
        values[i] = parquet.ValueOf(firstNames[i])
    }
    _, err := buffer.ColumnBuffers()[0].WriteValues(values)
    return err
}
type RowType struct { ID int64; Value float32 }

func writeColumns(buffer *parquet.GenericBuffer[RowType], ids []int64, values []float32) error {
    if len(ids) != len(values) {
        return fmt.Errorf("number of ids and values mismatch: ids=%d values=%d", len(ids), len(values))
    }
    columns := buffer.ColumnBuffers()
    if err := columns[0].(parquet.Int64Writer).WriteInt64s(ids); err != nil {
        return err
    }
    if err := columns[1].(parquet.FloatWriter).WriteFloats(values); err != nil {
        return err
    }
    return nil
}

The latter is more efficient as it does not require boxing the input into an intermediary array of parquet.Value. However, it may not always be the right model depending on the situation, sometimes the generic abstraction can be a more expressive model.

B. Implementing parquet.RowGroup

Programs that need full control over the construction of row groups can choose to provide their own implementation of the parquet.RowGroup interface, which includes defining implementations of parquet.ColumnChunk and parquet.Page to expose column values of the row group.

This model can be preferable when the underlying storage or in-memory representation of the data needs to be optimized further than what can be achieved by using an intermediary buffering layer with parquet.GenericBuffer[T].

See parquet.RowGroup for the full interface documentation.

C. Using on-disk page buffers

When generating parquet files, the writer needs to buffer all pages before it can create the row group. This may require significant amounts of memory as the entire file content must be buffered prior to generating it. In some cases, the files might even be larger than the amount of memory available to the program.

The parquet.GenericWriter[T] can be configured to use disk storage instead as a scratch buffer when generating files, by configuring a different page buffer pool using the parquet.ColumnPageBuffers option and parquet.PageBufferPool interface.

The parquet-go/parquet-go package provides an implementation of the interface which uses temporary files to store pages while a file is generated, allowing programs to use local storage as swap space to hold pages and keep memory utilization to a minimum. The following example demonstrates how to configure a parquet writer to use on-disk page buffers:

type RowType struct { ... }

writer := parquet.NewGenericWriter[RowType](output,
    parquet.ColumnPageBuffers(
        parquet.NewFileBufferPool("", "buffers.*"),
    ),
)

When a row group is complete, pages buffered to disk need to be copied back to the output file. This results in doubling I/O operations and storage space requirements (the system needs to have enough free disk space to hold two copies of the file). The resulting write amplification can often be optimized away by the kernel if the file system supports copy-on-write of disk pages since copies between os.File instances are optimized using copy_file_range(2) (on linux).

See parquet.PageBufferPool for the full interface documentation.

D. Parallel Column Writes

For applications that need to maximize throughput when writing large columnar datasets, the library supports writing columns in parallel. This is especially useful when each column can be prepared independently and written concurrently, leveraging multiple CPU cores.

You can use Go's goroutines to write to each column's ColumnWriter in parallel. Each column's values can be written using WriteRowValues, and the column must be closed after writing. It is the application's responsibility to ensure that all columns receive the same number of rows, as mismatched row counts will result in malformed files.

Example:

columns   = writer.ColumnWriters()
var (
    wg        sync.WaitGroup
    errs      = make([]error, len(columns))
    rowCounts = make([]int, len(columns))
)
for i, col := range columns {
    wg.Add(1)
    go func(i int, col parquet.ColumnWriter) {
        defer wg.Done()
        n, err := col.WriteRowValues(values[i]) // values[i] is []parquet.Value for column i
        if err != nil {
            errs[i] = err
            return
        }
        rowCounts[i] = n
        errs[i] = col.Close()
    }(i, col)
}
wg.Wait()
// Check errs and rowCounts for consistency

This approach can significantly reduce the time required to write wide tables or large datasets, especially on multi-core systems. However, you should ensure proper error handling and synchronization, as shown above.

SIMD Acceleration with GOAMD64 and GOEXPERIMENT=simd

On amd64, performance-critical kernels of this library (bloom filters, byte scanning, page statistics, encodings) are accelerated with SIMD. For the best results, we encourage building applications with a GOAMD64 microarchitecture level that matches the deployment target:

GOAMD64=v3 go build ...  # CPUs with AVX2 (Haswell/Zen 1 and newer)
GOAMD64=v4 go build ...  # CPUs with AVX-512 (Skylake-SP, Ice Lake, Zen 4 and newer)

At the default GOAMD64=v1, functions like math/bits.OnesCount compile to a runtime CPU feature check with a fallback call; in hot loops the check and the register spills it forces are measurable (we observed some kernels run near 2x slower than at v3/v4). Raising the level makes these intrinsics unconditional single instructions.

The library also ships experimental implementations of its assembly kernels written with the simd/archsimd package, enabled by building with:

GOEXPERIMENT=simd GOAMD64=v4 go build ...

These pure Go implementations benchmark at parity with the hand-written assembly for most kernels (and ahead of it for some, e.g. bloom filter checks and large byte broadcasts), while remaining maintainable Go code with explicit, correct CPU feature gating.

Once CL 813420 is merged, archsimd CPU feature checks become compile-time constants at matching GOAMD64 levels, allowing the compiler to eliminate the runtime feature branches and dead fallback paths entirely: in our benchmarks this improved the GOEXPERIMENT=simd bloom filter kernels by a further 12% on average, bringing them to overall parity with the assembly.

Note that the GOEXPERIMENT=simd build is currently supported with Go 1.26 only: the simd/archsimd package is experimental and Go 1.27 introduced breaking renames to its API (for example LoadUint8x64Slice became LoadUint8x64, StoreSlice became Store, and SumAbsDiff became SumOf8AbsDiff). Support for newer Go versions will follow as the API stabilizes. Builds without GOEXPERIMENT=simd are unaffected and continue to use the assembly kernels on all supported Go versions.

Maintenance

While initial design and development occurred at Twilio Segment, the project is now maintained by the open source community. We welcome external contributors. to participate in the form of discussions or code changes. Please review to the Contribution guidelines as well as the Code of Conduct before submitting contributions.

Continuous Integration

The project uses Github Actions for CI.

Debugging

The package has debugging capabilities built in which can be turned on using the PARQUETGODEBUG environment variable. The value follows a model similar to GODEBUG, it must be formatted as a comma-separated list of key=value pairs.

The following debug flag are currently supported:

  • tracebuf=1 turns on tracing of internal buffers, which validates that reference counters are set to zero when buffers are reclaimed by the garbage collector. When the package detects that a buffer was leaked, it logs an error message along with the stack trace captured when the buffer was last used.

Documentation

Overview

Package parquet is a library for working with parquet files. For an overview of Parquet's qualities as a storage format, see this blog post: https://blog.twitter.com/engineering/en_us/a/2013/dremel-made-simple-with-parquet

Or see the Parquet documentation: https://parquet.apache.org/docs/

Example
package main

import (
	"fmt"
	"io"
	"io/ioutil"
	"log"
	"os"

	"github.com/AlexAslan/parquet-go"
)

func main() {
	// parquet-go uses the same struct-tag definition style as JSON and XML
	type Contact struct {
		Name string `parquet:"name"`
		// "zstd" specifies the compression for this column
		PhoneNumber string `parquet:"phoneNumber,optional,zstd"`
	}

	type AddressBook struct {
		Owner             string    `parquet:"owner,zstd"`
		OwnerPhoneNumbers []string  `parquet:"ownerPhoneNumbers,gzip"`
		Contacts          []Contact `parquet:"contacts"`
	}

	f, _ := ioutil.TempFile("", "parquet-example-")
	writer := parquet.NewWriter(f)
	rows := []AddressBook{
		{Owner: "UserA", Contacts: []Contact{
			{Name: "Alice", PhoneNumber: "+15505551234"},
			{Name: "Bob"},
		}},
		// Add more rows here.
	}
	for _, row := range rows {
		if err := writer.Write(row); err != nil {
			log.Fatal(err)
		}
	}
	_ = writer.Close()
	_ = f.Close()

	// Now, we can read from the file.
	rf, _ := os.Open(f.Name())
	pf := parquet.NewReader(rf)
	addrs := make([]AddressBook, 0)
	for {
		var addr AddressBook
		err := pf.Read(&addr)
		if err == io.EOF {
			break
		}
		if err != nil {
			log.Fatal(err)
		}
		addrs = append(addrs, addr)
	}
	fmt.Println(addrs[0].Owner)
}
Output:
UserA

Index

Examples

Constants

View Source
const (
	DefaultColumnIndexSizeLimit = 16
	DefaultColumnBufferCapacity = 16 * 1024
	DefaultPageBufferSize       = 256 * 1024
	DefaultWriteBufferSize      = 32 * 1024
	DefaultDataPageVersion      = 2
	DefaultDataPageStatistics   = true
	DefaultSkipMagicBytes       = false
	DefaultSkipPageIndex        = false
	DefaultSkipBloomFilters     = false
	DefaultPrefetchBloomFilters = false
	DefaultMaxRowsPerRowGroup   = math.MaxInt64
	DefaultReadMode             = ReadModeSync
)
View Source
const (
	// MaxColumnDepth is the maximum column depth supported by this package.
	MaxColumnDepth = math.MaxUint8

	// MaxColumnIndex is the maximum column index supported by this package.
	MaxColumnIndex = math.MaxUint16 - 1

	// MaxRepetitionLevel is the maximum repetition level supported by this
	// package.
	MaxRepetitionLevel = math.MaxUint8

	// MaxDefinitionLevel is the maximum definition level supported by this
	// package.
	MaxDefinitionLevel = math.MaxUint8

	// MaxRowGroups is the maximum number of row groups which can be contained
	// in a single parquet file.
	//
	// This limit is enforced by the use of 16 bits signed integers in the file
	// metadata footer of parquet files. It is part of the parquet specification
	// and therefore cannot be changed.
	MaxRowGroups = math.MaxInt16
)

Variables

View Source
var (
	ErrMissingBloomFilter = errors.New("missing bloom filter")
	ErrMissingColumnIndex = errors.New("missing column index")
	ErrMissingOffsetIndex = errors.New("missing offset index")
)
View Source
var (
	// Uncompressed is a parquet compression codec representing uncompressed
	// pages.
	Uncompressed uncompressed.Codec

	// Snappy is the SNAPPY parquet compression codec.
	Snappy snappy.Codec

	// Gzip is the GZIP parquet compression codec.
	Gzip = gzip.Codec{
		Level: gzip.DefaultCompression,
	}

	// Brotli is the BROTLI parquet compression codec.
	Brotli = brotli.Codec{
		Quality: brotli.DefaultQuality,
		LGWin:   brotli.DefaultLGWin,
	}

	// Zstd is the ZSTD parquet compression codec.
	Zstd = zstd.Codec{
		Level: zstd.DefaultLevel,
	}

	// Lz4Raw is the LZ4_RAW parquet compression codec.
	Lz4Raw = lz4.Codec{
		Level: lz4.DefaultLevel,
	}
)
View Source
var (
	// Plain is the default parquet encoding.
	Plain plain.Encoding

	// RLE is the hybrid bit-pack/run-length parquet encoding.
	RLE rle.Encoding

	// BitPacked is the deprecated bit-packed encoding for repetition and
	// definition levels.
	BitPacked bitpacked.Encoding

	// PlainDictionary is the plain dictionary parquet encoding.
	//
	// This encoding should not be used anymore in parquet 2.0 and later,
	// it is implemented for backwards compatibility to support reading
	// files that were encoded with older parquet libraries.
	PlainDictionary plain.DictionaryEncoding

	// RLEDictionary is the RLE dictionary parquet encoding.
	RLEDictionary rle.DictionaryEncoding

	// DeltaBinaryPacked is the delta binary packed parquet encoding.
	DeltaBinaryPacked delta.BinaryPackedEncoding

	// DeltaLengthByteArray is the delta length byte array parquet encoding.
	DeltaLengthByteArray delta.LengthByteArrayEncoding

	// DeltaByteArray is the delta byte array parquet encoding.
	DeltaByteArray delta.ByteArrayEncoding

	// ByteStreamSplit is an encoding for numeric and fixed-length binary data.
	ByteStreamSplit bytestreamsplit.Encoding
)
View Source
var (
	// ErrCorrupted is an error returned by the Err method of ColumnPages
	// instances when they encountered a mismatch between the CRC checksum
	// recorded in a page header and the one computed while reading the page
	// data. Also returned on invalid headers / footer decoding.
	ErrCorrupted = errors.New("corrupted parquet page")

	// ErrMissingRootColumn is an error returned when opening an invalid parquet
	// file which does not have a root column.
	ErrMissingRootColumn = errors.New("parquet file is missing a root column")

	// ErrRowGroupSchemaMissing is an error returned when attempting to write a
	// row group but the source has no schema.
	ErrRowGroupSchemaMissing = errors.New("cannot write rows to a row group which has no schema")

	// ErrRowGroupSchemaMismatch is an error returned when attempting to write a
	// row group but the source and destination schemas differ.
	ErrRowGroupSchemaMismatch = errors.New("cannot write row groups with mismatching schemas")

	// ErrRowGroupSortingColumnsMismatch is an error returned when attempting to
	// write a row group but the sorting columns differ in the source and
	// destination.
	ErrRowGroupSortingColumnsMismatch = errors.New("cannot write row groups with mismatching sorting columns")

	// ErrSeekOutOfRange is an error returned when seeking to a row index which
	// is less than the first row of a page.
	ErrSeekOutOfRange = errors.New("seek to row index out of page range")

	// ErrUnexpectedDictionaryPage is an error returned when a page reader
	// encounters a dictionary page after the first page, or in a column
	// which does not use a dictionary encoding.
	ErrUnexpectedDictionaryPage = errors.New("unexpected dictionary page")

	// ErrMissingPageHeader is an error returned when a page reader encounters
	// a malformed page header which is missing page-type-specific information.
	ErrMissingPageHeader = errors.New("missing page header")

	// ErrUnexpectedRepetitionLevels is an error returned when attempting to
	// decode repetition levels into a page which is not part of a repeated
	// column.
	ErrUnexpectedRepetitionLevels = errors.New("unexpected repetition levels")

	// ErrUnexpectedDefinitionLevels is an error returned when attempting to
	// decode definition levels into a page which is part of a required column.
	ErrUnexpectedDefinitionLevels = errors.New("unexpected definition levels")

	// ErrTooManyRowGroups is returned when attempting to generate a parquet
	// file with more than MaxRowGroups row groups.
	ErrTooManyRowGroups = errors.New("the limit of 32767 row groups has been reached")

	// ErrConversion is used to indicate that a conversion betwen two values
	// cannot be done because there are no rules to translate between their
	// physical types.
	ErrInvalidConversion = errors.New("invalid conversion between parquet values")

	// ErrMalformedRepetitionLevel is returned when a page reader encounters
	// a repetition level which does not start at the beginning of a row.
	ErrMalformedRepetitionLevel = errors.New("parquet-go encountered a malformed data page which does not start at the beginning of a row")
)
View Source
var ErrKeyNotFound = errors.New("parquet: encryption key not found")

ErrKeyNotFound is the sentinel error that a KeyRetriever should return (or wrap with %w) when the caller intentionally does not have access to a particular column key. OpenFile treats this as a non-fatal signal and leaves that column inaccessible rather than aborting the open. Any other error from ColumnKey is treated as a hard failure and propagated to the caller.

Functions

func CompareDescending

func CompareDescending(cmp func(Value, Value) int) func(Value, Value) int

CompareDescending constructs a comparison function which inverses the order of values.

func CompareNullsFirst

func CompareNullsFirst(cmp func(Value, Value) int) func(Value, Value) int

CompareNullsFirst constructs a comparison function which assumes that null values are smaller than all other values.

func CompareNullsLast

func CompareNullsLast(cmp func(Value, Value) int) func(Value, Value) int

CompareNullsLast constructs a comparison function which assumes that null values are greater than all other values.

func CopyPages

func CopyPages(dst PageWriter, src PageReader) (numValues int64, err error)

CopyPages copies pages from src to dst, returning the number of values that were copied.

The function returns any error it encounters reading or writing pages, except for io.EOF from the reader which indicates that there were no more pages to read.

func CopyRows

func CopyRows(dst RowWriter, src RowReader) (int64, error)

CopyRows copies rows from src to dst.

The underlying types of src and dst are tested to determine if they expose information about the schema of rows that are read and expected to be written. If the schema information are available but do not match, the function will attempt to automatically convert the rows from the source schema to the destination.

As an optimization, the src argument may implement RowWriterTo to bypass the default row copy logic and provide its own. The dst argument may also implement RowReaderFrom for the same purpose.

The function returns the number of rows written, or any error encountered other than io.EOF.

func CopyValues

func CopyValues(dst ValueWriter, src ValueReader) (int64, error)

CopyValues copies values from src to dst, returning the number of values that were written.

As an optimization, the reader and writer may choose to implement ValueReaderFrom and ValueWriterTo to provide their own copy logic.

The function returns any error it encounters reading or writing pages, except for io.EOF from the reader which indicates that there were no more values to read.

func DeepEqual

func DeepEqual(v1, v2 Value) bool

DeepEqual returns true if v1 and v2 are equal, including their repetition levels, definition levels, and column indexes.

See Equal for details about how value equality is determined.

func Equal

func Equal(v1, v2 Value) bool

Equal returns true if v1 and v2 are equal.

Values are considered equal if they are of the same physical type and hold the same Go values. For BYTE_ARRAY and FIXED_LEN_BYTE_ARRAY, the content of the underlying byte arrays are tested for equality.

Note that the repetition levels, definition levels, and column indexes are not compared by this function, use DeepEqual instead.

func EqualNodes

func EqualNodes(node1, node2 Node) bool

EqualNodes returns true if node1 and node2 are equal.

Nodes that are not of the same repetition type (optional, required, repeated) or of the same hierarchical type (leaf, group) are considered not equal. Leaf nodes are considered equal if they are of the same data type.

Groups are compared recursively, taking the order of fields into account (because it influences the column index of each leaf node), and comparing their logical types: for example, a MAP node is not equal to a GROUP node with the same fields, because MAP nodes have a specific logical type.

Note that the encoding and compression of the nodes are not considered by this function.

func EqualSortingColumns

func EqualSortingColumns(a, b []SortingColumn) bool

EqualSortingColumns compares two slices of sorting columns for equality.

Two sorting column slices are considered equal if they have the same length and each corresponding pair of sorting columns is equal. Two sorting columns are equal if they have:

  • The same column path (including nested field paths)
  • The same sort direction (ascending or descending)
  • The same nulls handling (nulls first or nulls last)

The comparison is order-sensitive, meaning that [A, B] is not equal to [B, A]. Both nil and empty slices are considered equal.

This function is useful for:

  • Validating that merged row groups maintain expected sorting
  • Comparing sorting configurations between different row groups
  • Testing sorting column propagation in merge operations

Example:

cols1 := []SortingColumn{Ascending("name"), Descending("age")}
cols2 := []SortingColumn{Ascending("name"), Descending("age")}
equal := EqualSortingColumns(cols1, cols2) // returns true

cols3 := []SortingColumn{Descending("age"), Ascending("name")}
equal = EqualSortingColumns(cols1, cols3) // returns false (different order)

cols4 := []SortingColumn{Ascending("name"), Ascending("age")}
equal = EqualSortingColumns(cols1, cols4) // returns false (different direction)

func EqualTypes

func EqualTypes(type1, type2 Type) bool

EqualTypes returns true if type1 and type2 are equal.

Types are considered equal if they have the same Kind, Length, and LogicalType. The comparison uses reflect.DeepEqual for LogicalType comparison.

Note: This function is designed for leaf types. For complex group types like MAP and LIST, use EqualNodes instead, as those types require structural comparison of their nested fields.

func Find

func Find(index ColumnIndex, value Value, cmp func(Value, Value) int) int

Find uses the ColumnIndex passed as argument to find the page in a column chunk (determined by the given ColumnIndex) that the given value is expected to be found in.

The function returns the index of the first page that might contain the value. If the function determines that the value does not exist in the index, NumPages is returned.

If you want to search the entire parquet file, you must iterate over the RowGroups and search each one individually, if there are multiple in the file. If you call writer.Flush before closing the file, then you will have multiple RowGroups to iterate over, otherwise Flush is called once on Close.

The comparison function passed as last argument is used to determine the relative order of values. This should generally be the Compare method of the column type, but can sometimes be customized to modify how null values are interpreted, for example:

pageIndex := parquet.Find(columnIndex, value,
	parquet.CompareNullsFirst(typ.Compare),
)

func LookupCompressionCodec

func LookupCompressionCodec(codec format.CompressionCodec) compress.Codec

LookupCompressionCodec returns the compression codec associated with the given code.

The function never returns nil. If the encoding is not supported, an "unsupported" codec is returned.

func LookupEncoding

func LookupEncoding(enc format.Encoding) encoding.Encoding

LookupEncoding returns the parquet encoding associated with the given code.

The function never returns nil. If the encoding is not supported, encoding.NotSupported is returned.

func PrintColumnChunk

func PrintColumnChunk(w io.Writer, columnChunk ColumnChunk) error

func PrintPage

func PrintPage(w io.Writer, page Page) error

func PrintRowGroup

func PrintRowGroup(w io.Writer, rowGroup RowGroup) error

func PrintSchema

func PrintSchema(w io.Writer, name string, node Node) error

func PrintSchemaIndent

func PrintSchemaIndent(w io.Writer, name string, node Node, pattern, newline string) error

func Read

func Read[T any](r io.ReaderAt, size int64, options ...ReaderOption) (rows []T, err error)

Read reads and returns rows from the parquet file in the given reader.

The type T defines the type of rows read from r. T must be compatible with the file's schema or an error will be returned. The row type might represent a subset of the full schema, in which case only a subset of the columns will be loaded from r.

This function is provided for convenience to facilitate reading of parquet files from arbitrary locations in cases where the data set fit in memory.

Example (Any)
package main

import (
	"bytes"
	"fmt"
	"log"

	"github.com/AlexAslan/parquet-go"
)

func main() {
	type Row struct{ FirstName, LastName string }

	buf := new(bytes.Buffer)
	err := parquet.Write(buf, []Row{
		{FirstName: "Luke", LastName: "Skywalker"},
		{FirstName: "Han", LastName: "Solo"},
		{FirstName: "R2", LastName: "D2"},
	})
	if err != nil {
		log.Fatal(err)
	}

	file := bytes.NewReader(buf.Bytes())

	rows, err := parquet.Read[any](file, file.Size())
	if err != nil {
		log.Fatal(err)
	}

	for _, row := range rows {
		fmt.Printf("%q\n", row)
	}

}
Output:
map["FirstName":"Luke" "LastName":"Skywalker"]
map["FirstName":"Han" "LastName":"Solo"]
map["FirstName":"R2" "LastName":"D2"]

func ReadFile

func ReadFile[T any](path string, options ...ReaderOption) (rows []T, err error)

ReadFile reads rows of the parquet file at the given path.

The type T defines the type of rows read from r. T must be compatible with the file's schema or an error will be returned. The row type might represent a subset of the full schema, in which case only a subset of the columns will be loaded from the file.

This function is provided for convenience to facilitate reading of parquet files from the file system in cases where the data set fit in memory.

Example
package main

import (
	"fmt"
	"log"

	"github.com/AlexAslan/parquet-go"
)

func main() {
	type Row struct {
		ID   int64  `parquet:"id"`
		Name string `parquet:"name,zstd"`
	}

	ExampleWriteFile()

	rows, err := parquet.ReadFile[Row]("/tmp/file.parquet")
	if err != nil {
		log.Fatal(err)
	}

	for _, row := range rows {
		fmt.Printf("%d: %q\n", row.ID, row.Name)
	}

}

func ExampleWriteFile() {
	type Row struct {
		ID   int64  `parquet:"id"`
		Name string `parquet:"name,zstd"`
	}

	if err := parquet.WriteFile("/tmp/file.parquet", []Row{
		{ID: 0, Name: "Bob"},
		{ID: 1, Name: "Alice"},
		{ID: 2, Name: "Franky"},
	}); err != nil {
		log.Fatal(err)
	}

}
Output:
0: "Bob"
1: "Alice"
2: "Franky"

func RegisterEncoding

func RegisterEncoding(enc encoding.Encoding)

func Release

func Release(page Page)

Release is a helper function to decrement the reference counter of pages backed by memory which can be granularly managed by the application.

Usage of this is optional and with Retain, is intended to allow finer grained memory management in the application, at the expense of potentially causing panics if the page is used after its reference count has reached zero. Most programs should be able to rely on automated memory management provided by the Go garbage collector instead.

The function should be called to return a page to the internal buffer pool, when a goroutine "releases ownership" it acquired either by being the single owner (e.g. capturing the return value from a ReadPage call) or having gotten shared ownership by calling Retain.

Calling this function on pages that do not embed a reference counter does nothing.

func Retain

func Retain(page Page)

Retain is a helper function to increment the reference counter of pages backed by memory which can be granularly managed by the application.

Usage of this function is optional and with Release, is intended to allow finer grain memory management in the application. Most programs should be able to rely on automated memory management provided by the Go garbage collector instead.

The function should be called when a page lifetime is about to be shared between multiple goroutines or layers of an application, and the program wants to express "sharing ownership" of the page.

Calling this function on pages that do not embed a reference counter does nothing.

func SameNodes

func SameNodes(node1, node2 Node) bool

SameNodes returns true if node1 and node2 are equivalent, ignoring field order.

Unlike EqualNodes, this function considers nodes with the same fields in different orders as equivalent. This is useful when comparing schemas that may have been reordered by operations like MergeNodes.

For leaf nodes, this behaves identically to EqualNodes. For group nodes, this compares fields by name rather than position.

func Search(index ColumnIndex, value Value, typ Type) int

Search is like Find, but uses the default ordering of the given type. Search and Find are scoped to a given ColumnChunk and find the pages within a ColumnChunk which might contain the result. See Find for more details.

func Write

func Write[T any](w io.Writer, rows []T, options ...WriterOption) error

Write writes the given list of rows to a parquet file written to w.

This function is provided for convenience to facilitate the creation of parquet files.

Example (Any)
package main

import (
	"bytes"
	"fmt"
	"log"

	"github.com/AlexAslan/parquet-go"
)

func main() {
	schema := parquet.SchemaOf(struct {
		FirstName string
		LastName  string
	}{})

	buf := new(bytes.Buffer)
	err := parquet.Write[any](
		buf,
		[]any{
			map[string]string{"FirstName": "Luke", "LastName": "Skywalker"},
			map[string]string{"FirstName": "Han", "LastName": "Solo"},
			map[string]string{"FirstName": "R2", "LastName": "D2"},
		},
		schema,
	)
	if err != nil {
		log.Fatal(err)
	}

	file := bytes.NewReader(buf.Bytes())

	rows, err := parquet.Read[any](file, file.Size())
	if err != nil {
		log.Fatal(err)
	}

	for _, row := range rows {
		fmt.Printf("%q\n", row)
	}

}
Output:
map["FirstName":"Luke" "LastName":"Skywalker"]
map["FirstName":"Han" "LastName":"Solo"]
map["FirstName":"R2" "LastName":"D2"]

func WriteFile

func WriteFile[T any](path string, rows []T, options ...WriterOption) error

Write writes the given list of rows to a parquet file written to w.

This function is provided for convenience to facilitate writing parquet files to the file system.

Example
package main

import (
	"log"

	"github.com/AlexAslan/parquet-go"
)

func main() {
	type Row struct {
		ID   int64  `parquet:"id"`
		Name string `parquet:"name,zstd"`
	}

	if err := parquet.WriteFile("/tmp/file.parquet", []Row{
		{ID: 0, Name: "Bob"},
		{ID: 1, Name: "Alice"},
		{ID: 2, Name: "Franky"},
	}); err != nil {
		log.Fatal(err)
	}

}

Types

type BE128Reader

type BE128Reader interface {
	// Read 128-bit big-endian values into the buffer passed as argument,
	// returning the number of values read.
	//
	// The method returns io.EOF when all values have been read.
	ReadBE128s(values [][16]byte) (int, error)
}

BE128Reader is an interface implemented by ValueReader instances which expose the content of a column of 128-bit big-endian byte array values.

type BE128Writer

type BE128Writer interface {
	// Write 128-bit big-endian values.
	//
	// The method returns the number of values written, and any error that
	// occurred while writing the values.
	WriteBE128s(values [][16]byte) (int, error)
}

BE128Writer is an interface implemented by ValueWriter instances which support writing columns of 128-bit big-endian byte array values.

type BloomFilter

type BloomFilter interface {
	// Implement the io.ReaderAt interface as a mechanism to allow reading the
	// raw bits of the filter.
	io.ReaderAt

	// Returns the size of the bloom filter (in bytes).
	Size() int64

	// Tests whether the given value is present in the filter.
	//
	// A non-nil error may be returned if reading the filter failed. This may
	// happen if the filter was lazily loaded from a storage medium during the
	// call to Check for example. Applications that can guarantee that the
	// filter was in memory at the time Check was called can safely ignore the
	// error, which would always be nil in this case.
	Check(value Value) (bool, error)
}

BloomFilter is an interface allowing applications to test whether a key exists in a bloom filter.

type BloomFilterColumn

type BloomFilterColumn interface {
	// Returns the path of the column that the filter applies to.
	Path() []string

	// Returns the hashing algorithm used when inserting values into a bloom
	// filter.
	Hash() bloom.Hash

	// Returns an encoding which can be used to write columns of values to the
	// filter.
	Encoding() encoding.Encoding

	// Returns the size of the filter needed to encode values in the filter,
	// assuming each value will be encoded with the given number of bits.
	Size(numValues int64) int
}

The BloomFilterColumn interface is a declarative representation of bloom filters used when configuring filters on a parquet writer.

func SplitBlockFilter

func SplitBlockFilter(bitsPerValue uint, path ...string) BloomFilterColumn

SplitBlockFilter constructs a split block bloom filter object for the column at the given path, with the given bitsPerValue.

If you are unsure what number of bitsPerValue to use, 10 is a reasonable tradeoff between size and error rate for common datasets.

For more information on the tradeoff between size and error rate, consult this website: https://hur.st/bloomfilter/?n=4000&p=0.1&m=&k=1

type BooleanReader

type BooleanReader interface {
	// Read boolean values into the buffer passed as argument.
	//
	// The method returns io.EOF when all values have been read.
	ReadBooleans(values []bool) (int, error)
}

BooleanReader is an interface implemented by ValueReader instances which expose the content of a column of boolean values.

type BooleanWriter

type BooleanWriter interface {
	// Write boolean values.
	//
	// The method returns the number of values written, and any error that
	// occurred while writing the values.
	WriteBooleans(values []bool) (int, error)
}

BooleanWriter is an interface implemented by ValueWriter instances which support writing columns of boolean values.

type Buffer

type Buffer struct {
	// contains filtered or unexported fields
}

Buffer represents an in-memory group of parquet rows.

The main purpose of the Buffer type is to provide a way to sort rows before writing them to a parquet file. Buffer implements sort.Interface as a way to support reordering the rows that have been written to it.

func NewBuffer

func NewBuffer(options ...RowGroupOption) *Buffer

NewBuffer constructs a new buffer, using the given list of buffer options to configure the buffer returned by the function.

The function panics if the buffer configuration is invalid. Programs that cannot guarantee the validity of the options passed to NewBuffer should construct the buffer configuration independently prior to calling this function:

config, err := parquet.NewRowGroupConfig(options...)
if err != nil {
	// handle the configuration error
	...
} else {
	// this call to create a buffer is guaranteed not to panic
	buffer := parquet.NewBuffer(config)
	...
}

func (*Buffer) ColumnBuffers

func (buf *Buffer) ColumnBuffers() []ColumnBuffer

ColumnBuffers returns the buffer columns.

This method is similar to ColumnChunks, but returns a list of ColumnBuffer instead of a list of ColumnChunk (the latter being read-only); calling ColumnBuffers or ColumnChunks with the same index returns the same underlying objects, but with different types, which removes the need for making a type assertion if the program needed to write directly to the column buffers. The presence of the ColumnChunks method is still required to satisfy the RowGroup interface.

func (*Buffer) ColumnChunks

func (buf *Buffer) ColumnChunks() []ColumnChunk

ColumnChunks returns the buffer columns.

func (*Buffer) Len

func (buf *Buffer) Len() int

Len returns the number of rows written to the buffer.

func (*Buffer) Less

func (buf *Buffer) Less(i, j int) bool

Less returns true if row[i] < row[j] in the buffer.

func (*Buffer) NumRows

func (buf *Buffer) NumRows() int64

NumRows returns the number of rows written to the buffer.

func (*Buffer) Reset

func (buf *Buffer) Reset()

Reset clears the content of the buffer, allowing it to be reused.

func (*Buffer) Rows

func (buf *Buffer) Rows() Rows

Rows returns a reader exposing the current content of the buffer.

The buffer and the returned reader share memory. Mutating the buffer concurrently to reading rows may result in non-deterministic behavior.

func (*Buffer) Schema

func (buf *Buffer) Schema() *Schema

Schema returns the schema of the buffer.

The schema is either configured by passing a Schema in the option list when constructing the buffer, or lazily discovered when the first row is written.

func (*Buffer) Size

func (buf *Buffer) Size() int64

Size returns the estimated size of the buffer in memory (in bytes).

func (*Buffer) SortingColumns

func (buf *Buffer) SortingColumns() []SortingColumn

SortingColumns returns the list of columns by which the buffer will be sorted.

The sorting order is configured by passing a SortingColumns option when constructing the buffer.

func (*Buffer) Swap

func (buf *Buffer) Swap(i, j int)

Swap exchanges the rows at indexes i and j.

func (*Buffer) Write

func (buf *Buffer) Write(row any) error

Write writes a row held in a Go value to the buffer.

func (*Buffer) WriteRowGroup

func (buf *Buffer) WriteRowGroup(rowGroup RowGroup) (int64, error)

WriteRowGroup satisfies the RowGroupWriter interface.

func (*Buffer) WriteRows

func (buf *Buffer) WriteRows(rows []Row) (int, error)

WriteRows writes parquet rows to the buffer.

type BufferPool

type BufferPool interface {
	// GetBuffer is called when a parquet writer needs to acquire a new
	// page buffer from the pool.
	GetBuffer() io.ReadWriteSeeker

	// PutBuffer is called when a parquet writer releases a page buffer to
	// the pool.
	//
	// The parquet.Writer type guarantees that the buffers it calls this method
	// with were previously acquired by a call to GetBuffer on the same
	// pool, and that it will not use them anymore after the call.
	PutBuffer(io.ReadWriteSeeker)
}

BufferPool is an interface abstracting the underlying implementation of page buffer pools.

The parquet-go package provides two implementations of this interface, one backed by in-memory buffers (on the Go heap), and the other using temporary files on disk.

Applications which need finer grain control over the allocation and retention of page buffers may choose to provide their own implementation and install it via the parquet.ColumnPageBuffers writer option.

BufferPool implementations must be safe to use concurrently from multiple goroutines.

func NewBufferPool

func NewBufferPool() BufferPool

NewBufferPool creates a new in-memory page buffer pool.

The implementation is backed by sync.Pool and allocates memory buffers on the Go heap.

func NewChunkBufferPool

func NewChunkBufferPool(chunkSize int) BufferPool

NewChunkBufferPool creates a new in-memory page buffer pool.

The implementation is backed by sync.Pool and allocates memory buffers on the Go heap in fixed-size chunks.

func NewFileBufferPool

func NewFileBufferPool(tempdir, pattern string) BufferPool

NewFileBufferPool creates a new on-disk page buffer pool.

type ByteArrayReader

type ByteArrayReader interface {
	// Read values into the byte buffer passed as argument, returning the number
	// of values written to the buffer (not the number of bytes). Values are
	// written using the PLAIN encoding, each byte array prefixed with its
	// length encoded as a 4 bytes little endian unsigned integer.
	//
	// The method returns io.EOF when all values have been read.
	//
	// If the buffer was not empty, but too small to hold at least one value,
	// io.ErrShortBuffer is returned.
	ReadByteArrays(values []byte) (int, error)
}

ByteArrayReader is an interface implemented by ValueReader instances which expose the content of a column of variable length byte array values.

type ByteArrayWriter

type ByteArrayWriter interface {
	// Write variable length byte array values.
	//
	// The values passed as input must be laid out using the PLAIN encoding,
	// with each byte array prefixed with the four bytes little endian unsigned
	// integer length.
	//
	// The method returns the number of values written to the underlying column
	// (not the number of bytes), or any error that occurred while attempting to
	// write the values.
	WriteByteArrays(values []byte) (int, error)
}

ByteArrayWriter is an interface implemented by ValueWriter instances which support writing columns of variable length byte array values.

type Column

type Column struct {
	// contains filtered or unexported fields
}

Column represents a column in a parquet file.

Methods of Column values are safe to call concurrently from multiple goroutines.

Column instances satisfy the Node interface.

func (*Column) Column

func (c *Column) Column(name string) *Column

Column returns the child column matching the given name.

func (*Column) Columns

func (c *Column) Columns() []*Column

Columns returns the list of child columns.

The method returns the same slice across multiple calls, the program must treat it as a read-only value.

func (*Column) Compression

func (c *Column) Compression() compress.Codec

Compression returns the compression codecs used by this column.

func (*Column) DecodeDataPageV1

func (c *Column) DecodeDataPageV1(header DataPageHeaderV1, page []byte, dict Dictionary) (Page, error)

DecodeDataPageV1 decodes a data page from the header, compressed data, and optional dictionary passed as arguments.

func (*Column) DecodeDataPageV2

func (c *Column) DecodeDataPageV2(header DataPageHeaderV2, page []byte, dict Dictionary) (Page, error)

DecodeDataPageV2 decodes a data page from the header, compressed data, and optional dictionary passed as arguments.

func (*Column) DecodeDictionary

func (c *Column) DecodeDictionary(header DictionaryPageHeader, page []byte) (Dictionary, error)

DecodeDictionary decodes a data page from the header and compressed data passed as arguments.

func (*Column) Depth

func (c *Column) Depth() int

Depth returns the position of the column relative to the root.

func (*Column) Encoding

func (c *Column) Encoding() encoding.Encoding

Encoding returns the encodings used by this column.

func (*Column) Fields

func (c *Column) Fields() []Field

Fields returns the list of fields on the column.

func (*Column) GoType

func (c *Column) GoType() reflect.Type

GoType returns the Go type that best represents the parquet column.

func (*Column) ID

func (c *Column) ID() int

ID returns column field id

func (*Column) Index

func (c *Column) Index() int

Index returns the position of the column in a row. Only leaf columns have a column index, the method returns -1 when called on non-leaf columns.

func (*Column) Leaf

func (c *Column) Leaf() bool

Leaf returns true if c is a leaf column.

func (*Column) MaxDefinitionLevel

func (c *Column) MaxDefinitionLevel() int

MaxDefinitionLevel returns the maximum value of definition levels on this column.

func (*Column) MaxRepetitionLevel

func (c *Column) MaxRepetitionLevel() int

MaxRepetitionLevel returns the maximum value of repetition levels on this column.

func (*Column) Name

func (c *Column) Name() string

Name returns the column name.

func (*Column) Optional

func (c *Column) Optional() bool

Optional returns true if the column is optional.

func (*Column) Pages

func (c *Column) Pages() Pages

Pages returns a reader exposing all pages in this column, across row groups.

func (*Column) PagesFrom

func (c *Column) PagesFrom(reader io.ReaderAt) Pages

func (*Column) Path

func (c *Column) Path() []string

Path of the column in the parquet schema.

func (*Column) Repeated

func (c *Column) Repeated() bool

Repeated returns true if the column may repeat.

func (*Column) Required

func (c *Column) Required() bool

Required returns true if the column is required.

func (*Column) String

func (c *Column) String() string

String returns a human-readable string representation of the column.

func (*Column) Type

func (c *Column) Type() Type

Type returns the type of the column.

The returned value is unspecified if c is not a leaf column.

func (*Column) Value

func (c *Column) Value(base reflect.Value) reflect.Value

Value returns the sub-value in base for the child column at the given index.

type ColumnBuffer

type ColumnBuffer interface {
	// Exposes a read-only view of the column buffer.
	ColumnChunk

	// The column implements ValueReaderAt as a mechanism to read values at
	// specific locations within the buffer.
	ValueReaderAt

	// The column implements ValueWriter as a mechanism to optimize the copy
	// of values into the buffer in contexts where the row information is
	// provided by the values because the repetition and definition levels
	// are set.
	ValueWriter

	// For indexed columns, returns the underlying dictionary holding the column
	// values. If the column is not indexed, nil is returned.
	Dictionary() Dictionary

	// Returns a copy of the column. The returned copy shares no memory with
	// the original, mutations of either column will not modify the other.
	Clone() ColumnBuffer

	// Returns the column as a Page.
	Page() Page

	// Clears all rows written to the column.
	Reset()

	// Returns the current capacity of the column (rows).
	Cap() int

	// Returns the number of rows currently written to the column.
	Len() int

	// Compares rows at index i and j and reports whether i < j.
	Less(i, j int) bool

	// Swaps rows at index i and j.
	Swap(i, j int)

	// Returns the size of the column buffer in bytes.
	Size() int64
	// contains filtered or unexported methods
}

ColumnBuffer is an interface representing columns of a row group.

ColumnBuffer implements sort.Interface as a way to support reordering the rows that have been written to it.

The current implementation has a limitation which prevents applications from providing custom versions of this interface because it contains unexported methods. The only way to create ColumnBuffer values is to call the NewColumnBuffer of Type instances. This limitation may be lifted in future releases.

type ColumnChunk

type ColumnChunk interface {
	// Returns the column type.
	Type() Type

	// Returns the index of this column in its parent row group.
	Column() int

	// Returns a reader exposing the pages of the column.
	Pages() Pages

	// Returns the components of the page index for this column chunk,
	// containing details about the content and location of pages within the
	// chunk.
	//
	// Note that the returned value may be the same across calls to these
	// methods, programs must treat those as read-only.
	//
	// If the column chunk does not have a column or offset index, the methods return
	// ErrMissingColumnIndex or ErrMissingOffsetIndex respectively.
	//
	// Prior to v0.20, these methods did not return an error because the page index
	// for a file was either fully read when the file was opened, or skipped
	// completely using the parquet.SkipPageIndex option. Version v0.20 introduced a
	// change that the page index can be read on-demand at any time, even if a file
	// was opened with the parquet.SkipPageIndex option. Since reading the page index
	// can fail, these methods now return an error.
	ColumnIndex() (ColumnIndex, error)
	OffsetIndex() (OffsetIndex, error)
	BloomFilter() BloomFilter

	// Returns the number of values in the column chunk.
	//
	// This quantity may differ from the number of rows in the parent row group
	// because repeated columns may hold zero or more values per row.
	NumValues() int64
}

The ColumnChunk interface represents individual columns of a row group.

func AsyncColumnChunk

func AsyncColumnChunk(columnChunk ColumnChunk) ColumnChunk

AsyncColumnChunk returns a ColumnChunk that reads pages asynchronously.

type ColumnChunkValueReader

type ColumnChunkValueReader interface {
	ValueReader
	RowSeeker
	io.Closer
}

ColumnChunkValueReader is an interface for reading values from a column chunk.

func NewColumnChunkValueReader

func NewColumnChunkValueReader(column ColumnChunk) ColumnChunkValueReader

NewColumnChunkValueReader creates a new ColumnChunkValueReader for the given column chunk.

type ColumnIndex

type ColumnIndex interface {
	// NumPages returns the number of paged in the column index.
	NumPages() int

	// Returns the number of null values in the page at the given index.
	NullCount(int) int64

	// Tells whether the page at the given index contains null values only.
	NullPage(int) bool

	// PageIndex return min/max bounds for the page at the given index in the
	// column.
	MinValue(int) Value
	MaxValue(int) Value

	// IsAscending returns true if the column index min/max values are sorted
	// in ascending order (based on the ordering rules of the column's logical
	// type).
	IsAscending() bool

	// IsDescending returns true if the column index min/max values are sorted
	// in descending order (based on the ordering rules of the column's logical
	// type).
	IsDescending() bool
}

func NewColumnIndex

func NewColumnIndex(kind Kind, index *format.ColumnIndex) ColumnIndex

NewColumnIndex constructs a ColumnIndex instance from the given parquet format column index. The kind argument configures the type of values

type ColumnIndexer

type ColumnIndexer interface {
	// Resets the column indexer state.
	Reset()

	// Add a page to the column indexer.
	IndexPage(numValues, numNulls int64, min, max Value)

	// Generates a format.ColumnIndex value from the current state of the
	// column indexer.
	//
	// The returned value may reference internal buffers, in which case the
	// values remain valid until the next call to IndexPage or Reset on the
	// column indexer.
	ColumnIndex() format.ColumnIndex
}

The ColumnIndexer interface is implemented by types that support generating parquet column indexes.

The package does not export any types that implement this interface, programs must call NewColumnIndexer on a Type instance to construct column indexers.

type ColumnWriter

type ColumnWriter struct {
	// contains filtered or unexported fields
}

ColumnWriter writes values for a single column to underlying medium.

func (*ColumnWriter) Close

func (c *ColumnWriter) Close() (err error)

Close closes the column writer and resets all dependent resources. It can be reused after Close is called.

func (*ColumnWriter) Flush

func (c *ColumnWriter) Flush() (err error)

Flush writes any buffered data to the underlying io.Writer.

func (*ColumnWriter) WriteRowValues

func (c *ColumnWriter) WriteRowValues(rows []Value) (int, error)

WriteRowValues writes entire rows to the column. On success, this returns the number of rows written (not the number of values).

Unlike ValueWriter, where arbitrary values may be written regardless of row boundaries, this method requires whole rows. This is because the written values may be automatically flushed to a data page, based on the writer's configured page buffer size, and a single row is not permitted to span two pages.

type ConcurrentRowGroupWriter

type ConcurrentRowGroupWriter struct {
	// contains filtered or unexported fields
}

ConcurrentRowGroupWriter is a row group writer that can be used to write row groups in parallel. Multiple row groups can be created concurrently and written to independently, but they must be committed serially to maintain the order of row groups in the file.

See BeginRowGroup for more information on how this can be used.

While multiple row groups can be created concurrently, a single row group must be written sequentially.

func (*ConcurrentRowGroupWriter) ColumnWriters

func (rg *ConcurrentRowGroupWriter) ColumnWriters() []*ColumnWriter

ColumnWriters returns the column writers for this row group, allowing direct access to write values to individual columns.

func (*ConcurrentRowGroupWriter) Commit

func (rg *ConcurrentRowGroupWriter) Commit() (int64, error)

Commit commits the row group to the parent writer, returning the number of rows written and an error if any. This method must be called serially (not concurrently) to maintain row group order in the file.

If the parent writer has any pending rows buffered, they will be flushed before this row group is written.

After Commit returns successfully, the row group will be empty and can be reused.

func (*ConcurrentRowGroupWriter) Flush

func (rg *ConcurrentRowGroupWriter) Flush() error

Flush flushes any buffered data in the row group's column writers. This could be called before Commit to ensure all data pages are flushed.

func (*ConcurrentRowGroupWriter) Schema

func (rg *ConcurrentRowGroupWriter) Schema() *Schema

Schema returns the schema for this row group.

func (*ConcurrentRowGroupWriter) Size

func (rg *ConcurrentRowGroupWriter) Size() int64

Size returns an estimate of the current row group size in bytes.

The estimate sums encoded pages, buffered values not yet encoded, dictionary pages, and bloom filters for each column. Because buffered values have not been compressed, this is an upper-bound estimate.

func (*ConcurrentRowGroupWriter) WriteRows

func (rg *ConcurrentRowGroupWriter) WriteRows(rows []Row) (int, error)

WriteRows writes rows to the row group.

type Conversion

type Conversion interface {
	// Applies the conversion logic on the src row, returning the result
	// appended to dst.
	Convert(rows []Row) (int, error)
	// Converts the given column index in the target schema to the original
	// column index in the source schema of the conversion.
	Column(int) int
	// Returns the target schema of the conversion.
	Schema() *Schema
}

Conversion is an interface implemented by types that provide conversion of parquet rows from one schema to another.

Conversion instances must be safe to use concurrently from multiple goroutines.

func Convert

func Convert(to, from Node) (conv Conversion, err error)

Convert constructs a conversion function from one parquet schema to another.

The function supports converting between schemas where the source or target have extra columns; if there are more columns in the source, they will be stripped out of the rows. Extra columns in the target schema will be set to null or zero values.

Variant columns that the target declares unshredded but the source stores shredded are reconstructed with the construct_variant algorithm of the Variant Shredding specification instead of being mapped column-by-column (see convert_variant.go).

The returned function is intended to be used to append the converted source row to the destination buffer.

type ConvertError

type ConvertError struct {
	Path []string
	From Node
	To   Node
}

ConvertError is an error type returned by calls to Convert when the conversion of parquet schemas is impossible or the input row for the conversion is malformed.

func (*ConvertError) Error

func (e *ConvertError) Error() string

Error satisfies the error interface.

type DataPageHeader

type DataPageHeader interface {
	PageHeader

	// Returns the encoding of the repetition level section.
	RepetitionLevelEncoding() format.Encoding

	// Returns the encoding of the definition level section.
	DefinitionLevelEncoding() format.Encoding

	// Returns the number of null values in the page.
	NullCount() int64

	// Returns the minimum value in the page based on the ordering rules of the
	// column's logical type.
	//
	// As an optimization, the method may return the same slice across multiple
	// calls. Programs must treat the returned value as immutable to prevent
	// unpredictable behaviors.
	//
	// If the page only contains only null values, an empty slice is returned.
	MinValue() []byte

	// Returns the maximum value in the page based on the ordering rules of the
	// column's logical type.
	//
	// As an optimization, the method may return the same slice across multiple
	// calls. Programs must treat the returned value as immutable to prevent
	// unpredictable behaviors.
	//
	// If the page only contains only null values, an empty slice is returned.
	MaxValue() []byte
}

DataPageHeader is a specialization of the PageHeader interface implemented by data pages.

type DataPageHeaderV1

type DataPageHeaderV1 struct {
	// contains filtered or unexported fields
}

DataPageHeaderV1 is an implementation of the DataPageHeader interface representing data pages version 1.

func (DataPageHeaderV1) DefinitionLevelEncoding

func (v1 DataPageHeaderV1) DefinitionLevelEncoding() format.Encoding

func (DataPageHeaderV1) Encoding

func (v1 DataPageHeaderV1) Encoding() format.Encoding

func (DataPageHeaderV1) MaxValue

func (v1 DataPageHeaderV1) MaxValue() []byte

func (DataPageHeaderV1) MinValue

func (v1 DataPageHeaderV1) MinValue() []byte

func (DataPageHeaderV1) NullCount

func (v1 DataPageHeaderV1) NullCount() int64

func (DataPageHeaderV1) NumValues

func (v1 DataPageHeaderV1) NumValues() int64

func (DataPageHeaderV1) PageType

func (v1 DataPageHeaderV1) PageType() format.PageType

func (DataPageHeaderV1) RepetitionLevelEncoding

func (v1 DataPageHeaderV1) RepetitionLevelEncoding() format.Encoding

func (DataPageHeaderV1) String

func (v1 DataPageHeaderV1) String() string

type DataPageHeaderV2

type DataPageHeaderV2 struct {
	// contains filtered or unexported fields
}

DataPageHeaderV2 is an implementation of the DataPageHeader interface representing data pages version 2.

func (DataPageHeaderV2) DefinitionLevelEncoding

func (v2 DataPageHeaderV2) DefinitionLevelEncoding() format.Encoding

func (DataPageHeaderV2) DefinitionLevelsByteLength

func (v2 DataPageHeaderV2) DefinitionLevelsByteLength() int64

func (DataPageHeaderV2) Encoding

func (v2 DataPageHeaderV2) Encoding() format.Encoding

func (DataPageHeaderV2) IsCompressed

func (v2 DataPageHeaderV2) IsCompressed() bool

func (DataPageHeaderV2) MaxValue

func (v2 DataPageHeaderV2) MaxValue() []byte

func (DataPageHeaderV2) MinValue

func (v2 DataPageHeaderV2) MinValue() []byte

func (DataPageHeaderV2) NullCount

func (v2 DataPageHeaderV2) NullCount() int64

func (DataPageHeaderV2) NumNulls

func (v2 DataPageHeaderV2) NumNulls() int64

func (DataPageHeaderV2) NumRows

func (v2 DataPageHeaderV2) NumRows() int64

func (DataPageHeaderV2) NumValues

func (v2 DataPageHeaderV2) NumValues() int64

func (DataPageHeaderV2) PageType

func (v2 DataPageHeaderV2) PageType() format.PageType

func (DataPageHeaderV2) RepetitionLevelEncoding

func (v2 DataPageHeaderV2) RepetitionLevelEncoding() format.Encoding

func (DataPageHeaderV2) RepetitionLevelsByteLength

func (v2 DataPageHeaderV2) RepetitionLevelsByteLength() int64

func (DataPageHeaderV2) String

func (v2 DataPageHeaderV2) String() string

type DecryptionConfig

type DecryptionConfig struct {
	Keys KeyRetriever
}

DecryptionConfig holds the read-side decryption configuration.

type Dictionary

type Dictionary interface {
	// Returns the type that the dictionary was created from.
	Type() Type

	// Returns the number of value indexed in the dictionary.
	Len() int

	// Returns the total size in bytes of all values stored in the dictionary.
	// This is used for tracking dictionary memory usage and enforcing size limits.
	Size() int64

	// Returns the dictionary value at the given index.
	Index(index int32) Value

	// Inserts values from the second slice to the dictionary and writes the
	// indexes at which each value was inserted to the first slice.
	//
	// The method panics if the length of the indexes slice is smaller than the
	// length of the values slice.
	Insert(indexes []int32, values []Value)

	// Given an array of dictionary indexes, lookup the values into the array
	// of values passed as second argument.
	//
	// The method panics if len(indexes) > len(values), or one of the indexes
	// is negative or greater than the highest index in the dictionary.
	Lookup(indexes []int32, values []Value)

	// Returns the min and max values found in the given indexes.
	Bounds(indexes []int32) (min, max Value)

	// Resets the dictionary to its initial state, removing all values.
	Reset()

	// Returns a Page representing the content of the dictionary.
	//
	// The returned page shares the underlying memory of the buffer, it remains
	// valid to use until the dictionary's Reset method is called.
	Page() Page
	// contains filtered or unexported methods
}

The Dictionary interface represents type-specific implementations of parquet dictionaries.

Programs can instantiate dictionaries by call the NewDictionary method of a Type object.

The current implementation has a limitation which prevents applications from providing custom versions of this interface because it contains unexported methods. The only way to create Dictionary values is to call the NewDictionary of Type instances. This limitation may be lifted in future releases.

type DictionaryPageHeader

type DictionaryPageHeader struct {
	// contains filtered or unexported fields
}

DictionaryPageHeader is an implementation of the PageHeader interface representing dictionary pages.

func (DictionaryPageHeader) Encoding

func (dict DictionaryPageHeader) Encoding() format.Encoding

func (DictionaryPageHeader) IsSorted

func (dict DictionaryPageHeader) IsSorted() bool

func (DictionaryPageHeader) NumValues

func (dict DictionaryPageHeader) NumValues() int64

func (DictionaryPageHeader) PageType

func (dict DictionaryPageHeader) PageType() format.PageType

func (DictionaryPageHeader) String

func (dict DictionaryPageHeader) String() string

type DoubleReader

type DoubleReader interface {
	// Read double-precision floating point values into the buffer passed as
	// argument.
	//
	// The method returns io.EOF when all values have been read.
	ReadDoubles(values []float64) (int, error)
}

DoubleReader is an interface implemented by ValueReader instances which expose the content of a column of double-precision float point values.

type DoubleWriter

type DoubleWriter interface {
	// Write double-precision floating point values.
	//
	// The method returns the number of values written, and any error that
	// occurred while writing the values.
	WriteDoubles(values []float64) (int, error)
}

DoubleWriter is an interface implemented by ValueWriter instances which support writing columns of double-precision floating point values.

type EncryptionAlgorithmType

type EncryptionAlgorithmType int

EncryptionAlgorithmType selects the Parquet encryption algorithm.

const (
	// AES_GCM_V1 is the default algorithm (AES-GCM authenticated encryption).
	AES_GCM_V1 EncryptionAlgorithmType = iota
	// AES_GCM_CTR_V1 is the CTR variant (deprecated, included for spec completeness).
	AES_GCM_CTR_V1
)

type EncryptionConfig

type EncryptionConfig struct {
	// FooterKey is the AES key used to encrypt the footer and any columns
	// that do not have a per-column key. Must be 16, 24, or 32 bytes.
	FooterKey []byte

	// ColumnKeys maps a dot-joined column path (e.g. "a.b.c") to its AES key.
	// Columns not listed here use FooterKey.
	ColumnKeys map[string][]byte

	// Algorithm selects AES_GCM_V1 (default) or AES_GCM_CTR_V1.
	Algorithm EncryptionAlgorithmType

	// EncryptedFooter controls the file layout:
	//   true  → encrypted footer, file uses "PARE" magic.
	//   false → plaintext footer with GCM signature appended, file uses "PAR1" magic.
	EncryptedFooter bool

	// AadPrefix is an optional byte prefix that is prepended to every AAD.
	AadPrefix []byte

	// FileIdentifier is an 8-byte per-file unique value embedded in every AAD.
	// If nil, 8 random bytes are generated when the writer is created.
	FileIdentifier []byte
}

EncryptionConfig holds the write-side encryption configuration.

type Field

type Field interface {
	Node

	// Returns the name of this field in its parent node.
	Name() string

	// Given a reference to the Go value matching the structure of the parent
	// node, returns the Go value of the field.
	Value(base reflect.Value) reflect.Value
}

Field instances represent fields of a parquet node, which associate a node to their name in their parent node.

type File

type File struct {
	// contains filtered or unexported fields
}

File represents a parquet file. The layout of a Parquet file can be found here: https://github.com/apache/parquet-format#file-format

func OpenFile

func OpenFile(r io.ReaderAt, size int64, options ...FileOption) (*File, error)

OpenFile opens a parquet file and reads the content between offset 0 and the given size in r.

Only the parquet magic bytes and footer are read, column chunks and other parts of the file are left untouched; this means that successfully opening a file does not validate that the pages have valid checksums.

func (*File) ColumnIndexes

func (f *File) ColumnIndexes() []format.ColumnIndex

ColumnIndexes returns the page index of the parquet file f.

If the file did not contain a column index, the method returns an empty slice.

func (*File) Lookup

func (f *File) Lookup(key string) (value string, ok bool)

Lookup returns the value associated with the given key in the file key/value metadata.

The ok boolean will be true if the key was found, false otherwise.

func (*File) Metadata

func (f *File) Metadata() *format.FileMetaData

Metadata returns the metadata of f.

func (*File) NumRows

func (f *File) NumRows() int64

NumRows returns the number of rows in the file.

func (*File) OffsetIndexes

func (f *File) OffsetIndexes() []format.OffsetIndex

OffsetIndexes returns the page index of the parquet file f.

If the file did not contain an offset index, the method returns an empty slice.

func (*File) ReadAt

func (f *File) ReadAt(b []byte, off int64) (int, error)

ReadAt reads bytes into b from f at the given offset.

The method satisfies the io.ReaderAt interface.

func (*File) ReadPageIndex

func (f *File) ReadPageIndex() ([]format.ColumnIndex, []format.OffsetIndex, error)

ReadPageIndex reads the page index section of the parquet file f.

If the file did not contain a page index, the method returns two empty slices and a nil error.

Only leaf columns have indexes, the returned indexes are arranged using the following layout:

------------------
| col 0: chunk 0 |
------------------
| col 1: chunk 0 |
------------------
| ...            |
------------------
| col 0: chunk 1 |
------------------
| col 1: chunk 1 |
------------------
| ...            |
------------------

This method is useful in combination with the SkipPageIndex option to delay reading the page index section until after the file was opened. Note that in this case the page index is not cached within the file, programs are expected to make use of independently from the parquet package.

func (*File) Root

func (f *File) Root() *Column

Root returns the root column of f.

func (*File) RowGroups

func (f *File) RowGroups() []RowGroup

RowGroups returns the list of row groups in the file.

Elements of the returned slice are guaranteed to be of type *FileRowGroup.

func (*File) Schema

func (f *File) Schema() *Schema

Schema returns the schema of f.

func (*File) Size

func (f *File) Size() int64

Size returns the size of f (in bytes).

type FileBloomFilter

type FileBloomFilter struct {
	io.SectionReader
	// contains filtered or unexported fields
}

func (*FileBloomFilter) Check

func (f *FileBloomFilter) Check(v Value) (bool, error)

type FileColumnChunk

type FileColumnChunk struct {
	// contains filtered or unexported fields
}

FileColumnChunk is an implementation of the ColumnChunk interface on parquet files returned by OpenFile.

func (*FileColumnChunk) BloomFilter

func (c *FileColumnChunk) BloomFilter() BloomFilter

BloomFilter returns the bloom filter of the column chunk, or nil if it didn't have one.

func (*FileColumnChunk) BloomFilterFrom

func (c *FileColumnChunk) BloomFilterFrom(reader io.ReaderAt) (*FileBloomFilter, error)

BloomFilterFrom is like BloomFilter but uses the reader passed as argument to read the bloom filter.

func (*FileColumnChunk) Bounds

func (c *FileColumnChunk) Bounds() (min, max Value, ok bool)

Bounds returns the min and max values found in the column chunk.

func (*FileColumnChunk) Column

func (c *FileColumnChunk) Column() int

Column returns the column index of this chunk in its parent row group.

func (*FileColumnChunk) ColumnIndex

func (c *FileColumnChunk) ColumnIndex() (ColumnIndex, error)

ColumnIndex returns the column index of the column chunk, or an error if it didn't exist or couldn't be read.

func (*FileColumnChunk) ColumnIndexFrom

func (c *FileColumnChunk) ColumnIndexFrom(reader io.ReaderAt) (*FileColumnIndex, error)

ColumnIndexFrom is like ColumnIndex but uses the reader passed as argument to read the column index.

func (*FileColumnChunk) File

func (c *FileColumnChunk) File() *File

File returns the file that this column chunk belongs to.

func (*FileColumnChunk) Node

func (c *FileColumnChunk) Node() Node

Node returns the node that this column chunk belongs to in the parquet schema.

func (*FileColumnChunk) NullCount

func (c *FileColumnChunk) NullCount() int64

NullCount returns the number of null values in the column chunk.

This value is extracted from the column chunk statistics, parquet writers are not required to populate it.

func (*FileColumnChunk) NumValues

func (c *FileColumnChunk) NumValues() int64

NumValues returns the number of values in the column chunk.

func (*FileColumnChunk) OffsetIndex

func (c *FileColumnChunk) OffsetIndex() (OffsetIndex, error)

OffsetIndex returns the offset index of the column chunk, or an error if it didn't exist or couldn't be read.

func (*FileColumnChunk) OffsetIndexFrom

func (c *FileColumnChunk) OffsetIndexFrom(reader io.ReaderAt) (*FileOffsetIndex, error)

OffsetIndexFrom is like OffsetIndex but uses the reader passed as argument to read the offset index.

func (*FileColumnChunk) Pages

func (c *FileColumnChunk) Pages() Pages

Pages returns a page reader for the column chunk.

func (*FileColumnChunk) PagesFrom

func (c *FileColumnChunk) PagesFrom(reader io.ReaderAt) *FilePages

PagesFrom returns a page reader for the column chunk, using the reader passed as argument instead of the one that the file was originally opened from.

Note that unlike when calling Pages, the returned reader is not wrapped in an AsyncPages reader if the file was opened in async mode.

func (*FileColumnChunk) Type

func (c *FileColumnChunk) Type() Type

Type returns the type of the column chunk.

type FileColumnIndex

type FileColumnIndex struct {
	// contains filtered or unexported fields
}

func (*FileColumnIndex) IsAscending

func (i *FileColumnIndex) IsAscending() bool

func (*FileColumnIndex) IsDescending

func (i *FileColumnIndex) IsDescending() bool

func (*FileColumnIndex) MaxValue

func (i *FileColumnIndex) MaxValue(j int) Value

func (*FileColumnIndex) MinValue

func (i *FileColumnIndex) MinValue(j int) Value

func (*FileColumnIndex) NullCount

func (i *FileColumnIndex) NullCount(j int) int64

func (*FileColumnIndex) NullPage

func (i *FileColumnIndex) NullPage(j int) bool

func (*FileColumnIndex) NumPages

func (i *FileColumnIndex) NumPages() int

type FileConfig

type FileConfig struct {
	SkipMagicBytes       bool
	SkipPageIndex        bool
	SkipBloomFilters     bool
	PrefetchBloomFilters bool
	OptimisticRead       bool
	ReadBufferSize       int
	ReadMode             ReadMode
	Schema               *Schema
	Decryption           *DecryptionConfig
}

The FileConfig type carries configuration options for parquet files.

FileConfig implements the FileOption interface so it can be used directly as argument to the OpenFile function when needed, for example:

f, err := parquet.OpenFile(reader, size, &parquet.FileConfig{
	SkipPageIndex:    true,
	SkipBloomFilters: true,
	ReadMode:         ReadModeAsync,
})

func DefaultFileConfig

func DefaultFileConfig() *FileConfig

DefaultFileConfig returns a new FileConfig value initialized with the default file configuration.

func NewFileConfig

func NewFileConfig(options ...FileOption) (*FileConfig, error)

NewFileConfig constructs a new file configuration applying the options passed as arguments.

The function returns an non-nil error if some of the options carried invalid configuration values.

func (*FileConfig) Apply

func (c *FileConfig) Apply(options ...FileOption)

Apply applies the given list of options to c.

func (*FileConfig) ConfigureFile

func (c *FileConfig) ConfigureFile(config *FileConfig)

ConfigureFile applies configuration options from c to config.

func (*FileConfig) Validate

func (c *FileConfig) Validate() error

Validate returns a non-nil error if the configuration of c is invalid.

type FileOffsetIndex

type FileOffsetIndex struct {
	// contains filtered or unexported fields
}

func (*FileOffsetIndex) CompressedPageSize

func (i *FileOffsetIndex) CompressedPageSize(j int) int64

func (*FileOffsetIndex) FirstRowIndex

func (i *FileOffsetIndex) FirstRowIndex(j int) int64

func (*FileOffsetIndex) NumPages

func (i *FileOffsetIndex) NumPages() int

func (*FileOffsetIndex) Offset

func (i *FileOffsetIndex) Offset(j int) int64

type FileOption

type FileOption interface {
	ConfigureFile(*FileConfig)
}

FileOption is an interface implemented by types that carry configuration options for parquet files.

func FileReadMode

func FileReadMode(mode ReadMode) FileOption

FileReadMode is a file configuration option which controls the way pages are read. Currently the only two options are ReadModeAsync and ReadModeSync which control whether or not pages are loaded asynchronously. It can be advantageous to use ReadModeAsync if your reader is backed by network storage.

Defaults to ReadModeSync.

func FileSchema

func FileSchema(schema *Schema) FileOption

FileSchema is used to pass a known schema in while opening a Parquet file. This optimization is only useful if your application is currently opening an extremely large number of parquet files with the same, known schema.

Defaults to nil.

func OptimisticRead

func OptimisticRead(enabled bool) FileOption

OptimisticRead configures a file to optimistically perform larger buffered reads to improve performance. This is useful when reading from remote storage and amortize the cost of network round trips.

This is an option instead of enabled by default because dependents of this package have historically relied on the read patterns to provide external caches and achieve similar results (e.g., Tempo).

func PrefetchBloomFilters

func PrefetchBloomFilters(prefetch bool) FileOption

PrefetchBloomFilters is a file configuration option that controls whether the bloom filter contents are loaded into memory when a file is opened. By default, only the headers are parsed, requiring further reads to the file to probe the filter. Using this option with OptimisticRead can be useful when reading from remote storage, reducing network round trips.

Defaults to false.

func ReadBufferSize

func ReadBufferSize(size int) FileOption

ReadBufferSize is a file configuration option which controls the default buffer sizes for reads made to the provided io.Reader. The default of 4096 is appropriate for disk based access but if your reader is backed by network storage it can be advantageous to increase this value to something more like 4 MiB.

Defaults to 4096.

func SkipBloomFilters

func SkipBloomFilters(skip bool) FileOption

SkipBloomFilters is a file configuration option which prevents automatically reading the bloom filter headers when opening a parquet file, when set to true. This is useful as an optimization when programs know that they will not need to consume the bloom filters.

Defaults to false.

func SkipMagicBytes

func SkipMagicBytes(skip bool) FileOption

SkipMagicBytes is a file configuration option which prevents automatically reading the magic bytes when opening a parquet file, when set to true. This is useful as an optimization when programs can trust that they are dealing with parquet files and do not need to verify the first 4 bytes.

func SkipPageIndex

func SkipPageIndex(skip bool) FileOption

SkipPageIndex is a file configuration option which prevents automatically reading the page index when opening a parquet file, when set to true. This is useful as an optimization when programs know that they will not need to consume the page index.

Defaults to false.

func WithDecryption

func WithDecryption(keys KeyRetriever) FileOption

WithDecryption returns a FileOption that configures decryption with the given KeyRetriever.

type FilePages

type FilePages struct {
	// contains filtered or unexported fields
}

func (*FilePages) Close

func (f *FilePages) Close() error

Close closes the page reader.

func (*FilePages) ReadDictionary

func (f *FilePages) ReadDictionary() (Dictionary, error)

ReadDictionary returns the dictionary of the column chunk, or nil if the column chunk did not have one.

The program is not required to call this method before calling ReadPage, the dictionary is read automatically when needed. It is exposed to allow programs to access the dictionary without reading the first page.

func (*FilePages) ReadPage

func (f *FilePages) ReadPage() (Page, error)

ReadPages reads the next from from f.

func (*FilePages) SeekToRow

func (f *FilePages) SeekToRow(rowIndex int64) error

SeekToRow seeks to the given row index in the column chunk.

type FileRowGroup

type FileRowGroup struct {
	// contains filtered or unexported fields
}

FileRowGroup is an implementation of the RowGroup interface on parquet files returned by OpenFile.

func (*FileRowGroup) ColumnChunks

func (g *FileRowGroup) ColumnChunks() []ColumnChunk

ColumnChunks returns the list of column chunks in the row group.

Elements of the returned slice are guaranteed to be of type *FileColumnChunk.

func (*FileRowGroup) File

func (g *FileRowGroup) File() *File

File returns the file that this row group belongs to.

func (*FileRowGroup) NumRows

func (g *FileRowGroup) NumRows() int64

NumRows returns the number of rows in the row group.

func (*FileRowGroup) Rows

func (g *FileRowGroup) Rows() Rows

Rows returns a row reader for the row group.

func (*FileRowGroup) Schema

func (g *FileRowGroup) Schema() *Schema

Schema returns the schema of the row group.

func (*FileRowGroup) SortingColumns

func (g *FileRowGroup) SortingColumns() []SortingColumn

SortingColumns returns the list of sorting columns in the row group.

type FileView

type FileView interface {
	Metadata() *format.FileMetaData
	Schema() *Schema
	NumRows() int64
	Lookup(key string) (string, bool)
	Size() int64
	Root() *Column
	RowGroups() []RowGroup
	ColumnIndexes() []format.ColumnIndex
	OffsetIndexes() []format.OffsetIndex
}

type FixedLenByteArrayReader

type FixedLenByteArrayReader interface {
	// Read values into the byte buffer passed as argument, returning the number
	// of values written to the buffer (not the number of bytes).
	//
	// The method returns io.EOF when all values have been read.
	//
	// If the buffer was not empty, but too small to hold at least one value,
	// io.ErrShortBuffer is returned.
	ReadFixedLenByteArrays(values []byte) (int, error)
}

FixedLenByteArrayReader is an interface implemented by ValueReader instances which expose the content of a column of fixed length byte array values.

type FixedLenByteArrayWriter

type FixedLenByteArrayWriter interface {
	// Writes the fixed length byte array values.
	//
	// The size of the values is assumed to be the same as the expected size of
	// items in the column. The method errors if the length of the input values
	// is not a multiple of the expected item size.
	WriteFixedLenByteArrays(values []byte) (int, error)
}

FixedLenByteArrayWriter is an interface implemented by ValueWriter instances which support writing columns of fixed length byte array values.

type FloatReader

type FloatReader interface {
	// Read single-precision floating point values into the buffer passed as
	// argument.
	//
	// The method returns io.EOF when all values have been read.
	ReadFloats(values []float32) (int, error)
}

FloatReader is an interface implemented by ValueReader instances which expose the content of a column of single-precision floating point values.

type FloatWriter

type FloatWriter interface {
	// Write single-precision floating point values.
	//
	// The method returns the number of values written, and any error that
	// occurred while writing the values.
	WriteFloats(values []float32) (int, error)
}

FloatWriter is an interface implemented by ValueWriter instances which support writing columns of single-precision floating point values.

type GenericBuffer

type GenericBuffer[T any] struct {
	// contains filtered or unexported fields
}

GenericBuffer is similar to a Buffer but uses a type parameter to define the Go type representing the schema of rows in the buffer.

See GenericWriter for details about the benefits over the classic Buffer API.

func NewGenericBuffer

func NewGenericBuffer[T any](options ...RowGroupOption) *GenericBuffer[T]

NewGenericBuffer is like NewBuffer but returns a GenericBuffer[T] suited to write rows of Go type T.

The type parameter T should be a map, struct, or any. Any other types will cause a panic at runtime. Type checking is a lot more effective when the generic parameter is a struct type, using map and interface types is somewhat similar to using a Writer. If using an interface type for the type parameter, then providing a schema at instantiation is required.

If the option list may explicitly declare a schema, it must be compatible with the schema generated from T.

func (*GenericBuffer[T]) ColumnBuffers

func (buf *GenericBuffer[T]) ColumnBuffers() []ColumnBuffer

func (*GenericBuffer[T]) ColumnChunks

func (buf *GenericBuffer[T]) ColumnChunks() []ColumnChunk

func (*GenericBuffer[T]) Len

func (buf *GenericBuffer[T]) Len() int

func (*GenericBuffer[T]) Less

func (buf *GenericBuffer[T]) Less(i, j int) bool

func (*GenericBuffer[T]) NumRows

func (buf *GenericBuffer[T]) NumRows() int64

func (*GenericBuffer[T]) Reset

func (buf *GenericBuffer[T]) Reset()

func (*GenericBuffer[T]) Rows

func (buf *GenericBuffer[T]) Rows() Rows

func (*GenericBuffer[T]) Schema

func (buf *GenericBuffer[T]) Schema() *Schema

func (*GenericBuffer[T]) Size

func (buf *GenericBuffer[T]) Size() int64

func (*GenericBuffer[T]) SortingColumns

func (buf *GenericBuffer[T]) SortingColumns() []SortingColumn

func (*GenericBuffer[T]) Swap

func (buf *GenericBuffer[T]) Swap(i, j int)

func (*GenericBuffer[T]) Write

func (buf *GenericBuffer[T]) Write(rows []T) (int, error)

func (*GenericBuffer[T]) WriteRowGroup

func (buf *GenericBuffer[T]) WriteRowGroup(rowGroup RowGroup) (int64, error)

func (*GenericBuffer[T]) WriteRows

func (buf *GenericBuffer[T]) WriteRows(rows []Row) (int, error)

type GenericReader

type GenericReader[T any] struct {
	// contains filtered or unexported fields
}

GenericReader is similar to a Reader but uses a type parameter to define the Go type representing the schema of rows being read.

See GenericWriter for details about the benefits over the classic Reader API.

func NewGenericReader

func NewGenericReader[T any](input io.ReaderAt, options ...ReaderOption) *GenericReader[T]

NewGenericReader is like NewReader but returns GenericReader[T] suited to write rows of Go type T.

The type parameter T should be a map, struct, or any. Any other types will cause a panic at runtime. Type checking is a lot more effective when the generic parameter is a struct type, using map and interface types is somewhat similar to using a Writer.

If the option list may explicitly declare a schema, it must be compatible with the schema generated from T.

func NewGenericRowGroupReader

func NewGenericRowGroupReader[T any](rowGroup RowGroup, options ...ReaderOption) *GenericReader[T]

func (*GenericReader[T]) Close

func (r *GenericReader[T]) Close() error

func (*GenericReader[T]) File

func (r *GenericReader[T]) File() FileView

File returns a FileView of the underlying parquet file.

func (*GenericReader[T]) NumRows

func (r *GenericReader[T]) NumRows() int64

func (*GenericReader[T]) Read

func (r *GenericReader[T]) Read(rows []T) (int, error)

Read reads the next rows from the reader into the given rows slice up to len(rows).

The returned values are safe to reuse across Read calls and do not share memory with the reader's underlying page buffers.

The method returns the number of rows read and io.EOF when no more rows can be read from the reader.

func (*GenericReader[T]) ReadRows

func (r *GenericReader[T]) ReadRows(rows []Row) (int, error)

func (*GenericReader[T]) Reset

func (r *GenericReader[T]) Reset()

func (*GenericReader[T]) Schema

func (r *GenericReader[T]) Schema() *Schema

func (*GenericReader[T]) SeekToRow

func (r *GenericReader[T]) SeekToRow(rowIndex int64) error

type GenericWriter

type GenericWriter[T any] struct {
	// contains filtered or unexported fields
}

GenericWriter is similar to a Writer but uses a type parameter to define the Go type representing the schema of rows being written.

Using this type over Writer has multiple advantages:

  • By leveraging type information, the Go compiler can provide greater guarantees that the code is correct. For example, the parquet.Writer.Write method accepts an argument of type interface{}, which delays type checking until runtime. The parquet.GenericWriter[T].Write method ensures at compile time that the values it receives will be of type T, reducing the risk of introducing errors.

  • Since type information is known at compile time, the implementation of parquet.GenericWriter[T] can make safe assumptions, removing the need for runtime validation of how the parameters are passed to its methods. Optimizations relying on type information are more effective, some of the writer's state can be precomputed at initialization, which was not possible with parquet.Writer.

  • The parquet.GenericWriter[T].Write method uses a data-oriented design, accepting an slice of T instead of a single value, creating more opportunities to amortize the runtime cost of abstractions. This optimization is not available for parquet.Writer because its Write method's argument would be of type []interface{}, which would require conversions back and forth from concrete types to empty interfaces (since a []T cannot be interpreted as []interface{} in Go), would make the API more difficult to use and waste compute resources in the type conversions, defeating the purpose of the optimization in the first place.

Note that this type is only available when compiling with Go 1.18 or later.

func NewGenericWriter

func NewGenericWriter[T any](output io.Writer, options ...WriterOption) *GenericWriter[T]

NewGenericWriter is like NewWriter but returns a GenericWriter[T] suited to write rows of Go type T.

The type parameter T should be a map, struct, or any. Any other types will cause a panic at runtime. Type checking is a lot more effective when the generic parameter is a struct type, using map and interface types is somewhat similar to using a Writer.

If the option list may explicitly declare a schema, it must be compatible with the schema generated from T.

Sorting columns may be set on the writer to configure the generated row groups metadata. However, rows are always written in the order they were seen, no reordering is performed, the writer expects the application to ensure proper correlation between the order of rows and the list of sorting columns. See SortingWriter[T] for a writer which handles reordering rows based on the configured sorting columns.

func (*GenericWriter[T]) BeginRowGroup

func (w *GenericWriter[T]) BeginRowGroup() *ConcurrentRowGroupWriter

BeginRowGroup returns a new ConcurrentRowGroupWriter that can be written to in parallel with other row groups. However these need to be committed back to the writer serially using the Commit method on the row group.

Example usage could look something like:

writer := parquet.NewGenericWriter[any](...)
rgs := make([]*parquet.ConcurrentRowGroupWriter, 5)
var wg sync.WaitGroup
for i := range rgs {
  rg := writer.BeginRowGroup()
  rgs[i] = rg
  wg.Add(1)
  go func() {
    defer wg.Done()
    writeChunkRows(i, rg)
  }()
}
wg.Wait()
for _, rg := range rgs {
  if _, err := rg.Commit(); err != nil {
    return err
  }
}
return writer.Close()

func (*GenericWriter[T]) Close

func (w *GenericWriter[T]) Close() error

Close must be called after all values were produced to the writer in order to flush all buffers and write the parquet footer. The writer can only be reused if Reset is called first. Failure to do so will result in defined behavior.

func (*GenericWriter[T]) ColumnWriters

func (w *GenericWriter[T]) ColumnWriters() []*ColumnWriter

func (*GenericWriter[T]) File

func (w *GenericWriter[T]) File() FileView

File returns a FileView of the written parquet file. Only available after Close is called.

func (*GenericWriter[T]) Flush

func (w *GenericWriter[T]) Flush() error

func (*GenericWriter[T]) ReadRowsFrom

func (w *GenericWriter[T]) ReadRowsFrom(rows RowReader) (int64, error)

func (*GenericWriter[T]) Reset

func (w *GenericWriter[T]) Reset(output io.Writer)

Reset clears the state of the writer without flushing any of the buffers, and setting the output to the io.Writer passed as argument, allowing the writer to be reused to produce another parquet file.

func (*GenericWriter[T]) Schema

func (w *GenericWriter[T]) Schema() *Schema

func (*GenericWriter[T]) SetKeyValueMetadata

func (w *GenericWriter[T]) SetKeyValueMetadata(key, value string)

SetKeyValueMetadata sets a key/value pair in the Parquet file metadata.

Keys are assumed to be unique, if the same key is repeated multiple times the last value is retained. While the parquet format does not require unique keys, this design decision was made to optimize for the most common use case where applications leverage this extension mechanism to associate single values to keys. This may create incompatibilities with other parquet libraries, or may cause some key/value pairs to be lost when open parquet files written with repeated keys. We can revisit this decision if it ever becomes a blocker.

func (*GenericWriter[T]) Size

func (w *GenericWriter[T]) Size() int64

Size returns an estimate of the current file size in bytes. See Writer.Size for details.

func (*GenericWriter[T]) Write

func (w *GenericWriter[T]) Write(rows []T) (written int, err error)

func (*GenericWriter[T]) WriteRowGroup

func (w *GenericWriter[T]) WriteRowGroup(rowGroup RowGroup) (int64, error)

func (*GenericWriter[T]) WriteRows

func (w *GenericWriter[T]) WriteRows(rows []Row) (int, error)

type Group

type Group map[string]Node

func (Group) Compression

func (g Group) Compression() compress.Codec

func (Group) Encoding

func (g Group) Encoding() encoding.Encoding

func (Group) Fields

func (g Group) Fields() []Field

func (Group) GoType

func (g Group) GoType() reflect.Type

func (Group) ID

func (g Group) ID() int

func (Group) Leaf

func (g Group) Leaf() bool

func (Group) Optional

func (g Group) Optional() bool

func (Group) Repeated

func (g Group) Repeated() bool

func (Group) Required

func (g Group) Required() bool

func (Group) String

func (g Group) String() string

func (Group) Type

func (g Group) Type() Type

type Int32Reader

type Int32Reader interface {
	// Read 32 bits integer values into the buffer passed as argument.
	//
	// The method returns io.EOF when all values have been read.
	ReadInt32s(values []int32) (int, error)
}

Int32Reader is an interface implemented by ValueReader instances which expose the content of a column of int32 values.

type Int32Writer

type Int32Writer interface {
	// Write 32 bits signed integer values.
	//
	// The method returns the number of values written, and any error that
	// occurred while writing the values.
	WriteInt32s(values []int32) (int, error)
}

Int32Writer is an interface implemented by ValueWriter instances which support writing columns of 32 bits signed integer values.

type Int64Reader

type Int64Reader interface {
	// Read 64 bits integer values into the buffer passed as argument.
	//
	// The method returns io.EOF when all values have been read.
	ReadInt64s(values []int64) (int, error)
}

Int64Reader is an interface implemented by ValueReader instances which expose the content of a column of int64 values.

type Int64Writer

type Int64Writer interface {
	// Write 64 bits signed integer values.
	//
	// The method returns the number of values written, and any error that
	// occurred while writing the values.
	WriteInt64s(values []int64) (int, error)
}

Int64Writer is an interface implemented by ValueWriter instances which support writing columns of 64 bits signed integer values.

type Int96Reader

type Int96Reader interface {
	// Read 96 bits integer values into the buffer passed as argument.
	//
	// The method returns io.EOF when all values have been read.
	ReadInt96s(values []deprecated.Int96) (int, error)
}

Int96Reader is an interface implemented by ValueReader instances which expose the content of a column of int96 values.

type Int96Writer

type Int96Writer interface {
	// Write 96 bits signed integer values.
	//
	// The method returns the number of values written, and any error that
	// occurred while writing the values.
	WriteInt96s(values []deprecated.Int96) (int, error)
}

Int96Writer is an interface implemented by ValueWriter instances which support writing columns of 96 bits signed integer values.

type Interval

type Interval struct {
	Months       uint32
	Days         uint32
	Milliseconds uint32
}

Interval represents a Parquet INTERVAL value.

The physical representation is a FIXED_LEN_BYTE_ARRAY(12) containing three little-endian unsigned 32-bit integers: months, days, and milliseconds.

https://github.com/apache/parquet-format/blob/master/LogicalTypes.md#interval

type KeyRetriever

type KeyRetriever interface {
	// FooterKey returns the AES key for the file footer.
	// keyMetadata is the optional bytes stored in FileCryptoMetaData.KeyMetadata
	// (may be nil).
	FooterKey(keyMetadata []byte) ([]byte, error)

	// ColumnKey returns the AES key for an encrypted column.
	// path is the column's path in the schema; keyMetadata is the optional bytes
	// stored in EncryptionWithColumnKey.KeyMetadata (may be nil).
	//
	// To signal that a particular column's key is intentionally unavailable
	// (e.g. the caller only holds keys for a subset of columns), return an
	// error that wraps ErrKeyNotFound:
	//
	//	return nil, fmt.Errorf("no key for %v: %w", path, parquet.ErrKeyNotFound)
	//
	// OpenFile treats ErrKeyNotFound as non-fatal and leaves that column
	// inaccessible; any other non-nil error is propagated as a hard failure.
	ColumnKey(path []string, keyMetadata []byte) ([]byte, error)
}

KeyRetriever resolves AES keys from the metadata bytes stored in the file. Implement this interface to supply keys when opening an encrypted parquet file.

type Kind

type Kind int8

Kind is an enumeration type representing the physical types supported by the parquet type system.

const (
	Boolean           Kind = Kind(format.Boolean)
	Int32             Kind = Kind(format.Int32)
	Int64             Kind = Kind(format.Int64)
	Int96             Kind = Kind(format.Int96)
	Float             Kind = Kind(format.Float)
	Double            Kind = Kind(format.Double)
	ByteArray         Kind = Kind(format.ByteArray)
	FixedLenByteArray Kind = Kind(format.FixedLenByteArray)
)

func (Kind) String

func (k Kind) String() string

String returns a human-readable representation of the physical type.

func (Kind) Value

func (k Kind) Value(v []byte) Value

Value constructs a value from k and v.

The method panics if the data is not a valid representation of the value kind; for example, if the kind is Int32 but the data is not 4 bytes long.

type LeafColumn

type LeafColumn struct {
	Node               Node
	Path               []string
	ColumnIndex        int
	MaxRepetitionLevel int
	MaxDefinitionLevel int
}

LeafColumn is a struct type representing leaf columns of a parquet schema.

type Node

type Node interface {
	// The id of this node in its parent node. Zero value is treated as id is not
	// set. ID only needs to be unique within its parent context.
	//
	// This is the same as parquet field_id
	ID() int

	// Returns a human-readable representation of the parquet node.
	String() string

	// For leaf nodes, returns the type of values of the parquet column.
	//
	// Calling this method on non-leaf nodes will panic.
	Type() Type

	// Returns whether the parquet column is optional.
	Optional() bool

	// Returns whether the parquet column is repeated.
	Repeated() bool

	// Returns whether the parquet column is required.
	Required() bool

	// Returns true if this a leaf node.
	Leaf() bool

	// Returns a mapping of the node's fields.
	//
	// As an optimization, the same slices may be returned by multiple calls to
	// this method, programs must treat the returned values as immutable.
	//
	// This method returns an empty mapping when called on leaf nodes.
	Fields() []Field

	// Returns the encoding used by the node.
	//
	// The method may return nil to indicate that no specific encoding was
	// configured on the node, in which case a default encoding might be used.
	Encoding() encoding.Encoding

	// Returns compression codec used by the node.
	//
	// The method may return nil to indicate that no specific compression codec
	// was configured on the node, in which case a default compression might be
	// used.
	Compression() compress.Codec

	// Returns the Go type that best represents the parquet node.
	//
	// For leaf nodes, this will be one of bool, int32, int64, deprecated.Int96,
	// float32, float64, string, []byte, or [N]byte.
	//
	// For groups, the method returns a struct type.
	//
	// If the method is called on a repeated node, the method returns a slice of
	// the underlying type.
	//
	// For optional nodes, the method returns a pointer of the underlying type.
	//
	// For nodes that were constructed from Go values (e.g. using SchemaOf), the
	// method returns the original Go type.
	GoType() reflect.Type
}

Node values represent nodes of a parquet schema.

Nodes carry the type of values, as well as properties like whether the values are optional or repeat. Nodes with one or more children represent parquet groups and therefore do not have a logical type.

Nodes are immutable values and therefore safe to use concurrently from multiple goroutines.

func Compressed

func Compressed(node Node, codec compress.Codec) Node

Compressed wraps the node passed as argument to use the given compression codec.

If the codec is nil, the node's compression is left unchanged.

The function panics if it is called on a non-leaf node.

func Decimal

func Decimal(scale, precision int, typ Type) Node

Decimal constructs a leaf node of decimal logical type with the given scale, precision, and underlying type.

https://github.com/apache/parquet-format/blob/master/LogicalTypes.md#decimal

func Encoded

func Encoded(node Node, encoding encoding.Encoding) Node

Encoded wraps the node passed as argument to use the given encoding.

The function panics if it is called on a non-leaf node, or if the encoding does not support the node type.

func Enum

func Enum() Node

Enum constructs a leaf node with a logical type representing enumerations.

https://github.com/apache/parquet-format/blob/master/LogicalTypes.md#enum

func FieldID

func FieldID(node Node, id int) Node

FieldID wraps a node to provide node field id

func Geography

func Geography(crs string, algorithm format.EdgeInterpolationAlgorithm) Node

func Geometry

func Geometry(crs string) Node

func Int

func Int(bitWidth int) Node

Int constructs a leaf node of signed integer logical type of the given bit width.

The bit width must be one of 8, 16, 32, 64, or the function will panic.

func IntervalNode

func IntervalNode() Node

IntervalNode constructs a leaf node of INTERVAL logical type.

https://github.com/apache/parquet-format/blob/master/LogicalTypes.md#interval

func Leaf

func Leaf(typ Type) Node

Leaf returns a leaf node of the given type.

func MergeNodes

func MergeNodes(nodes ...Node) Node

MergeNodes takes a list of nodes and greedily retains properties of the schemas: - keeps last compression that is not nil - keeps last non-plain encoding that is not nil - keeps last non-zero field id - union of all columns for group nodes - retains the most permissive repetition (required < optional < repeated)

func Optional

func Optional(node Node) Node

Optional wraps the given node to make it optional.

func Repeated

func Repeated(node Node) Node

Repeated wraps the given node to make it repeated.

func Required

func Required(node Node) Node

Required wraps the given node to make it required.

func ShreddedVariant

func ShreddedVariant(shreddedType Node) (Node, error)

ShreddedVariant constructs a node of shredded VARIANT logical type. It is a group with a required byte array "metadata" field, an optional byte array "value" field (for any unshredded values), and an optional "typed_value" field whose type is that of the shredded value. If the given node is a group or list (or contains a list), the resulting "typed_value" field will have some additional structure to allow each group or element in a list to have a mix of shredded and unshredded data.

The given node may only contain types that map to valid variant value types. Therefore, it may not contain ENUM, FLOAT16, INTERVAL, JSON, BSON, VARIANT, GEOMETRY, GEOGRAPHY, MAP, or UNKNOWN logical types. It may only use signed INT logical types. It may only use DECIMAL logical types whose precision is less than or equal to 38. It may not contain any repeated fields unless they are the middle level of a 3-level LIST logical type. Any "required" settings on fields will be ignored: shredded fields must always be optional to represent values that may not conform to the shredded type. It also may not contain empty groups: any groups must have at least one field.

More information on shredded variants can be found in the Parquet documentation.

*Experimental*: Support for the VARIANT type is still being developed and subject to change.

Go values written to a shredded variant column with the "variant" struct tag are shredded per the Parquet documentation: values matching the shredded type are stored in the typed_value columns, and everything else is encoded into the appropriate value column. Reading reconstructs the variant from all columns.

Example

ExampleShreddedVariant demonstrates reading and writing shredded variant data. When a Go value matches the typed_value schema, it is stored efficiently in the typed column. Otherwise, it falls back to the variant-encoded value column.

package main

import (
	"bytes"
	"fmt"

	"github.com/AlexAslan/parquet-go"
)

func main() {
	type Record struct {
		ID   int64 `parquet:"id"`
		Data any   `parquet:"data"`
	}

	shreddedType, _ := parquet.ShreddedVariant(parquet.String())
	schema := parquet.NewSchema("Record", parquet.Group{
		"id":   parquet.Int(64),
		"data": shreddedType,
	})

	buf := new(bytes.Buffer)
	writer := parquet.NewGenericWriter[Record](buf, schema)
	_, _ = writer.Write([]Record{
		{ID: 1, Data: "typed-value"},          // Shredded: stored in typed_value column
		{ID: 2, Data: int32(42)},              // Not shredded: stored in value column
		{ID: 3, Data: map[string]any{"k": 1}}, // Not shredded: stored in value column
	})
	_ = writer.Close()

	reader := parquet.NewGenericReader[Record](bytes.NewReader(buf.Bytes()), schema)
	records := make([]Record, 3)
	n, _ := reader.Read(records)
	_ = reader.Close()

	for _, r := range records[:n] {
		fmt.Printf("Record %d: %v (%T)\n", r.ID, r.Data, r.Data)
	}
}
Output:
Record 1: typed-value (string)
Record 2: 42 (int32)
Record 3: map[k:1] (map[string]interface {})

func Time

func Time(unit TimeUnit) Node

Time constructs a leaf node of TIME logical type. IsAdjustedToUTC is true by default.

https://github.com/apache/parquet-format/blob/master/LogicalTypes.md#time

func TimeAdjusted

func TimeAdjusted(unit TimeUnit, isAdjustedToUTC bool) Node

TimeAdjusted constructs a leaf node of TIME logical type with the IsAdjustedToUTC property explicitly set.

https://github.com/apache/parquet-format/blob/master/LogicalTypes.md#time

func Timestamp

func Timestamp(unit TimeUnit) Node

Timestamp constructs of leaf node of TIMESTAMP logical type. IsAdjustedToUTC is true by default.

https://github.com/apache/parquet-format/blob/master/LogicalTypes.md#timestamp

func TimestampAdjusted

func TimestampAdjusted(unit TimeUnit, isAdjustedToUTC bool) Node

TimestampAdjusted constructs a leaf node of TIMESTAMP logical type with the IsAdjustedToUTC property explicitly set.

https://github.com/apache/parquet-format/blob/master/LogicalTypes.md#time

func Uint

func Uint(bitWidth int) Node

Uint constructs a leaf node of unsigned integer logical type of the given bit width.

The bit width must be one of 8, 16, 32, 64, or the function will panic.

func Variant

func Variant() Node

Variant constructs a node of unshredded VARIANT logical type. It is a group with two required fields, "metadata" and "value", both byte arrays.

*Experimental*: Support for the VARIANT type is still being developed and subject to change.

Go values written to a variant column with the "variant" struct tag are encoded with the variant binary encoding (see the variant subpackage); structs with Metadata and Value []byte fields pass the raw encoding through unchanged. When reading, files that store the column shredded are reconstructed automatically.

Example

ExampleVariant demonstrates reading and writing unshredded variant data using Go structs with the "variant" struct tag. The variant field accepts any Go value, which is automatically marshaled to/from the Parquet variant binary format.

package main

import (
	"bytes"
	"fmt"

	"github.com/AlexAslan/parquet-go"
)

func main() {
	type Event struct {
		ID   int64 `parquet:"id"`
		Data any   `parquet:"data,variant"`
	}

	buf := new(bytes.Buffer)
	writer := parquet.NewGenericWriter[Event](buf)
	_, _ = writer.Write([]Event{
		{ID: 1, Data: "hello"},
		{ID: 2, Data: int32(42)},
		{ID: 3, Data: map[string]any{"key": "value"}},
	})
	_ = writer.Close()

	reader := parquet.NewGenericReader[Event](bytes.NewReader(buf.Bytes()))
	events := make([]Event, 3)
	_, _ = reader.Read(events)
	_ = reader.Close()

	for _, e := range events {
		fmt.Printf("Event %d: %v (%T)\n", e.ID, e.Data, e.Data)
	}
}
Output:
Event 1: hello (string)
Event 2: 42 (int32)
Event 3: map[key:value] (map[string]interface {})
Example (Schema)

ExampleVariant_schema demonstrates creating and printing variant schemas.

package main

import (
	"fmt"

	"github.com/AlexAslan/parquet-go"
)

func main() {
	// Unshredded variant: group with metadata + value byte arrays.
	unshredded := parquet.NewSchema("Unshredded", parquet.Group{
		"data": parquet.Optional(parquet.Variant()),
	})
	fmt.Println(unshredded)

	// Shredded variant: adds a typed_value column for efficient typed access.
	shreddedType, _ := parquet.ShreddedVariant(parquet.String())
	shredded := parquet.NewSchema("Shredded", parquet.Group{
		"data": parquet.Optional(shreddedType),
	})
	fmt.Println(shredded)
}
Output:
message Unshredded {
	optional group data (VARIANT) {
		required binary metadata;
		required binary value;
	}
}
message Shredded {
	optional group data (VARIANT) {
		required binary metadata;
		optional binary value;
		optional binary typed_value (STRING);
	}
}

type OffsetIndex

type OffsetIndex interface {
	// NumPages returns the number of pages in the offset index.
	NumPages() int

	// Offset returns the offset starting from the beginning of the file for the
	// page at the given index.
	Offset(int) int64

	// CompressedPageSize returns the size of the page at the given index
	// (in bytes).
	CompressedPageSize(int) int64

	// FirstRowIndex returns the the first row in the page at the given index.
	//
	// The returned row index is based on the row group that the page belongs
	// to, the first row has index zero.
	FirstRowIndex(int) int64
}

type Page

type Page interface {
	// Returns the type of values read from this page.
	//
	// The returned type can be used to encode the page data, in the case of
	// an indexed page (which has a dictionary), the type is configured to
	// encode the indexes stored in the page rather than the plain values.
	Type() Type

	// Returns the column index that this page belongs to.
	Column() int

	// If the page contains indexed values, calling this method returns the
	// dictionary in which the values are looked up. Otherwise, the method
	// returns nil.
	Dictionary() Dictionary

	// Returns the number of rows, values, and nulls in the page. The number of
	// rows may be less than the number of values in the page if the page is
	// part of a repeated column.
	NumRows() int64
	NumValues() int64
	NumNulls() int64

	// Returns the page's min and max values.
	//
	// The third value is a boolean indicating whether the page bounds were
	// available. Page bounds may not be known if the page contained no values
	// or only nulls, or if they were read from a parquet file which had neither
	// page statistics nor a page index.
	Bounds() (min, max Value, ok bool)

	// Returns the size of the page in bytes (uncompressed).
	Size() int64

	// Returns a reader exposing the values contained in the page.
	//
	// Depending on the underlying implementation, the returned reader may
	// support reading an array of typed Go values by implementing interfaces
	// like parquet.Int32Reader. Applications should use type assertions on
	// the returned reader to determine whether those optimizations are
	// available.
	//
	// In the data page format version 1, it wasn't specified whether pages
	// must start with a new row. Legacy writers have produced parquet files
	// where row values were overlapping between two consecutive pages.
	// As a result, the values read must not be assumed to start at the
	// beginning of a row, unless the program knows that it is only working
	// with parquet files that used the data page format version 2 (which is
	// the default behavior for parquet-go).
	Values() ValueReader

	// Returns a new page which is as slice of the receiver between row indexes
	// i and j.
	Slice(i, j int64) Page

	// Expose the lists of repetition and definition levels of the page.
	//
	// The returned slices may be empty when the page has no repetition or
	// definition levels.
	RepetitionLevels() []byte
	DefinitionLevels() []byte

	// Returns the in-memory buffer holding the page values.
	//
	// The intent is for the returned value to be used as input parameter when
	// calling the Encode method of the associated Type.
	//
	// The slices referenced by the encoding.Values may be the same across
	// multiple calls to this method, applications must treat the content as
	// immutable.
	Data() encoding.Values
}

Page values represent sequences of parquet values. From the Parquet documentation: "Column chunks are a chunk of the data for a particular column. They live in a particular row group and are guaranteed to be contiguous in the file. Column chunks are divided up into pages. A page is conceptually an indivisible unit (in terms of compression and encoding). There can be multiple page types which are interleaved in a column chunk."

https://github.com/apache/parquet-format#glossary

type PageHeader interface {
	// Returns the number of values in the page (including nulls).
	NumValues() int64

	// Returns the page encoding.
	Encoding() format.Encoding

	// Returns the parquet format page type.
	PageType() format.PageType
}

PageHeader is an interface implemented by parquet page headers.

type PageReader

type PageReader interface {
	// Reads and returns the next page from the sequence. When all pages have
	// been read, or if the sequence was closed, the method returns io.EOF.
	ReadPage() (Page, error)
}

PageReader is an interface implemented by types that support producing a sequence of pages.

type PageWriter

type PageWriter interface {
	WritePage(Page) (int64, error)
}

PageWriter is an interface implemented by types that support writing pages to an underlying storage medium.

type Pages

type Pages interface {
	PageReader
	RowSeeker
	io.Closer
}

Pages is an interface implemented by page readers returned by calling the Pages method of ColumnChunk instances.

func AsyncPages

func AsyncPages(pages Pages) Pages

AsyncPages wraps the given Pages instance to perform page reads asynchronously in a separate goroutine.

Performing page reads asynchronously is important when the application may be reading pages from a high latency backend, and the last page read may be processed while initiating reading of the next page.

type ReadMode

type ReadMode int

ReadMode is an enum that is used to configure the way that a File reads pages.

const (
	ReadModeSync  ReadMode = iota // ReadModeSync reads pages synchronously on demand (Default).
	ReadModeAsync                 // ReadModeAsync reads pages asynchronously in the background.
)

type Reader deprecated

type Reader struct {
	// contains filtered or unexported fields
}

Deprecated: A Reader reads Go values from parquet files.

This example showcases a typical use of parquet readers:

reader := parquet.NewReader(file)
rows := []RowType{}
for {
	row := RowType{}
	err := reader.Read(&row)
	if err != nil {
		if err == io.EOF {
			break
		}
		...
	}
	rows = append(rows, row)
}
if err := reader.Close(); err != nil {
	...
}

For programs building with Go 1.18 or later, the GenericReader[T] type supersedes this one.

func NewReader

func NewReader(input io.ReaderAt, options ...ReaderOption) *Reader

NewReader constructs a parquet reader reading rows from the given io.ReaderAt.

In order to read parquet rows, the io.ReaderAt must be converted to a parquet.File. If r is already a parquet.File it is used directly; otherwise, the io.ReaderAt value is expected to either have a `Size() int64` method or implement io.Seeker in order to determine its size.

The function panics if the reader configuration is invalid. Programs that cannot guarantee the validity of the options passed to NewReader should construct the reader configuration independently prior to calling this function:

config, err := parquet.NewReaderConfig(options...)
if err != nil {
	// handle the configuration error
	...
} else {
	// this call to create a reader is guaranteed not to panic
	reader := parquet.NewReader(input, config)
	...
}

func NewRowGroupReader

func NewRowGroupReader(rowGroup RowGroup, options ...ReaderOption) *Reader

NewRowGroupReader constructs a new Reader which reads rows from the RowGroup passed as argument.

func (*Reader) Close

func (r *Reader) Close() error

Close closes the reader, preventing more rows from being read.

func (*Reader) File

func (r *Reader) File() FileView

File returns a FileView of the parquet file being read. Only available if Reader was created with a File.

func (*Reader) NumRows

func (r *Reader) NumRows() int64

NumRows returns the number of rows that can be read from r.

func (*Reader) Read

func (r *Reader) Read(row any) error

Read reads the next row from r. The type of the row must match the schema of the underlying parquet file or an error will be returned.

The method returns io.EOF when no more rows can be read from r.

func (*Reader) ReadRows

func (r *Reader) ReadRows(rows []Row) (int, error)

ReadRows reads the next rows from r into the given Row buffer.

The returned values are laid out in the order expected by the parquet.(*Schema).Reconstruct method.

The method returns io.EOF when no more rows can be read from r.

func (*Reader) Reset

func (r *Reader) Reset()

Reset repositions the reader at the beginning of the underlying parquet file.

func (*Reader) Schema

func (r *Reader) Schema() *Schema

Schema returns the schema of rows read by r.

func (*Reader) SeekToRow

func (r *Reader) SeekToRow(rowIndex int64) error

SeekToRow positions r at the given row index.

type ReaderConfig

type ReaderConfig struct {
	Schema       *Schema
	SchemaConfig *SchemaConfig
}

The ReaderConfig type carries configuration options for parquet readers.

ReaderConfig implements the ReaderOption interface so it can be used directly as argument to the NewReader function when needed, for example:

reader := parquet.NewReader(output, schema, &parquet.ReaderConfig{
	// ...
})

func DefaultReaderConfig

func DefaultReaderConfig() *ReaderConfig

DefaultReaderConfig returns a new ReaderConfig value initialized with the default reader configuration.

func NewReaderConfig

func NewReaderConfig(options ...ReaderOption) (*ReaderConfig, error)

NewReaderConfig constructs a new reader configuration applying the options passed as arguments.

The function returns an non-nil error if some of the options carried invalid configuration values.

func (*ReaderConfig) Apply

func (c *ReaderConfig) Apply(options ...ReaderOption)

Apply applies the given list of options to c.

func (*ReaderConfig) ConfigureReader

func (c *ReaderConfig) ConfigureReader(config *ReaderConfig)

ConfigureReader applies configuration options from c to config.

func (*ReaderConfig) Validate

func (c *ReaderConfig) Validate() error

Validate returns a non-nil error if the configuration of c is invalid.

type ReaderOption

type ReaderOption interface {
	ConfigureReader(*ReaderConfig)
}

ReaderOption is an interface implemented by types that carry configuration options for parquet readers.

type Row

type Row []Value

Row represents a parquet row as a slice of values.

Each value should embed a column index, repetition level, and definition level allowing the program to determine how to reconstruct the original object from the row.

func AppendRow

func AppendRow(row Row, columns ...[]Value) Row

AppendRow appends to row the given list of column values.

AppendRow can be used to construct a Row value from columns, while retaining the underlying memory buffer to avoid reallocation; for example:

The function panics if the column indexes of values in each column do not match their position in the argument list.

func MakeRow

func MakeRow(columns ...[]Value) Row

MakeRow constructs a Row from a list of column values.

The function panics if the column indexes of values in each column do not match their position in the argument list.

func (Row) Clone

func (row Row) Clone() Row

Clone creates a copy of the row which shares no pointers.

This method is useful to capture rows after a call to RowReader.ReadRows when values need to be retained before the next call to ReadRows or after the lifespan of the reader.

func (Row) Equal

func (row Row) Equal(other Row) bool

Equal returns true if row and other contain the same sequence of values.

func (Row) Range

func (row Row) Range(f func(columnIndex int, columnValues []Value) bool)

Range calls f for each column of row.

type RowBuffer

type RowBuffer[T any] struct {
	// contains filtered or unexported fields
}

RowBuffer is an implementation of the RowGroup interface which stores parquet rows in memory.

Unlike GenericBuffer which uses a column layout to store values in memory buffers, RowBuffer uses a row layout. The use of row layout provides greater efficiency when sorting the buffer, which is the primary use case for the RowBuffer type. Applications which intend to sort rows prior to writing them to a parquet file will often see lower CPU utilization from using a RowBuffer than a GenericBuffer.

RowBuffer values are not safe to use concurrently from multiple goroutines.

func NewRowBuffer

func NewRowBuffer[T any](options ...RowGroupOption) *RowBuffer[T]

NewRowBuffer constructs a new row buffer.

func (*RowBuffer[T]) ColumnChunks

func (buf *RowBuffer[T]) ColumnChunks() []ColumnChunk

ColumnChunks returns a view of the buffer's columns.

Note that reading columns of a RowBuffer will be less efficient than reading columns of a GenericBuffer since the latter uses a column layout. This method is mainly exposed to satisfy the RowGroup interface, applications which need compute-efficient column scans on in-memory buffers should likely use a GenericBuffer instead.

The returned column chunks are snapshots at the time the method is called, they remain valid until the next call to Reset on the buffer.

func (*RowBuffer[T]) Len

func (buf *RowBuffer[T]) Len() int

Len returns the number of rows in the buffer.

The method contributes to satisfying sort.Interface.

func (*RowBuffer[T]) Less

func (buf *RowBuffer[T]) Less(i, j int) bool

Less compares the rows at index i and j according to the sorting columns configured on the buffer.

The method contributes to satisfying sort.Interface.

func (*RowBuffer[T]) NumRows

func (buf *RowBuffer[T]) NumRows() int64

NumRows returns the number of rows currently written to the buffer.

func (*RowBuffer[T]) Reset

func (buf *RowBuffer[T]) Reset()

Reset clears the content of the buffer without releasing its memory.

func (*RowBuffer[T]) Rows

func (buf *RowBuffer[T]) Rows() Rows

Rows returns a Rows instance exposing rows stored in the buffer.

The rows returned are a snapshot at the time the method is called. The returned rows and values read from it remain valid until the next call to Reset on the buffer.

func (*RowBuffer[T]) Schema

func (buf *RowBuffer[T]) Schema() *Schema

Schema returns the schema of rows in the buffer.

func (*RowBuffer[T]) SortingColumns

func (buf *RowBuffer[T]) SortingColumns() []SortingColumn

SortingColumns returns the list of columns that rows are expected to be sorted by.

The list of sorting columns is configured when the buffer is created and used when it is sorted.

Note that unless the buffer is explicitly sorted, there are no guarantees that the rows it contains will be in the order specified by the sorting columns.

func (*RowBuffer[T]) Swap

func (buf *RowBuffer[T]) Swap(i, j int)

Swap exchanges the rows at index i and j in the buffer.

The method contributes to satisfying sort.Interface.

func (*RowBuffer[T]) Write

func (buf *RowBuffer[T]) Write(rows []T) (int, error)

Write writes rows to the buffer, returning the number of rows written.

func (*RowBuffer[T]) WriteRows

func (buf *RowBuffer[T]) WriteRows(rows []Row) (int, error)

WriteRows writes parquet rows to the buffer, returing the number of rows written.

type RowBuilder

type RowBuilder struct {
	// contains filtered or unexported fields
}

RowBuilder is a type which helps build parquet rows incrementally by adding values to columns.

Example
package main

import (
	"fmt"

	"github.com/AlexAslan/parquet-go"
)

func main() {
	builder := parquet.NewRowBuilder(parquet.Group{
		"birth_date": parquet.Optional(parquet.Date()),
		"first_name": parquet.String(),
		"last_name":  parquet.String(),
	})

	builder.Add(1, parquet.ByteArrayValue([]byte("Luke")))
	builder.Add(2, parquet.ByteArrayValue([]byte("Skywalker")))

	row := builder.Row()
	row.Range(func(columnIndex int, columnValues []parquet.Value) bool {
		fmt.Printf("%+v\n", columnValues[0])
		return true
	})

}
Output:
C:0 D:0 R:0 V:<null>
C:1 D:0 R:0 V:Luke
C:2 D:0 R:0 V:Skywalker

func NewRowBuilder

func NewRowBuilder(schema Node) *RowBuilder

NewRowBuilder constructs a RowBuilder which builds rows for the parquet schema passed as argument.

func (*RowBuilder) Add

func (b *RowBuilder) Add(columnIndex int, columnValue Value)

Add adds columnValue to the column at columnIndex.

func (*RowBuilder) AppendRow

func (b *RowBuilder) AppendRow(row Row) Row

AppendRow appends the current state of b to row and returns it.

func (*RowBuilder) Next

func (b *RowBuilder) Next(columnIndex int)

Next must be called to indicate the start of a new repeated record for the column at the given index.

If the column index is part of a repeated group, the builder automatically starts a new record for all adjacent columns, the application does not need to call this method for each column of the repeated group.

Next must be called after adding a sequence of records.

func (*RowBuilder) Reset

func (b *RowBuilder) Reset()

Reset clears the internal state of b, making it possible to reuse while retaining the internal buffers.

func (*RowBuilder) Row

func (b *RowBuilder) Row() Row

Row materializes the current state of b into a parquet row.

type RowGroup

type RowGroup interface {
	// Returns the number of rows in the group.
	NumRows() int64

	// Returns the list of column chunks in this row group. The chunks are
	// ordered in the order of leaf columns from the row group's schema.
	//
	// If the underlying implementation is not read-only, the returned
	// parquet.ColumnChunk may implement other interfaces: for example,
	// parquet.ColumnBuffer if the chunk is backed by an in-memory buffer,
	// or typed writer interfaces like parquet.Int32Writer depending on the
	// underlying type of values that can be written to the chunk.
	//
	// As an optimization, the row group may return the same slice across
	// multiple calls to this method. Applications should treat the returned
	// slice as read-only.
	ColumnChunks() []ColumnChunk

	// Returns the schema of rows in the group.
	Schema() *Schema

	// Returns the list of sorting columns describing how rows are sorted in the
	// group.
	//
	// The method will return an empty slice if the rows are not sorted.
	SortingColumns() []SortingColumn

	// Returns a reader exposing the rows of the row group.
	//
	// As an optimization, the returned parquet.Rows object may implement
	// parquet.RowWriterTo, and test the RowWriter it receives for an
	// implementation of the parquet.RowGroupWriter interface.
	//
	// This optimization mechanism is leveraged by the parquet.CopyRows function
	// to skip the generic row-by-row copy algorithm and delegate the copy logic
	// to the parquet.Rows object.
	Rows() Rows
}

RowGroup is an interface representing a parquet row group. From the Parquet docs, a RowGroup is "a logical horizontal partitioning of the data into rows. There is no physical structure that is guaranteed for a row group. A row group consists of a column chunk for each column in the dataset."

https://github.com/apache/parquet-format#glossary

func AsyncRowGroup

func AsyncRowGroup(base RowGroup) RowGroup

func ConvertRowGroup

func ConvertRowGroup(rowGroup RowGroup, conv Conversion) RowGroup

ConvertRowGroup constructs a wrapper of the given row group which applies the given schema conversion to its rows.

func MergeRowGroups

func MergeRowGroups(rowGroups []RowGroup, options ...RowGroupOption) (RowGroup, error)

MergeRowGroups constructs a row group which is a merged view of rowGroups. If rowGroups are sorted and the passed options include sorting, the merged row group will also be sorted.

The function validates the input to ensure that the merge operation is possible, ensuring that the schemas match or can be converted to an optionally configured target schema passed as argument in the option list.

The sorting columns of each row group are also consulted to determine whether the output can be represented. If sorting columns are configured on the merge they must be a prefix of sorting columns of all row groups being merged.

func MultiRowGroup

func MultiRowGroup(rowGroups ...RowGroup) RowGroup

MultiRowGroup wraps multiple row groups to appear as if it was a single RowGroup. RowGroups must have the same schema or it will error.

type RowGroupConfig

type RowGroupConfig struct {
	ColumnBufferCapacity int
	Schema               *Schema
	Sorting              SortingConfig
}

The RowGroupConfig type carries configuration options for parquet row groups.

RowGroupConfig implements the RowGroupOption interface so it can be used directly as argument to the NewBuffer function when needed, for example:

buffer := parquet.NewBuffer(&parquet.RowGroupConfig{
	ColumnBufferCapacity: 10_000,
})

func DefaultRowGroupConfig

func DefaultRowGroupConfig() *RowGroupConfig

DefaultRowGroupConfig returns a new RowGroupConfig value initialized with the default row group configuration.

func NewRowGroupConfig

func NewRowGroupConfig(options ...RowGroupOption) (*RowGroupConfig, error)

NewRowGroupConfig constructs a new row group configuration applying the options passed as arguments.

The function returns an non-nil error if some of the options carried invalid configuration values.

func (*RowGroupConfig) Apply

func (c *RowGroupConfig) Apply(options ...RowGroupOption)

func (*RowGroupConfig) ConfigureRowGroup

func (c *RowGroupConfig) ConfigureRowGroup(config *RowGroupConfig)

func (*RowGroupConfig) Validate

func (c *RowGroupConfig) Validate() error

Validate returns a non-nil error if the configuration of c is invalid.

type RowGroupOption

type RowGroupOption interface {
	ConfigureRowGroup(*RowGroupConfig)
}

RowGroupOption is an interface implemented by types that carry configuration options for parquet row groups.

func ColumnBufferCapacity

func ColumnBufferCapacity(size int) RowGroupOption

ColumnBufferCapacity creates a configuration option which defines the size of row group column buffers.

Defaults to 16384.

func SortingRowGroupConfig

func SortingRowGroupConfig(options ...SortingOption) RowGroupOption

SortingRowGroupConfig is a row group option which applies configuration specific sorting row groups.

type RowGroupReader

type RowGroupReader interface {
	ReadRowGroup() (RowGroup, error)
}

RowGroupReader is an interface implemented by types that expose sequences of row groups to the application.

type RowGroupWriter

type RowGroupWriter interface {
	WriteRowGroup(RowGroup) (int64, error)
}

RowGroupWriter is an interface implemented by types that allow the program to write row groups.

type RowReadCloser

type RowReadCloser interface {
	RowReader
	io.Closer
}

RowReadCloser is an interface implemented by row readers which require closing when done.

type RowReadSeekCloser

type RowReadSeekCloser interface {
	RowReader
	RowSeeker
	io.Closer
}

RowReadSeekCloser is an interface implemented by row readers which support seeking to arbitrary row positions and required closing the reader when done.

func NewColumnChunkRowReader

func NewColumnChunkRowReader(columns []ColumnChunk) RowReadSeekCloser

NewColumnChunkRowReader creates a new ColumnChunkRowReader for the given column chunks.

type RowReadSeeker

type RowReadSeeker interface {
	RowReader
	RowSeeker
}

RowReadSeeker is an interface implemented by row readers which support seeking to arbitrary row positions.

type RowReader

type RowReader interface {
	// ReadRows reads rows from the reader, returning the number of rows read
	// into the buffer, and any error that occurred.
	//
	// When all rows have been read, the reader returns io.EOF to indicate the
	// end of the sequence. It is valid for the reader to return both a non-zero
	// number of rows and a non-nil error (including io.EOF).
	//
	// The buffer of rows passed as argument will be used to store values of
	// each row read from the reader. If the rows are not nil, the backing array
	// of the slices will be used as an optimization to avoid re-allocating new
	// arrays.
	//
	// The application is expected to handle the case where ReadRows returns
	// less rows than requested and no error, by looking at the first returned
	// value from ReadRows, which is the number of rows that were read.
	ReadRows([]Row) (int, error)
}

RowReader reads a sequence of parquet rows.

func DedupeRowReader

func DedupeRowReader(reader RowReader, compare func(Row, Row) int) RowReader

DedupeRowReader constructs a row reader which drops duplicated consecutive rows, according to the comparator function passed as argument.

If the underlying reader produces a sequence of rows sorted by the same comparison predicate, the output is guaranteed to produce unique rows only.

func FilterRowReader

func FilterRowReader(reader RowReader, predicate func(Row) bool) RowReader

FilterRowReader constructs a RowReader which exposes rows from reader for which the predicate has returned true.

func MergeRowReaders

func MergeRowReaders(rows []RowReader, compare func(Row, Row) int) RowReader

MergeRowReader constructs a RowReader which creates an ordered sequence of all the readers using the given compare function as the ordering predicate.

func ScanRowReader

func ScanRowReader(reader RowReader, predicate func(Row, int64) bool) RowReader

ScanRowReader constructs a RowReader which exposes rows from reader until the predicate returns false for one of the rows, or EOF is reached.

func TransformRowReader

func TransformRowReader(reader RowReader, transform func(dst, src Row) (Row, error)) RowReader

TransformRowReader constructs a RowReader which applies the given transform to each row rad from reader.

The transformation function appends the transformed src row to dst, returning dst and any error that occurred during the transformation. If dst is returned unchanged, the row is skipped.

type RowReaderFrom

type RowReaderFrom interface {
	ReadRowsFrom(RowReader) (int64, error)
}

RowReaderFrom reads parquet rows from reader.

type RowReaderFunc

type RowReaderFunc func([]Row) (int, error)

RowReaderFunc is a function type implementing the RowReader interface.

func (RowReaderFunc) ReadRows

func (f RowReaderFunc) ReadRows(rows []Row) (int, error)

type RowReaderWithSchema

type RowReaderWithSchema interface {
	RowReader
	Schema() *Schema
}

RowReaderWithSchema is an extension of the RowReader interface which advertises the schema of rows returned by ReadRow calls.

func ConvertRowReader

func ConvertRowReader(rows RowReader, conv Conversion) RowReaderWithSchema

ConvertRowReader constructs a wrapper of the given row reader which applies the given schema conversion to the rows.

type RowSeeker

type RowSeeker interface {
	// Positions the stream on the given row index.
	//
	// Some implementations of the interface may only allow seeking forward.
	//
	// The method returns io.ErrClosedPipe if the stream had already been closed.
	SeekToRow(int64) error
}

RowSeeker is an interface implemented by readers of parquet rows which can be positioned at a specific row index.

type RowWriter

type RowWriter interface {
	// Writes rows to the writer, returning the number of rows written and any
	// error that occurred.
	//
	// Because columnar operations operate on independent columns of values,
	// writes of rows may not be atomic operations, and could result in some
	// rows being partially written. The method returns the number of rows that
	// were successfully written, but if an error occurs, values of the row(s)
	// that failed to be written may have been partially committed to their
	// columns. For that reason, applications should consider a write error as
	// fatal and assume that they need to discard the state, they cannot retry
	// the write nor recover the underlying file.
	WriteRows([]Row) (int, error)
}

RowWriter writes parquet rows to an underlying medium.

func DedupeRowWriter

func DedupeRowWriter(writer RowWriter, compare func(Row, Row) int) RowWriter

DedupeRowWriter constructs a row writer which drops duplicated consecutive rows, according to the comparator function passed as argument.

If the writer is given a sequence of rows sorted by the same comparison predicate, the output is guaranteed to contain unique rows only.

func FilterRowWriter

func FilterRowWriter(writer RowWriter, predicate func(Row) bool) RowWriter

FilterRowWriter constructs a RowWriter which writes rows to writer for which the predicate has returned true.

func MultiRowWriter

func MultiRowWriter(writers ...RowWriter) RowWriter

MultiRowWriter constructs a RowWriter which dispatches writes to all the writers passed as arguments.

When writing rows, if any of the writers returns an error, the operation is aborted and the error returned. If one of the writers did not error, but did not write all the rows, the operation is aborted and io.ErrShortWrite is returned.

Rows are written sequentially to each writer in the order they are given to this function.

func TransformRowWriter

func TransformRowWriter(writer RowWriter, transform func(dst, src Row) (Row, error)) RowWriter

TransformRowWriter constructs a RowWriter which applies the given transform to each row writter to writer.

The transformation function appends the transformed src row to dst, returning dst and any error that occurred during the transformation. If dst is returned unchanged, the row is skipped.

type RowWriterFunc

type RowWriterFunc func([]Row) (int, error)

RowWriterFunc is a function type implementing the RowWriter interface.

func (RowWriterFunc) WriteRows

func (f RowWriterFunc) WriteRows(rows []Row) (int, error)

type RowWriterTo

type RowWriterTo interface {
	WriteRowsTo(RowWriter) (int64, error)
}

RowWriterTo writes parquet rows to a writer.

type RowWriterWithSchema

type RowWriterWithSchema interface {
	RowWriter
	Schema() *Schema
}

RowWriterWithSchema is an extension of the RowWriter interface which advertises the schema of rows expected to be passed to WriteRow calls.

type Rows

type Rows interface {
	RowReadSeekCloser
	Schema() *Schema
}

Rows is an interface implemented by row readers returned by calling the Rows method of RowGroup instances.

Applications should call Close when they are done using a Rows instance in order to release the underlying resources held by the row sequence.

After calling Close, all attempts to read more rows will return io.EOF.

func NewRowGroupRowReader

func NewRowGroupRowReader(rowGroup RowGroup) Rows

/ NewRowGroupRowReader constructs a new row reader for the given row group.

type Schema

type Schema struct {
	// contains filtered or unexported fields
}

Schema represents a parquet schema created from a Go value.

Schema implements the Node interface to represent the root node of a parquet schema.

Schema values are safe to use concurrently from multiple goroutines but must be passed by referenced after being created because their internal state contains synchronization primitives that are not safe to copy.

func NewSchema

func NewSchema(name string, root Node) *Schema

NewSchema constructs a new Schema object with the given name and root node.

The function panics if Node contains more leaf columns than supported by the package (see parquet.MaxColumnIndex).

func SchemaOf

func SchemaOf(model any, opts ...SchemaOption) *Schema

SchemaOf constructs a parquet schema from a Go value.

The function can construct parquet schemas from struct or pointer-to-struct values only. A panic is raised if a Go value of a different type is passed to this function.

When creating a parquet Schema from a Go value, the struct fields may contain a "parquet" tag to describe properties of the parquet node. The "parquet" tag follows the conventional format of Go struct tags: a comma-separated list of values describe the options, with the first one defining the name of the parquet column.

The following options are also supported in the "parquet" struct tag:

optional     | make the parquet column optional (any type)
snappy       | sets the parquet column compression codec to snappy (any type)
gzip         | sets the parquet column compression codec to gzip (any type)
brotli       | sets the parquet column compression codec to brotli (any type)
lz4          | sets the parquet column compression codec to lz4 (any type)
zstd         | sets the parquet column compression codec to zstd (any type)
uncompressed | explicitly sets no compression (any type)
plain        | enables the plain encoding (no-op default, any type)
dict         | enables dictionary encoding on the parquet column (any leaf type)
delta        | enables delta encoding: DeltaBinaryPacked for int32, int64, int, uint, uint32, uint64, time.Time; DeltaByteArray for string, []byte, [N]byte
list         | for slice types, use the parquet LIST logical type
enum         | for string types, use the parquet ENUM logical type
bytes        | for string types, use no parquet logical type
string       | for []byte types, use the parquet STRING logical type
json         | for string, []byte, and map types, use the parquet JSON logical type
uuid         | for string and [16]byte types, use the parquet UUID logical type
decimal      | for int32, int64, []byte, and [N]byte types, use the parquet DECIMAL logical type
date         | for int32, time.Time, and *time.Time types, use the DATE logical type
time         | for int32, int64, time.Duration, and *time.Duration types, use the TIME logical type
timestamp    | for int64, time.Time, and *time.Time types, use the TIMESTAMP logical type (default millisecond precision)
split        | for float32 and float64 types, use the BYTE_STREAM_SPLIT encoding
geometry     | for []byte types, use the GEOMETRY logical type; use geometry(crs) to set the CRS
geography    | for []byte types, use the GEOGRAPHY logical type; use geography(crs:algorithm) to set the CRS and edge algorithm
int(n)       | for integer types, use the parquet INT logical type with the given bit width (8, 16, 32, or 64)
uint(n)      | for integer types, use the parquet UINT logical type with the given bit width (8, 16, 32, or 64)
id(n)        | where n is int denoting a column field id. Example id(2) for a column with field id of 2

When "optional" is used on a bare slice (without the "list" tag), it applies to the repeated elements, not the slice itself. When combined with the "list" tag, "optional" applies to the list as a whole; use the parquet-element tag to make list elements optional (e.g. parquet-element:",optional").

The date logical type is an int32 value of the number of days since the unix epoch

The time and timestamp precision can be changed by defining which precision to use as an argument. Supported precisions are: nanosecond, millisecond and microsecond. Note that for the time tag, int32 only supports millisecond precision, while int64 supports microsecond and nanosecond precision. Example:

type Message struct {
  TimestampMicros int64 `parquet:"timestamp_micros,timestamp(microsecond)"
}

Both the time and timestamp tags accept an optional second parameter to set the `isAdjustedToUTC` annotation of the parquet logical type. Valid values are "utc" or "local". If not specified, the default value for this annotation will be "utc", which will set the `isAdjustedToUTC` annotation value to true. Example:

type Message struct {
  TimestampMicrosAdjusted int64 `parquet:"timestamp_micros_adjusted,timestamp(microsecond:utc)"
  TimestampMicrosNotAdjusted int64 `parquet:"timestamp_micros_not_adjusted,timestamp(microsecond:local)"
}

The decimal tag must be followed by two integer parameters, the first integer representing the scale and the second the precision; for example:

type Item struct {
	Cost int64 `parquet:"cost,decimal(0:3)"`
}

Invalid combination of struct tags and Go types, or repeating options will cause the function to panic.

As a special case, if the field tag is "-", the field is omitted from the schema and the data will not be written into the parquet file(s). Note that a field with name "-" can still be generated using the tag "-,".

The configuration of Parquet maps are done via two tags:

  • The `parquet-key` tag allows to configure the key of a map.
  • The `parquet-value` tag allows users to configure a map's values, for example to declare their native Parquet types.

When configuring a Parquet map, the `parquet` tag will configure the map itself.

For example, the following will set the int64 key of the map to be a timestamp:

type Actions struct {
  Action map[int64]string `parquet:"," parquet-key:",timestamp"`
}

To configure the element of a list, use the `parquet-element` tag. For example, the following will set the id of the element field to 2:

type Item struct {
  Attributes []string `parquet:",id(1),list" parquet-element:",id(2)"`
}

Note that the name of the element cannot be changed.

The schema name is the Go type name of the value.

func (*Schema) Columns

func (s *Schema) Columns() [][]string

Columns returns the list of column paths available in the schema.

The method always returns the same slice value across calls to ColumnPaths, applications should treat it as immutable.

func (*Schema) Comparator

func (s *Schema) Comparator(sortingColumns ...SortingColumn) func(Row, Row) int

Comparator constructs a comparator function which orders rows according to the list of sorting columns passed as arguments.

func (*Schema) Compression

func (s *Schema) Compression() compress.Codec

Compression returns the compression codec set on the root node of the parquet schema.

func (*Schema) ConfigureReader

func (s *Schema) ConfigureReader(config *ReaderConfig)

ConfigureReader satisfies the ReaderOption interface, allowing Schema instances to be passed to NewReader to pre-declare the schema of rows read from the reader.

func (*Schema) ConfigureRowGroup

func (s *Schema) ConfigureRowGroup(config *RowGroupConfig)

ConfigureRowGroup satisfies the RowGroupOption interface, allowing Schema instances to be passed to row group constructors to pre-declare the schema of the output parquet file.

func (*Schema) ConfigureWriter

func (s *Schema) ConfigureWriter(config *WriterConfig)

ConfigureWriter satisfies the WriterOption interface, allowing Schema instances to be passed to NewWriter to pre-declare the schema of the output parquet file.

func (*Schema) Deconstruct

func (s *Schema) Deconstruct(row Row, value any) Row

Deconstruct deconstructs a Go value and appends it to a row.

The method panics is the structure of the go value does not match the parquet schema.

func (*Schema) Encoding

func (s *Schema) Encoding() encoding.Encoding

Encoding returns the encoding set on the root node of the parquet schema.

func (*Schema) Fields

func (s *Schema) Fields() []Field

Fields returns the list of fields on the root node of the parquet schema.

func (*Schema) GoType

func (s *Schema) GoType() reflect.Type

GoType returns the Go type that best represents the schema.

func (*Schema) ID

func (s *Schema) ID() int

ID returns field id of the root node.

func (*Schema) Leaf

func (s *Schema) Leaf() bool

Leaf returns true if the root node of the parquet schema is a leaf column.

func (*Schema) Lookup

func (s *Schema) Lookup(path ...string) (LeafColumn, bool)

Lookup returns the leaf column at the given path.

The path is the sequence of column names identifying a leaf column (not including the root).

If the path was not found in the mapping, or if it did not represent a leaf column of the parquet schema, the boolean will be false.

Example
package main

import (
	"fmt"
	"strings"

	"github.com/AlexAslan/parquet-go"
)

func main() {
	schema := parquet.SchemaOf(struct {
		FirstName  string `parquet:"first_name"`
		LastName   string `parquet:"last_name"`
		Attributes []struct {
			Name  string `parquet:"name"`
			Value string `parquet:"value"`
		} `parquet:"attributes"`
	}{})

	for _, path := range schema.Columns() {
		leaf, _ := schema.Lookup(path...)
		fmt.Printf("%d => %q\n", leaf.ColumnIndex, strings.Join(path, "."))
	}

}
Output:
0 => "first_name"
1 => "last_name"
2 => "attributes.name"
3 => "attributes.value"

func (*Schema) Name

func (s *Schema) Name() string

Name returns the name of s.

func (*Schema) Optional

func (s *Schema) Optional() bool

Optional returns false since the root node of a parquet schema is always required.

func (*Schema) Reconstruct

func (s *Schema) Reconstruct(value any, row Row) error

Reconstruct reconstructs a Go value from a row.

The go value passed as first argument must be a non-nil pointer for the row to be decoded into.

The method panics if the structure of the go value and parquet row do not match.

func (*Schema) Repeated

func (s *Schema) Repeated() bool

Repeated returns false since the root node of a parquet schema is always required.

func (*Schema) Required

func (s *Schema) Required() bool

Required returns true since the root node of a parquet schema is always required.

func (*Schema) String

func (s *Schema) String() string

String returns a parquet schema representation of s.

func (*Schema) Type

func (s *Schema) Type() Type

Type returns the parquet type of s.

type SchemaConfig

type SchemaConfig struct {
	StructTags []StructTagOption
}

The SchemaConfig type carries configuration options for parquet schemas.

SchemaConfig implements the SchemaOption interface so it can be used directly as argument to the SchemaOf function when needed, for example:

schema := parquet.SchemaOf(obj, &parquet.SchemaConfig{
	...
})

func DefaultSchemaConfig

func DefaultSchemaConfig() *SchemaConfig

func (*SchemaConfig) ConfigureReader

func (c *SchemaConfig) ConfigureReader(config *ReaderConfig)

func (*SchemaConfig) ConfigureSchema

func (c *SchemaConfig) ConfigureSchema(config *SchemaConfig)

func (*SchemaConfig) ConfigureWriter

func (c *SchemaConfig) ConfigureWriter(config *WriterConfig)

type SchemaOption

type SchemaOption interface {
	ReaderOption
	WriterOption

	ConfigureSchema(*SchemaConfig)
}

SchemaOption is an interface implemented by types that carry configuration options for parquet schemas. SchemaOption also implements ReaderOption and WriterOption and may be used to configure the way NewGenericReader and NewGenericWriter derive schemas from the arguments.

func StructTag

func StructTag(tag reflect.StructTag, path ...string) SchemaOption

StructTag performs runtime replacement of struct tags when deriving a schema from a Go struct for the column at the given path. This option can be used anywhere a schema is derived from a Go struct including SchemaOf, NewGenericReader, and NewGenericWriter.

This option is additive, it may be used multiple times to affect multiple columns.

When renaming a column, configure the option by its original name.

type SortingColumn

type SortingColumn interface {
	// Returns the path of the column in the row group schema, omitting the name
	// of the root node.
	Path() []string

	// Returns true if the column will sort values in descending order.
	Descending() bool

	// Returns true if the column will put null values at the beginning.
	NullsFirst() bool
}

SortingColumn represents a column by which a row group is sorted.

func Ascending

func Ascending(path ...string) SortingColumn

Ascending constructs a SortingColumn value which dictates to sort the column at the path given as argument in ascending order.

func Descending

func Descending(path ...string) SortingColumn

Descending constructs a SortingColumn value which dictates to sort the column at the path given as argument in descending order.

func MergeSortingColumns

func MergeSortingColumns(sortingColumns ...[]SortingColumn) []SortingColumn

MergeSortingColumns returns the common prefix of all sorting columns passed as arguments. This function is used to determine the resulting sorting columns when merging multiple row groups that each have their own sorting columns.

The function returns the longest common prefix where all sorting columns match exactly (same path, same descending flag, same nulls first flag). If any row group has no sorting columns, or if there's no common prefix, an empty slice is returned.

Example:

columns1 := []SortingColumn{Ascending("A"), Ascending("B"), Descending("C")}
columns2 := []SortingColumn{Ascending("A"), Ascending("B"), Ascending("D")}
result := MergeSortingColumns(columns1, columns2)
// result will be []SortingColumn{Ascending("A"), Ascending("B")}

func NullsFirst

func NullsFirst(sortingColumn SortingColumn) SortingColumn

NullsFirst wraps the SortingColumn passed as argument so that it instructs the row group to place null values first in the column.

type SortingConfig

type SortingConfig struct {
	SortingBuffers     BufferPool
	SortingColumns     []SortingColumn
	DropDuplicatedRows bool
}

The SortingConfig type carries configuration options for parquet row groups.

SortingConfig implements the SortingOption interface so it can be used directly as argument to the NewSortingWriter function when needed, for example:

buffer := parquet.NewSortingWriter[Row](
	parquet.SortingWriterConfig(
		parquet.DropDuplicatedRows(true),
	),
})

func DefaultSortingConfig

func DefaultSortingConfig() *SortingConfig

DefaultSortingConfig returns a new SortingConfig value initialized with the default row group configuration.

func NewSortingConfig

func NewSortingConfig(options ...SortingOption) (*SortingConfig, error)

NewSortingConfig constructs a new sorting configuration applying the options passed as arguments.

The function returns an non-nil error if some of the options carried invalid configuration values.

func (*SortingConfig) Apply

func (c *SortingConfig) Apply(options ...SortingOption)

func (*SortingConfig) ConfigureSorting

func (c *SortingConfig) ConfigureSorting(config *SortingConfig)

func (*SortingConfig) Validate

func (c *SortingConfig) Validate() error

type SortingOption

type SortingOption interface {
	ConfigureSorting(*SortingConfig)
}

SortingOption is an interface implemented by types that carry configuration options for parquet sorting writers.

func DropDuplicatedRows

func DropDuplicatedRows(drop bool) SortingOption

DropDuplicatedRows configures whether a sorting writer will keep or remove duplicated rows.

Two rows are considered duplicates if the values of their all their sorting columns are equal.

Defaults to false

func SortingBuffers

func SortingBuffers(buffers BufferPool) SortingOption

SortingBuffers creates a configuration option which sets the pool of buffers used to hold intermediary state when sorting parquet rows.

Defaults to using in-memory buffers.

func SortingColumns

func SortingColumns(columns ...SortingColumn) SortingOption

SortingColumns creates a configuration option which defines the sorting order of columns in a row group.

The order of sorting columns passed as argument defines the ordering hierarchy; when elements are equal in the first column, the second column is used to order rows, etc...

type SortingWriter

type SortingWriter[T any] struct {
	// contains filtered or unexported fields
}

SortingWriter is a type similar to GenericWriter but it ensures that rows are sorted according to the sorting columns configured on the writer.

The writer accumulates rows in an in-memory buffer which is sorted when it reaches the target number of rows, then written to a temporary row group. When the writer is flushed or closed, the temporary row groups are merged into a row group in the output file, ensuring that rows remain sorted in the final row group.

Because row groups get encoded and compressed, they hold a lot less memory than if all rows were retained in memory. Sorting then merging rows chunks also tends to be a lot more efficient than sorting all rows in memory as it results in better CPU cache utilization since sorting multi-megabyte arrays causes a lot of cache misses since the data set cannot be held in CPU caches.

func NewSortingWriter

func NewSortingWriter[T any](output io.Writer, sortRowCount int64, options ...WriterOption) *SortingWriter[T]

NewSortingWriter constructs a new sorting writer which writes a parquet file where rows of each row group are ordered according to the sorting columns configured on the writer.

The sortRowCount argument defines the target number of rows that will be sorted in memory before being written to temporary row groups. The greater this value the more memory is needed to buffer rows in memory. Choosing a value that is too small limits the maximum number of rows that can exist in the output file since the writer cannot create more than 32K temporary row groups to hold the sorted row chunks.

func (*SortingWriter[T]) Close

func (w *SortingWriter[T]) Close() error

func (*SortingWriter[T]) File

func (w *SortingWriter[T]) File() FileView

File returns a FileView of the written parquet file. Only available after Close is called.

func (*SortingWriter[T]) Flush

func (w *SortingWriter[T]) Flush() error

Flush sorts any buffered rows and writes them to temporary storage. This can be called multiple times to manage memory usage. The actual merge and write to output happens on Close.

func (*SortingWriter[T]) Reset

func (w *SortingWriter[T]) Reset(output io.Writer)

func (*SortingWriter[T]) Schema

func (w *SortingWriter[T]) Schema() *Schema

func (*SortingWriter[T]) SetKeyValueMetadata

func (w *SortingWriter[T]) SetKeyValueMetadata(key, value string)

func (*SortingWriter[T]) Write

func (w *SortingWriter[T]) Write(rows []T) (int, error)

func (*SortingWriter[T]) WriteRows

func (w *SortingWriter[T]) WriteRows(rows []Row) (int, error)

type StructTagOption

type StructTagOption struct {
	ColumnPath []string
	StructTag  reflect.StructTag
}

StructTagOption performs runtime replacement of "parquet..." struct tags. This option can be used anywhere a schema is derived from a Go struct including SchemaOf, NewGenericReader, and NewGenericWriter.

func (*StructTagOption) ConfigureReader

func (f *StructTagOption) ConfigureReader(config *ReaderConfig)

func (*StructTagOption) ConfigureSchema

func (f *StructTagOption) ConfigureSchema(config *SchemaConfig)

func (*StructTagOption) ConfigureWriter

func (f *StructTagOption) ConfigureWriter(config *WriterConfig)

type TimeUnit

type TimeUnit interface {
	// Returns the precision of the time unit as a time.Duration value.
	Duration() time.Duration
	// Converts the TimeUnit value to its representation in the parquet thrift
	// format.
	TimeUnit() format.TimeUnit
}

TimeUnit represents units of time in the parquet type system.

var (
	Millisecond TimeUnit = &millisecond{}
	Microsecond TimeUnit = &microsecond{}
	Nanosecond  TimeUnit = &nanosecond{}
)

type Type

type Type interface {
	// Returns a human-readable representation of the parquet type.
	String() string

	// Returns the Kind value representing the underlying physical type.
	//
	// The method panics if it is called on a group type.
	Kind() Kind

	// For integer and floating point physical types, the method returns the
	// size of values in bits.
	//
	// For fixed-length byte arrays, the method returns the size of elements
	// in bytes.
	//
	// For other types, the value is zero.
	Length() int

	// Returns an estimation of the number of bytes required to hold the given
	// number of values of this type in memory.
	//
	// The method returns zero for group types.
	EstimateSize(numValues int) int

	// Returns an estimation of the number of values of this type that can be
	// held in the given byte size.
	//
	// The method returns zero for group types.
	EstimateNumValues(size int) int

	// Compares two values and returns a negative integer if a < b, positive if
	// a > b, or zero if a == b.
	//
	// The values' Kind must match the type, otherwise the result is undefined.
	//
	// The method panics if it is called on a group type.
	Compare(a, b Value) int

	// ColumnOrder returns the type's column order. For group types, this method
	// returns nil.
	//
	// The order describes the comparison logic implemented by the Less method.
	//
	// As an optimization, the method may return the same pointer across
	// multiple calls. Applications must treat the returned value as immutable,
	// mutating the value will result in undefined behavior.
	ColumnOrder() *format.ColumnOrder

	// Returns the physical type as a *format.Type value. For group types, this
	// method returns nil.
	//
	// As an optimization, the method may return the same pointer across
	// multiple calls. Applications must treat the returned value as immutable,
	// mutating the value will result in undefined behavior.
	PhysicalType() *format.Type

	// Returns the logical type as a *format.LogicalType value. When the logical
	// type is unknown, the method returns nil.
	//
	// As an optimization, the method may return the same pointer across
	// multiple calls. Applications must treat the returned value as immutable,
	// mutating the value will result in undefined behavior.
	LogicalType() *format.LogicalType

	// Returns the logical type's equivalent converted type. When there are
	// no equivalent converted type, the method returns nil.
	//
	// As an optimization, the method may return the same pointer across
	// multiple calls. Applications must treat the returned value as immutable,
	// mutating the value will result in undefined behavior.
	ConvertedType() *deprecated.ConvertedType

	// Creates a column indexer for values of this type.
	//
	// The size limit is a hint to the column indexer that it is allowed to
	// truncate the page boundaries to the given size. Only BYTE_ARRAY and
	// FIXED_LEN_BYTE_ARRAY types currently take this value into account.
	//
	// A value of zero or less means no limits.
	//
	// The method panics if it is called on a group type.
	NewColumnIndexer(sizeLimit int) ColumnIndexer

	// Creates a row group buffer column for values of this type.
	//
	// Column buffers are created using the index of the column they are
	// accumulating values in memory for (relative to the parent schema),
	// and the size of their memory buffer.
	//
	// The application may give an estimate of the number of values it expects
	// to write to the buffer as second argument. This estimate helps set the
	// initialize buffer capacity but is not a hard limit, the underlying memory
	// buffer will grown as needed to allow more values to be written. Programs
	// may use the Size method of the column buffer (or the parent row group,
	// when relevant) to determine how many bytes are being used, and perform a
	// flush of the buffers to a storage layer.
	//
	// The method panics if it is called on a group type.
	NewColumnBuffer(columnIndex, numValues int) ColumnBuffer

	// Creates a dictionary holding values of this type.
	//
	// The dictionary retains the data buffer, it does not make a copy of it.
	// If the application needs to share ownership of the memory buffer, it must
	// ensure that it will not be modified while the page is in use, or it must
	// make a copy of it prior to creating the dictionary.
	//
	// The method panics if the data type does not correspond to the parquet
	// type it is called on.
	NewDictionary(columnIndex, numValues int, data encoding.Values) Dictionary

	// Creates a page belonging to a column at the given index, backed by the
	// data buffer.
	//
	// The page retains the data buffer, it does not make a copy of it. If the
	// application needs to share ownership of the memory buffer, it must ensure
	// that it will not be modified while the page is in use, or it must make a
	// copy of it prior to creating the page.
	//
	// The method panics if the data type does not correspond to the parquet
	// type it is called on.
	NewPage(columnIndex, numValues int, data encoding.Values) Page

	// Creates an encoding.Values instance backed by the given buffers.
	//
	// The offsets is only used by BYTE_ARRAY types, where it represents the
	// positions of each variable length value in the values buffer.
	//
	// The following expression creates an empty instance for any type:
	//
	//		values := typ.NewValues(nil, nil)
	//
	// The method panics if it is called on group types.
	NewValues(values []byte, offsets []uint32) encoding.Values

	// Assuming the src buffer contains PLAIN encoded values of the type it is
	// called on, applies the given encoding and produces the output to the dst
	// buffer passed as first argument by dispatching the call to one of the
	// encoding methods.
	Encode(dst []byte, src encoding.Values, enc encoding.Encoding) ([]byte, error)

	// Assuming the src buffer contains values encoding in the given encoding,
	// decodes the input and produces the encoded values into the dst output
	// buffer passed as first argument by dispatching the call to one of the
	// encoding methods.
	Decode(dst encoding.Values, src []byte, enc encoding.Encoding) (encoding.Values, error)

	// Returns an estimation of the output size after decoding the values passed
	// as first argument with the given encoding.
	//
	// For most types, this is similar to calling EstimateSize with the known
	// number of encoded values. For variable size types, using this method may
	// provide a more precise result since it can inspect the input buffer.
	EstimateDecodeSize(numValues int, src []byte, enc encoding.Encoding) int

	// Assigns a Parquet value to a Go value. Returns an error if assignment is
	// not possible. The source Value must be an expected logical type for the
	// receiver. This can be accomplished using ConvertValue.
	AssignValue(dst reflect.Value, src Value) error

	// Convert a Parquet Value of the given Type into a Parquet Value that is
	// compatible with the receiver. The returned Value is suitable to be passed
	// to AssignValue.
	ConvertValue(val Value, typ Type) (Value, error)
}

The Type interface represents logical types of the parquet type system.

Types are immutable and therefore safe to access from multiple goroutines.

var (
	BooleanType   Type = booleanType{}
	Int32Type     Type = int32Type{}
	Int64Type     Type = int64Type{}
	Int96Type     Type = int96Type{}
	FloatType     Type = floatType{}
	DoubleType    Type = doubleType{}
	ByteArrayType Type = byteArrayType{}
	NullType      Type = &nullType{}
)

func FixedLenByteArrayType

func FixedLenByteArrayType(length int) Type

FixedLenByteArrayType constructs a type for fixed-length values of the given size (in bytes).

type UUIDReader

type UUIDReader interface {
	// Read UUID values into the buffer passed as argument, returning the
	// number of values read.
	//
	// The method returns io.EOF when all values have been read.
	ReadUUIDs(values []uuid.UUID) (int, error)
}

UUIDReader is an interface implemented by ValueReader instances which expose the content of a column of UUID values stored as 16-byte FIXED_LEN_BYTE_ARRAY with the UUID logical type.

type UUIDWriter

type UUIDWriter interface {
	// Write UUID values.
	//
	// The method returns the number of values written, and any error that
	// occurred while writing the values.
	WriteUUIDs(values []uuid.UUID) (int, error)
}

UUIDWriter is an interface implemented by ValueWriter instances which support writing columns of UUID values stored as 16-byte FIXED_LEN_BYTE_ARRAY with the UUID logical type.

type Value

type Value struct {
	// contains filtered or unexported fields
}

The Value type is similar to the reflect.Value abstraction of Go values, but for parquet values. Value instances wrap underlying Go values mapped to one of the parquet physical types.

Value instances are small, immutable objects, and usually passed by value between function calls.

The zero-value of Value represents the null parquet value.

func BooleanValue

func BooleanValue(value bool) Value

BooleanValue constructs a BOOLEAN parquet value from the bool passed as argument.

func ByteArrayValue

func ByteArrayValue(value []byte) Value

ByteArrayValue constructs a BYTE_ARRAY parquet value from the byte slice passed as argument.

func DoubleValue

func DoubleValue(value float64) Value

DoubleValue constructs a DOUBLE parquet value from the float64 passed as argument.

func FixedLenByteArrayValue

func FixedLenByteArrayValue(value []byte) Value

FixedLenByteArrayValue constructs a BYTE_ARRAY parquet value from the byte slice passed as argument.

func FloatValue

func FloatValue(value float32) Value

FloatValue constructs a FLOAT parquet value from the float32 passed as argument.

func Int32Value

func Int32Value(value int32) Value

Int32Value constructs a INT32 parquet value from the int32 passed as argument.

func Int64Value

func Int64Value(value int64) Value

Int64Value constructs a INT64 parquet value from the int64 passed as argument.

func Int96Value

func Int96Value(value deprecated.Int96) Value

Int96Value constructs a INT96 parquet value from the deprecated.Int96 passed as argument.

func NullValue

func NullValue() Value

NulLValue constructs a null value, which is the zero-value of the Value type.

func ValueOf

func ValueOf(v any) Value

ValueOf constructs a parquet value from a Go value v.

The physical type of the value is assumed from the Go type of v using the following conversion table:

Go type | Parquet physical type
------- | ---------------------
nil     | NULL
bool    | BOOLEAN
int8    | INT32
int16   | INT32
int32   | INT32
int64   | INT64
int     | INT64
uint8   | INT32
uint16  | INT32
uint32  | INT32
uint64  | INT64
uintptr | INT64
float32 | FLOAT
float64 | DOUBLE
string  | BYTE_ARRAY
[]byte  | BYTE_ARRAY
[*]byte | FIXED_LEN_BYTE_ARRAY

When converting a []byte or [*]byte value, the underlying byte array is not copied; instead, the returned parquet value holds a reference to it.

The repetition and definition levels of the returned value are both zero.

The function panics if the Go value cannot be represented in parquet.

func ZeroValue

func ZeroValue(kind Kind) Value

ZeroValue constructs a zero value of the given kind.

func (Value) AppendBytes

func (v Value) AppendBytes(b []byte) []byte

AppendBytes appends the binary representation of v to b.

If v is the null value, b is returned unchanged.

func (Value) Boolean

func (v Value) Boolean() bool

Boolean returns v as a bool, assuming the underlying type is BOOLEAN.

func (Value) Byte

func (v Value) Byte() byte

Byte returns v as a byte, which may truncate the underlying byte.

func (Value) ByteArray

func (v Value) ByteArray() []byte

ByteArray returns v as a []byte, assuming the underlying type is either BYTE_ARRAY or FIXED_LEN_BYTE_ARRAY.

The application must treat the returned byte slice as a read-only value, mutating the content will result in undefined behaviors.

func (Value) Bytes

func (v Value) Bytes() []byte

Bytes returns the binary representation of v.

If v is the null value, an nil byte slice is returned.

func (Value) Clone

func (v Value) Clone() Value

Clone returns a copy of v which does not share any pointers with it.

func (Value) Column

func (v Value) Column() int

Column returns the column index within the row that v was created from.

Returns -1 if the value does not carry a column index.

func (Value) DefinitionLevel

func (v Value) DefinitionLevel() int

DefinitionLevel returns the definition level of v.

func (Value) Double

func (v Value) Double() float64

Double returns v as a float64, assuming the underlying type is DOUBLE.

func (Value) Float

func (v Value) Float() float32

Float returns v as a float32, assuming the underlying type is FLOAT.

func (Value) Format

func (v Value) Format(w fmt.State, r rune)

Format outputs a human-readable representation of v to w, using r as the formatting verb to describe how the value should be printed.

The following formatting options are supported:

%c	prints the column index
%+c	prints the column index, prefixed with "C:"
%d	prints the definition level
%+d	prints the definition level, prefixed with "D:"
%r	prints the repetition level
%+r	prints the repetition level, prefixed with "R:"
%q	prints the quoted representation of v
%+q	prints the quoted representation of v, prefixed with "V:"
%s	prints the string representation of v
%+s	prints the string representation of v, prefixed with "V:"
%v	same as %s
%+v	prints a verbose representation of v
%#v	prints a Go value representation of v

Format satisfies the fmt.Formatter interface.

func (Value) GoString

func (v Value) GoString() string

GoString returns a Go value string representation of v.

func (Value) Int32

func (v Value) Int32() int32

Int32 returns v as a int32, assuming the underlying type is INT32.

func (Value) Int64

func (v Value) Int64() int64

Int64 returns v as a int64, assuming the underlying type is INT64.

func (Value) Int96

func (v Value) Int96() deprecated.Int96

Int96 returns v as a int96, assuming the underlying type is INT96.

func (Value) IsNull

func (v Value) IsNull() bool

IsNull returns true if v is the null value.

func (Value) Kind

func (v Value) Kind() Kind

Kind returns the kind of v, which represents its parquet physical type.

func (Value) Level

func (v Value) Level(repetitionLevel, definitionLevel, columnIndex int) Value

Level returns v with the repetition level, definition level, and column index set to the values passed as arguments.

The method panics if any argument is negative or exceeds its maximum (MaxRepetitionLevel, MaxDefinitionLevel, MaxColumnIndex respectively).

To mutate a single level field in place, see SetRepetitionLevel, SetDefinitionLevel, and SetColumnIndex, which the compiler can inline.

func (Value) RepetitionLevel

func (v Value) RepetitionLevel() int

RepetitionLevel returns the repetition level of v.

func (*Value) SetColumnIndex

func (v *Value) SetColumnIndex(columnIndex int)

SetColumnIndex sets the column index of v to the given value.

See SetRepetitionLevel for the rationale behind the single-field setters.

The method panics if columnIndex is negative or greater than MaxColumnIndex.

func (*Value) SetDefinitionLevel

func (v *Value) SetDefinitionLevel(level int)

SetDefinitionLevel sets the definition level of v to the given value.

See SetRepetitionLevel for the rationale behind the single-field setters.

The method panics if level is negative or greater than MaxDefinitionLevel.

func (*Value) SetRepetitionLevel

func (v *Value) SetRepetitionLevel(level int)

SetRepetitionLevel sets the repetition level of v to the given value.

Unlike Level, which sets all three level fields and returns a new Value, this method mutates a single field in place through a pointer receiver. Avoiding the value copy keeps it small enough for the compiler to inline, which lets the bounds check be eliminated entirely when the argument is provably in range (for example when forwarding an existing repetition level). Mutating a slice element in place (values[i].SetRepetitionLevel(...)) makes it a good fit for hot loops.

The method panics if level is negative or greater than MaxRepetitionLevel.

func (Value) String

func (v Value) String() string

String returns a string representation of v.

func (Value) Uint32

func (v Value) Uint32() uint32

Uint32 returns v as a uint32, assuming the underlying type is INT32.

func (Value) Uint64

func (v Value) Uint64() uint64

Uint64 returns v as a uint64, assuming the underlying type is INT64.

type ValueReader

type ValueReader interface {
	// Read values into the buffer passed as argument and return the number of
	// values read. When all values have been read, the error will be io.EOF.
	ReadValues([]Value) (int, error)
}

ValueReader is an interface implemented by types that support reading batches of values.

type ValueReaderAt

type ValueReaderAt interface {
	ReadValuesAt([]Value, int64) (int, error)
}

ValueReaderAt is an interface implemented by types that support reading values at offsets specified by the application.

type ValueReaderFrom

type ValueReaderFrom interface {
	ReadValuesFrom(ValueReader) (int64, error)
}

ValueReaderFrom is an interface implemented by value writers to read values from a reader.

type ValueReaderFunc

type ValueReaderFunc func([]Value) (int, error)

ValueReaderFunc is a function type implementing the ValueReader interface.

func (ValueReaderFunc) ReadValues

func (f ValueReaderFunc) ReadValues(values []Value) (int, error)

type ValueWriter

type ValueWriter interface {
	// Write values from the buffer passed as argument and returns the number
	// of values written.
	WriteValues([]Value) (int, error)
}

ValueWriter is an interface implemented by types that support reading batches of values.

type ValueWriterFunc

type ValueWriterFunc func([]Value) (int, error)

ValueWriterFunc is a function type implementing the ValueWriter interface.

func (ValueWriterFunc) WriteValues

func (f ValueWriterFunc) WriteValues(values []Value) (int, error)

type ValueWriterTo

type ValueWriterTo interface {
	WriteValuesTo(ValueWriter) (int64, error)
}

ValueWriterTo is an interface implemented by value readers to write values to a writer.

type VariantColumnWriter

type VariantColumnWriter struct {
	// contains filtered or unexported fields
}

VariantColumnWriter writes one VARIANT column of a parquet file in streaming, columnar form. It implements variant.ValueBuilder: each row is bracketed by BeginRow and EndRow, and the row's value is described by the event methods in between, following the event-sequence rules documented on variant.ValueBuilder. WriteValue and WriteNullRow are row-level conveniences. String and []byte event arguments, and Field names, are copied; callers may reuse their backing memory after the call returns.

The writer appends to the column buffers of the ColumnWriters backing the variant column's leaf columns and flushes full pages at row boundaries, like ColumnWriter.WriteRowValues. To keep row boundaries cheap the page-size check is polled every few rows while the buffers are under half the configured page size (and every row past that), so pages may moderately overshoot the configured size. Other columns of the file must be written separately (e.g. through Writer.ColumnWriters) with the same number of rows, and the parent writer must be closed (or its row group committed) as usual to flush buffered sub-page data to the file.

Errors are sticky: the first misuse or write failure is retained, subsequent events are ignored, and the error is reported by Err, EndRow, and row-level methods. After an error the underlying columns may hold a partial row, so the parent writer's current row group is inconsistent and must be abandoned as a whole, including the other columns' writers: for a ConcurrentRowGroupWriter target, do not Commit; for a Writer or GenericWriter target, discard the output.

A VariantColumnWriter is not safe for concurrent use, and the parent writer must not be flushed, reset, or closed between BeginRow and EndRow (rows cannot span pages or row groups).

func NewVariantColumnWriter

func NewVariantColumnWriter(w WriterTarget, path ...string) (*VariantColumnWriter, error)

NewVariantColumnWriter returns a streaming columnar writer for the VARIANT column at the given path in w's schema (e.g. "event" for a top-level column, "attrs", "v" for one nested in a group). The column may be unshredded (metadata, value) or shredded (metadata, value, typed_value) in any shape permitted by the Variant Shredding specification.

func (*VariantColumnWriter) BeginArray

func (w *VariantColumnWriter) BeginArray()

BeginArray starts an array value at the current position. If the position is shredded as a list, every value until the matching EndArray shreds as one element; any other shredded type is a conflict and the whole array streams into the position's value column.

func (*VariantColumnWriter) BeginObject

func (w *VariantColumnWriter) BeginObject()

BeginObject starts an object value at the current position. If the position is shredded as an object, its fields shred individually; any other shredded type is a conflict and the whole object streams into the position's value column.

func (*VariantColumnWriter) BeginRow

func (w *VariantColumnWriter) BeginRow() error

BeginRow starts a new row. Exactly one value — a single scalar event, or a single balanced object or array — must be written before the matching EndRow.

func (*VariantColumnWriter) Binary

func (w *VariantColumnWriter) Binary(v []byte)

func (*VariantColumnWriter) Bool

func (w *VariantColumnWriter) Bool(v bool)

func (*VariantColumnWriter) Date

func (w *VariantColumnWriter) Date(days int32)

func (*VariantColumnWriter) Decimal4

func (w *VariantColumnWriter) Decimal4(unscaled int32, scale byte)

func (*VariantColumnWriter) Decimal8

func (w *VariantColumnWriter) Decimal8(unscaled int64, scale byte)

func (*VariantColumnWriter) Decimal16

func (w *VariantColumnWriter) Decimal16(unscaled [16]byte, scale byte)

func (*VariantColumnWriter) Double

func (w *VariantColumnWriter) Double(v float64)

func (*VariantColumnWriter) EndArray

func (w *VariantColumnWriter) EndArray()

EndArray closes the innermost open array.

func (*VariantColumnWriter) EndObject

func (w *VariantColumnWriter) EndObject()

EndObject closes the innermost open object, recording shredded fields that received no value as missing and committing the partial residual, if any, to the object's value column.

func (*VariantColumnWriter) EndRow

func (w *VariantColumnWriter) EndRow() error

EndRow completes the current row: it writes the row's metadata dictionary and periodically flushes columns whose buffered page reached the writer's page size.

func (*VariantColumnWriter) Err

func (w *VariantColumnWriter) Err() error

Err returns the first error encountered, or nil.

func (*VariantColumnWriter) Field

func (w *VariantColumnWriter) Field(name string)

Field names the next value written as a field of the innermost open object. Values of shredded fields go to the field's columns; others stream into the object's partial residual.

func (*VariantColumnWriter) FieldByRef

func (w *VariantColumnWriter) FieldByRef(ref *VariantFieldRef)

FieldByRef is Field with a pre-resolved name; see NewVariantFieldRef.

func (*VariantColumnWriter) Float

func (w *VariantColumnWriter) Float(v float32)

func (*VariantColumnWriter) Int8

func (w *VariantColumnWriter) Int8(v int8)

func (*VariantColumnWriter) Int16

func (w *VariantColumnWriter) Int16(v int16)

func (*VariantColumnWriter) Int32

func (w *VariantColumnWriter) Int32(v int32)

func (*VariantColumnWriter) Int64

func (w *VariantColumnWriter) Int64(v int64)

func (*VariantColumnWriter) Null

func (w *VariantColumnWriter) Null()

Null writes the variant null value (distinct from a missing field and from a SQL null row; see WriteNullRow). Variant null never matches a shredded type: it encodes into the position's value column, or fails the row when the position is shredded without one.

func (*VariantColumnWriter) String

func (w *VariantColumnWriter) String(v string)

func (*VariantColumnWriter) Time

func (w *VariantColumnWriter) Time(micros int64)

func (*VariantColumnWriter) Timestamp

func (w *VariantColumnWriter) Timestamp(us int64)

func (*VariantColumnWriter) TimestampNTZ

func (w *VariantColumnWriter) TimestampNTZ(us int64)

func (*VariantColumnWriter) TimestampNTZNanos

func (w *VariantColumnWriter) TimestampNTZNanos(ns int64)

func (*VariantColumnWriter) TimestampNanos

func (w *VariantColumnWriter) TimestampNanos(ns int64)

func (*VariantColumnWriter) UUID

func (w *VariantColumnWriter) UUID(v uuid.UUID)

func (*VariantColumnWriter) WriteNullRow

func (w *VariantColumnWriter) WriteNullRow() error

WriteNullRow writes a row where the variant column itself is null (a SQL null, distinct from the variant null value). The variant group must be optional. Any optional ancestors of the variant group are written as present; the writer cannot express a row where an ancestor group is null.

func (*VariantColumnWriter) WriteValue

func (w *VariantColumnWriter) WriteValue(v variant.Value) error

WriteValue writes one row holding the given variant value: it brackets v.Write(w) with BeginRow and EndRow and reports any error either raised.

type VariantCursor

type VariantCursor struct {
	// contains filtered or unexported fields
}

VariantCursor is a position in the variant navigation tree: the root value, an object field, or a list element. See VariantReader.

func (*VariantCursor) Booleans

func (c *VariantCursor) Booleans() []bool

Booleans returns the dense typed values of a BOOLEAN leaf cursor.

func (*VariantCursor) ByteArrays

func (c *VariantCursor) ByteArrays() (slab []byte, offsets []uint32)

ByteArrays returns the dense typed values of a BYTE_ARRAY-backed leaf cursor (string, binary, and byte-array decimal variant types) as a shared slab plus len+1 offsets: value i is slab[offsets[i]:offsets[i+1]].

func (*VariantCursor) Doubles

func (c *VariantCursor) Doubles() []float64

Doubles returns the dense typed values of a DOUBLE leaf cursor.

func (*VariantCursor) Elements

func (c *VariantCursor) Elements() *VariantCursor

Elements returns the cursor for the elements of lists at this position. The element entry space concatenates, in row order, the elements of typed (shredded) lists and the elements of residual values that turn out to be arrays, so an engine can process mixed shredded/residual lists in one loop. ListOffsets on this cursor maps entries to element ranges.

func (*VariantCursor) Field

func (c *VariantCursor) Field(name string) *VariantCursor

Field returns the cursor for the object field with the given name below this cursor. Navigation is total: if the field is not part of the shredded schema at this position, the returned cursor is backed by per-row navigation of residual values.

func (*VariantCursor) Fields

func (c *VariantCursor) Fields() []string

Fields returns the names of the shredded fields at a cursor of kind VariantCursorObject, in schema order, and nil for other kinds. Residual fields of partially shredded objects are not listed; they are reached through Residual (or by name through Field).

func (*VariantCursor) FixedLenByteArrays

func (c *VariantCursor) FixedLenByteArrays() (slab []byte, size int)

FixedLenByteArrays returns the dense typed values of a FIXED_LEN_BYTE_ARRAY-backed leaf cursor (uuid and fixed decimal variant types) as a shared slab of size-byte values.

func (*VariantCursor) Floats

func (c *VariantCursor) Floats() []float32

Floats returns the dense typed values of a FLOAT leaf cursor.

func (*VariantCursor) Int32s

func (c *VariantCursor) Int32s() []int32

Int32s returns the dense typed values of an INT32-backed leaf cursor (int8, int16, int32, date, and decimal4 variant types).

func (*VariantCursor) Int64s

func (c *VariantCursor) Int64s() []int64

Int64s returns the dense typed values of an INT64-backed leaf cursor (int64, timestamps, time, and decimal8 variant types).

func (*VariantCursor) Kind

func (c *VariantCursor) Kind() VariantCursorKind

Kind returns the static kind of the cursor position, derived from the shredded schema.

func (*VariantCursor) LeafType

func (c *VariantCursor) LeafType() Type

LeafType returns the parquet type of the typed_value column for cursors of kind VariantCursorLeaf, and nil otherwise. The type's logical type identifies the variant type per the shredded types table (e.g. a STRING-annotated BYTE_ARRAY holds variant strings).

func (*VariantCursor) ListOffsets

func (c *VariantCursor) ListOffsets() []int32

ListOffsets returns len(entries)+1 offsets into the entry space of the Elements cursor: the elements of entry e span Elements() entries [offsets[e], offsets[e+1]). It is computed only when Elements was called before the window was read.

func (*VariantCursor) Locs

func (c *VariantCursor) Locs() []variant.Loc

Locs returns the location tag of every entry of the current window: one per row for row-level cursors, one per element for element cursors.

func (*VariantCursor) Residual

func (c *VariantCursor) Residual(i int) (variant.Value, bool, error)

Residual returns the residual variant value of entry i. It is valid for entries tagged variant.LocResidual (the value itself), variant.LocNull (the variant null value), and variant.LocTypedObject when the row is a partially shredded object (the leftover fields not covered by the shredded schema, as a variant object). The boolean result reports whether a residual value exists at the entry.

func (*VariantCursor) ResidualCount

func (c *VariantCursor) ResidualCount() int

ResidualCount returns the number of entries tagged variant.LocResidual in the current window. A zero count lets an engine skip per-entry residual decoding, though Residual may still return values for entries tagged variant.LocNull or variant.LocTypedObject (partial-object leftovers).

func (*VariantCursor) Rows

func (c *VariantCursor) Rows() []int32

Rows maps each entry of the current window to its row index within the window. For row-level cursors the mapping is the identity and Rows returns nil; for element cursors (and cursors below them) it returns one row index per entry, and several entries may map to one row.

func (*VariantCursor) TypedRows

func (c *VariantCursor) TypedRows() []int32

TypedRows maps each value of the typed vectors to its entry index in the current window: TypedRows()[i] is the entry holding the i-th typed value.

type VariantCursorKind

type VariantCursorKind uint8

VariantCursorKind is the static (schema-derived) kind of a cursor position. It tells the engine which access pattern the position supports; the per-row runtime disposition is in Locs.

const (
	// VariantCursorUnshredded is a position with no typed_value column:
	// either the schema stores only variant binary here, or the cursor
	// navigates below the shredded schema entirely. All present values are
	// variant.LocNull or variant.LocResidual.
	VariantCursorUnshredded VariantCursorKind = iota

	// VariantCursorLeaf is a position shredded as a primitive type; typed
	// values are read from the typed vector accessors.
	VariantCursorLeaf

	// VariantCursorObject is a position shredded as an object; descend into
	// shredded fields with Field.
	VariantCursorObject

	// VariantCursorList is a position shredded as a list; descend with
	// Elements and use ListOffsets.
	VariantCursorList
)

func (VariantCursorKind) String

func (k VariantCursorKind) String() string

type VariantFieldRef

type VariantFieldRef struct {
	// contains filtered or unexported fields
}

VariantFieldRef is a pre-resolved object field name for FieldByRef: a ref caches its resolution (shredded position and metadata intern) against the last writer and object it was used with, so hot ingest loops writing the same fields on every row skip the per-event name hashing of Field. The cache re-resolves whenever the writer or object position changes, but FieldByRef mutates it, so a ref must not be used by multiple goroutines concurrently — for concurrency and peak performance use one ref per writer per field.

func NewVariantFieldRef

func NewVariantFieldRef(name string) *VariantFieldRef

NewVariantFieldRef pre-resolves an object field name for use with VariantColumnWriter.FieldByRef. The name is copied; the caller may reuse its backing memory. The ref is not bound to any writer: resolution is cached lazily by FieldByRef.

type VariantReader

type VariantReader struct {
	// contains filtered or unexported fields
}

VariantReader reads one VARIANT column of a row group in columnar form.

All cursors created from a reader share one row window, advanced with Next. Vectors and residual views returned by cursors are valid until the next call to Next or SeekToRow.

Only the columns reachable from cursors that have been created are read: creating cursors is how an engine declares projection. Cursors may be created at any time; cursors created after a call to Next take effect on the following call.

func NewVariantReader

func NewVariantReader(rowGroup RowGroup, path ...string) (*VariantReader, error)

NewVariantReader creates a reader for the variant column at the given path in the row group. The path names the variant group node in the schema (e.g. "event" for a top-level column, "attrs", "v" for a nested one). The column may be unshredded (metadata, value) or shredded (metadata, value, typed_value) in any shape permitted by the Variant Shredding specification.

func (*VariantReader) Close

func (r *VariantReader) Close() error

Close releases the page readers of all columns opened by the reader.

func (*VariantReader) Next

func (r *VariantReader) Next(n int) (int, error)

Next advances the shared row window by up to n rows and recomputes the window state of every cursor. It returns the number of rows in the new window, or (0, io.EOF) when the row group is exhausted.

func (*VariantReader) NumRows

func (r *VariantReader) NumRows() int64

NumRows returns the number of rows of the row group.

func (*VariantReader) Path

func (r *VariantReader) Path(path ...string) *VariantCursor

Path returns the cursor for the object field path below the root, equivalent to chaining Field calls.

func (*VariantReader) Root

func (r *VariantReader) Root() *VariantCursor

Root returns the cursor for the variant value itself.

func (*VariantReader) SeekToRow

func (r *VariantReader) SeekToRow(rowIndex int64) error

SeekToRow positions the reader such that the next call to Next starts at the given row of the row group.

type Writer deprecated

type Writer struct {
	// contains filtered or unexported fields
}

Deprecated: A Writer uses a parquet schema and sequence of Go values to produce a parquet file to an io.Writer.

Use NewGenericWriter instead. To maintain dynamic behavior (schema unknown at compile time), use "any" as the type parameter:

This gives you the same behavior as the old Writer
w := parquet.NewGenericWriter[any](output, schema)

This example showcases a typical use of parquet writers:

writer := parquet.NewWriter(output)

for _, row := range rows {
	if err := writer.Write(row); err != nil {
		...
	}
}

if err := writer.Close(); err != nil {
	...
}

The Writer type optimizes for minimal memory usage, each page is written as soon as it has been filled so only a single page per column needs to be held in memory and as a result, there are no opportunities to sort rows within an entire row group. Programs that need to produce parquet files with sorted row groups should use the Buffer type to buffer and sort the rows prior to writing them to a Writer.

For programs building with Go 1.18 or later, the GenericWriter[T] type supersedes this one.

func NewWriter

func NewWriter(output io.Writer, options ...WriterOption) *Writer

NewWriter constructs a parquet writer writing a file to the given io.Writer.

The function panics if the writer configuration is invalid. Programs that cannot guarantee the validity of the options passed to NewWriter should construct the writer configuration independently prior to calling this function:

config, err := parquet.NewWriterConfig(options...)
if err != nil {
	// handle the configuration error
	...
} else {
	// this call to create a writer is guaranteed not to panic
	writer := parquet.NewWriter(output, config)
	...
}

func (*Writer) BeginRowGroup

func (w *Writer) BeginRowGroup() *ConcurrentRowGroupWriter

BeginRowGroup returns a new ConcurrentRowGroupWriter that can be written to in parallel with other row groups. However these need to be committed back to the writer serially using the Commit method on the row group.

func (*Writer) Close

func (w *Writer) Close() error

Close must be called after all values were produced to the writer in order to flush all buffers and write the parquet footer.

func (*Writer) ColumnWriters

func (w *Writer) ColumnWriters() []*ColumnWriter

ColumnWriters returns writers for each column. This allows applications to write values directly to each column instead of having to first assemble values into rows to use WriteRows.

func (*Writer) File

func (w *Writer) File() FileView

File returns a FileView of the written parquet file. Only available after Close is called.

func (*Writer) Flush

func (w *Writer) Flush() error

Flush flushes all buffers into a row group to the underlying io.Writer.

Flush is called automatically on Close, it is only useful to call explicitly if the application needs to limit the size of row groups or wants to produce multiple row groups per file.

If the writer attempts to create more than MaxRowGroups row groups the method returns ErrTooManyRowGroups.

func (*Writer) ReadRowsFrom

func (w *Writer) ReadRowsFrom(rows RowReader) (written int64, err error)

ReadRowsFrom reads rows from the reader passed as arguments and writes them to w.

This is similar to calling WriteRow repeatedly, but will be more efficient if optimizations are supported by the reader.

func (*Writer) Reset

func (w *Writer) Reset(output io.Writer)

Reset clears the state of the writer without flushing any of the buffers, and setting the output to the io.Writer passed as argument, allowing the writer to be reused to produce another parquet file.

Reset may be called at any time, including after a writer was closed.

func (*Writer) Schema

func (w *Writer) Schema() *Schema

Schema returns the schema of rows written by w.

The returned value will be nil if no schema has yet been configured on w.

func (*Writer) SetKeyValueMetadata

func (w *Writer) SetKeyValueMetadata(key, value string)

SetKeyValueMetadata sets a key/value pair in the Parquet file metadata.

Keys are assumed to be unique, if the same key is repeated multiple times the last value is retained. While the parquet format does not require unique keys, this design decision was made to optimize for the most common use case where applications leverage this extension mechanism to associate single values to keys. This may create incompatibilities with other parquet libraries, or may cause some key/value pairs to be lost when open parquet files written with repeated keys. We can revisit this decision if it ever becomes a blocker.

func (*Writer) Size

func (w *Writer) Size() int64

Size returns an estimate of the current file size in bytes, including all completed row groups and the current in-progress row group.

The estimate includes bytes already written to the underlying writer plus the estimated size of buffered data in the current row group. The buffered data has not been compressed yet, so this is an upper-bound estimate.

The size of the footer metadata is not included because it is variable and typically small relative to the data.

This method can be used to target a specific file size by calling Flush when Size reaches a threshold.

func (*Writer) Write

func (w *Writer) Write(row any) error

Write is called to write another row to the parquet file.

The method uses the parquet schema configured on w to traverse the Go value and decompose it into a set of columns and values. If no schema were passed to NewWriter, it is deducted from the Go type of the row, which then have to be a struct or pointer to struct.

func (*Writer) WriteRowGroup

func (w *Writer) WriteRowGroup(rowGroup RowGroup) (int64, error)

WriteRowGroup writes a row group to the parquet file.

Buffered rows will be flushed prior to writing rows from the group, unless the row group was empty in which case nothing is written to the file.

The content of the row group is flushed to the writer; after the method returns successfully, the row group will be empty and in ready to be reused.

func (*Writer) WriteRows

func (w *Writer) WriteRows(rows []Row) (int, error)

WriteRows is called to write rows to the parquet file.

The Writer must have been given a schema when NewWriter was called, otherwise the structure of the parquet file cannot be determined from the row only.

The row is expected to contain values for each column of the writer's schema, in the order produced by the parquet.(*Schema).Deconstruct method.

type WriterConfig

type WriterConfig struct {
	CreatedBy                    string
	ColumnPageBuffers            BufferPool
	ColumnIndexSizeLimit         func(path []string) int
	PageBufferSize               int
	WriteBufferSize              int
	DataPageVersion              int
	DataPageStatistics           bool
	DeprecatedDataPageStatistics bool
	MaxRowsPerRowGroup           int64
	KeyValueMetadata             map[string]string
	Schema                       *Schema
	BloomFilters                 []BloomFilterColumn
	DeferredBloomFiltersBuffers  BufferPool
	BloomFilterCompression       compress.Codec
	Compression                  compress.Codec
	Sorting                      SortingConfig
	SkipPageBounds               [][]string
	SkipPageStatistics           [][]string
	Encodings                    map[Kind]encoding.Encoding
	DictionaryMaxBytes           int64
	SchemaConfig                 *SchemaConfig
	Encryption                   *EncryptionConfig
}

The WriterConfig type carries configuration options for parquet writers.

WriterConfig implements the WriterOption interface so it can be used directly as argument to the NewWriter function when needed, for example:

writer := parquet.NewWriter(output, schema, &parquet.WriterConfig{
	CreatedBy: "my test program",
})

func DefaultWriterConfig

func DefaultWriterConfig() *WriterConfig

DefaultWriterConfig returns a new WriterConfig value initialized with the default writer configuration.

func NewWriterConfig

func NewWriterConfig(options ...WriterOption) (*WriterConfig, error)

NewWriterConfig constructs a new writer configuration applying the options passed as arguments.

The function returns an non-nil error if some of the options carried invalid configuration values.

func (*WriterConfig) Apply

func (c *WriterConfig) Apply(options ...WriterOption)

Apply applies the given list of options to c.

func (*WriterConfig) ConfigureWriter

func (c *WriterConfig) ConfigureWriter(config *WriterConfig)

ConfigureWriter applies configuration options from c to config.

func (*WriterConfig) Validate

func (c *WriterConfig) Validate() error

Validate returns a non-nil error if the configuration of c is invalid.

type WriterOption

type WriterOption interface {
	ConfigureWriter(*WriterConfig)
}

WriterOption is an interface implemented by types that carry configuration options for parquet writers.

func BloomFilterCompression

func BloomFilterCompression(codec compress.Codec) WriterOption

BloomFilterCompression creates a configuration option which sets the compression codec used when writing bloom filters. The default is uncompressed for backward compatibility.

func BloomFilters

func BloomFilters(filters ...BloomFilterColumn) WriterOption

BloomFilters creates a configuration option which defines the bloom filters that parquet writers should generate.

The compute and memory footprint of generating bloom filters for all columns of a parquet schema can be significant, so by default no filters are created and applications need to explicitly declare the columns that they want to create filters for.

func ColumnIndexSizeLimit

func ColumnIndexSizeLimit(f func(path []string) int) WriterOption

ColumnIndexSizeLimit creates a configuration option to customize the size limit of page boundaries recorded in column indexes. The result of the provided function must be larger then 0.

Defaults to the function that returns 16 for all paths.

func ColumnPageBuffers

func ColumnPageBuffers(buffers BufferPool) WriterOption

ColumnPageBuffers creates a configuration option to customize the buffer pool used when constructing row groups. This can be used to provide on-disk buffers as swap space to ensure that the parquet file creation will no be bottlenecked on the amount of memory available.

Defaults to using in-memory buffers.

func Compression

func Compression(codec compress.Codec) WriterOption

Compression creates a configuration option which sets the default compression codec used by a writer for columns where none were defined.

func CreatedBy

func CreatedBy(application, version, build string) WriterOption

CreatedBy creates a configuration option which sets the name of the application that created a parquet file.

The option formats the "CreatedBy" file metadata according to the convention described by the parquet spec:

"<application> version <version> (build <build>)"

By default, the option is set to the parquet-go module name, version, and build hash.

func DataPageStatistics

func DataPageStatistics(enabled bool) WriterOption

DataPageStatistics creates a configuration option which defines whether data page statistics are emitted. This option is useful when generating parquet files that intend to be backward compatible with older readers which may not have the ability to load page statistics from the column index.

This used to be disabled by default, but it was switched from false to true after v0.26.0, because the computation of page statistics is cheap and query engines do better the more statistics they have available, enabling is a more sensible default.

Defaults to true.

func DataPageVersion

func DataPageVersion(version int) WriterOption

DataPageVersion creates a configuration option which configures the version of data pages used when creating a parquet file.

Defaults to version 2.

func DefaultEncoding

func DefaultEncoding(enc encoding.Encoding) WriterOption

DefaultEncoding creates a configuration option which sets the default encoding used by a writer for columns where none were defined.

It will fail if the specified enconding isn't compatible with any of the primitive types.

func DefaultEncodingFor

func DefaultEncodingFor(kind Kind, enc encoding.Encoding) WriterOption

DefaultEncodingFor creates a configuration option which sets the default encoding used by a writer for columns with the specified primitive type where none were defined.

It will fail if the specified enconding isn't compatible with the specified primitive type.

func DeferBloomFiltersWithBuffers

func DeferBloomFiltersWithBuffers(buffer BufferPool) WriterOption

DeferBloomFiltersWithBuffers creates a configuration option which delays the writing of bloom filters until the end of the file. This can be beneficial for files read from remote storage with a custom reader, as an optimistic read can capture the file footer along with the bloom filters in a single request.

When this option is enabled, the accumulated bloom filters need to be retained until the file is closed; it is therefore required to provide a buffer when using this option.

Defaults to nil; bloom filters are written immediately after each row group.

func DeprecatedDataPageStatistics

func DeprecatedDataPageStatistics(enabled bool) WriterOption

DeprecatedDataPageStatistics creates a configuration option which defines whether to also write the deprecated Min/Max fields in column chunk statistics, in addition to the standard MinValue/MaxValue fields. This option is useful for backward compatibility with older readers that only understand the deprecated fields.

The deprecated Min/Max fields use signed comparison, while MinValue/MaxValue respect the column's logical type ordering. For columns where these orderings differ (e.g., unsigned integers), enabling this option may produce incorrect statistics for older readers.

Defaults to false.

func DictionaryMaxBytes

func DictionaryMaxBytes(size int64) WriterOption

DictionaryMaxBytes creates a configuration option which sets the maximum size in bytes for each column's dictionary.

When a column's dictionary exceeds this limit, that column will switch from dictionary encoding to PLAIN encoding for the remainder of the row group. Pages written before the limit was reached remain dictionary-encoded, while subsequent pages use PLAIN encoding.

A value of 0 (the default) means unlimited dictionary size.

func KeyValueMetadata

func KeyValueMetadata(key, value string) WriterOption

KeyValueMetadata creates a configuration option which adds key/value metadata to add to the metadata of parquet files.

This option is additive, it may be used multiple times to add more than one key/value pair.

Keys are assumed to be unique, if the same key is repeated multiple times the last value is retained. While the parquet format does not require unique keys, this design decision was made to optimize for the most common use case where applications leverage this extension mechanism to associate single values to keys. This may create incompatibilities with other parquet libraries, or may cause some key/value pairs to be lost when open parquet files written with repeated keys. We can revisit this decision if it ever becomes a blocker.

func MaxRowsPerRowGroup

func MaxRowsPerRowGroup(numRows int64) WriterOption

MaxRowsPerRowGroup configures the maximum number of rows that a writer will produce in each row group.

This limit is useful to control size of row groups in both number of rows and byte size. While controlling the byte size of a row group is difficult to achieve with parquet due to column encoding and compression, the number of rows remains a useful proxy.

Defaults to unlimited.

func PageBufferSize

func PageBufferSize(size int) WriterOption

PageBufferSize configures the size of column page buffers on parquet writers.

Note that the page buffer size refers to the in-memory buffers where pages are generated, not the size of pages after encoding and compression. This design choice was made to help control the amount of memory needed to read and write pages rather than controlling the space used by the encoded representation on disk.

Defaults to 256KiB.

func SkipPageBounds

func SkipPageBounds(path ...string) WriterOption

SkipPageBounds lists the path to a column that shouldn't have bounds written to the footer of the parquet file. This is useful for data blobs, like a raw html file, where the bounds are not meaningful.

This option is additive, it may be used multiple times to skip multiple columns.

func SkipPageStatistics

func SkipPageStatistics(path ...string) WriterOption

SkipPageStatistics lists the path to a column that shouldn't have statistics written for pages. This is useful for data blobs, like a raw html file, where the bounds are not meaningful.

This option has no effect if DataPageStatistics(false) is used.

This option is additive, it may be used multiple times to skip multiple columns.

func SortingWriterConfig

func SortingWriterConfig(options ...SortingOption) WriterOption

SortingWriterConfig is a writer option which applies configuration specific to sorting writers.

func WithEncryption

func WithEncryption(cfg *EncryptionConfig) WriterOption

WithEncryption returns a WriterOption that configures encryption.

func WriteBufferSize

func WriteBufferSize(size int) WriterOption

WriteBufferSize configures the size of the write buffer.

Setting the writer buffer size to zero deactivates buffering, all writes are immediately sent to the output io.Writer.

Defaults to 32KiB.

type WriterTarget

type WriterTarget interface {
	Schema() *Schema
	ColumnWriters() []*ColumnWriter
}

WriterTarget is the writer surface VariantColumnWriter binds to; it is satisfied by *Writer, *GenericWriter[T], and *ConcurrentRowGroupWriter.

Source Files

Directories

Path Synopsis
Package bloom implements parquet bloom filters.
Package bloom implements parquet bloom filters.
xxhash
Package xxhash is an extension of github.com/cespare/xxhash which adds routines optimized to hash arrays of fixed size elements.
Package xxhash is an extension of github.com/cespare/xxhash which adds routines optimized to hash arrays of fixed size elements.
Package compress provides the generic APIs implemented by parquet compression codecs.
Package compress provides the generic APIs implemented by parquet compression codecs.
brotli
Package brotli implements the BROTLI parquet compression codec.
Package brotli implements the BROTLI parquet compression codec.
gzip
Package gzip implements the GZIP parquet compression codec.
Package gzip implements the GZIP parquet compression codec.
lz4
Package lz4 implements the LZ4_RAW parquet compression codec.
Package lz4 implements the LZ4_RAW parquet compression codec.
snappy
Package snappy implements the SNAPPY parquet compression codec.
Package snappy implements the SNAPPY parquet compression codec.
uncompressed
Package uncompressed provides implementations of the compression codec interfaces as pass-through without applying any compression nor decompression.
Package uncompressed provides implementations of the compression codec interfaces as pass-through without applying any compression nor decompression.
zstd
Package zstd implements the ZSTD parquet compression codec.
Package zstd implements the ZSTD parquet compression codec.
Package encoding provides the generic APIs implemented by parquet encodings in its sub-packages.
Package encoding provides the generic APIs implemented by parquet encodings in its sub-packages.
fuzz
Package fuzz contains functions to help fuzz test parquet encodings.
Package fuzz contains functions to help fuzz test parquet encodings.
plain
Package plain implements the PLAIN parquet encoding.
Package plain implements the PLAIN parquet encoding.
rle
Package rle implements the hybrid RLE/Bit-Packed encoding employed in repetition and definition levels, dictionary indexed data pages, and boolean values in the PLAIN encoding.
Package rle implements the hybrid RLE/Bit-Packed encoding employed in repetition and definition levels, dictionary indexed data pages, and boolean values in the PLAIN encoding.
Package hashprobe provides implementations of probing tables for various data types.
Package hashprobe provides implementations of probing tables for various data types.
aeshash
Package aeshash implements hashing functions derived from the Go runtime's internal hashing based on the support of AES encryption in CPU instructions.
Package aeshash implements hashing functions derived from the Go runtime's internal hashing based on the support of AES encryption in CPU instructions.
wyhash
Package wyhash implements a hashing algorithm derived from the Go runtime's internal hashing fallback, which uses a variation of the wyhash algorithm.
Package wyhash implements a hashing algorithm derived from the Go runtime's internal hashing fallback, which uses a variation of the wyhash algorithm.
internal
bytealg
Package bytealg contains optimized algorithms operating on byte slices.
Package bytealg contains optimized algorithms operating on byte slices.
unsafecast
Package unsafecast exposes functions to bypass the Go type system and perform conversions between types that would otherwise not be possible.
Package unsafecast exposes functions to bypass the Go type system and perform conversions between types that would otherwise not be possible.
Package sparse contains abstractions to help work on arrays of values in sparse memory locations.
Package sparse contains abstractions to help work on arrays of values in sparse memory locations.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL