mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-08-25 00:27:11 +00:00
format: add parquet adapter cutting extents at row-group starts
The footer already names every row-group byte range, so the adapter only reads metadata: one extent per row group, the leading magic riding with the first, and a trailing extent for the page indexes and footer. Engines that fetch row groups by offset then read exactly the covering chunks. Parquet needs no view; alignment alone delivers the benefit.
This commit is contained in:
@@ -0,0 +1,93 @@
|
||||
// Package parquet adapts parquet files: extents are cut at row-group starts,
|
||||
// with one trailing extent for the page indexes and footer. Readers that fetch
|
||||
// row groups by offset then hit exactly the covering chunks. No view is
|
||||
// needed; the whole benefit is delivered by alignment.
|
||||
package parquet
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"math"
|
||||
|
||||
parquetgo "github.com/parquet-go/parquet-go"
|
||||
"github.com/seaweedfs/seaweedfs/weed/format"
|
||||
)
|
||||
|
||||
const FormatName = "parquet"
|
||||
|
||||
var magic = []byte("PAR1")
|
||||
|
||||
func init() {
|
||||
format.Register(Adapter{})
|
||||
}
|
||||
|
||||
type Adapter struct{}
|
||||
|
||||
func (Adapter) Name() string { return FormatName }
|
||||
|
||||
func (Adapter) Sniff(h format.Hint) bool {
|
||||
return bytes.HasPrefix(h.Head, magic) && bytes.HasSuffix(h.Tail, magic)
|
||||
}
|
||||
|
||||
// Index reads the footer and cuts one extent per row group. The leading magic
|
||||
// rides with the first row group; everything after the last row group (page
|
||||
// indexes, footer) forms the final extent.
|
||||
func (Adapter) Index(ctx context.Context, r io.ReaderAt, size int64) (*format.Layout, error) {
|
||||
file, err := parquetgo.OpenFile(r, size, parquetgo.SkipPageIndex(true), parquetgo.SkipBloomFilters(true))
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("open parquet: %w", err)
|
||||
}
|
||||
rowGroups := file.Metadata().RowGroups
|
||||
if len(rowGroups) == 0 {
|
||||
return nil, fmt.Errorf("parquet file has no row groups")
|
||||
}
|
||||
if len(rowGroups) >= format.MaxExtentCount {
|
||||
return nil, fmt.Errorf("parquet file has too many row groups: %d", len(rowGroups))
|
||||
}
|
||||
|
||||
starts := make([]int64, 0, len(rowGroups))
|
||||
var lastEnd int64
|
||||
for i, rowGroup := range rowGroups {
|
||||
start, end := int64(math.MaxInt64), int64(0)
|
||||
for _, column := range rowGroup.Columns {
|
||||
columnStart := column.MetaData.DataPageOffset
|
||||
if dictionary := column.MetaData.DictionaryPageOffset; dictionary > 0 && dictionary < columnStart {
|
||||
columnStart = dictionary
|
||||
}
|
||||
if columnStart < start {
|
||||
start = columnStart
|
||||
}
|
||||
if columnEnd := columnStart + column.MetaData.TotalCompressedSize; columnEnd > end {
|
||||
end = columnEnd
|
||||
}
|
||||
}
|
||||
if len(rowGroup.Columns) == 0 || start >= end || start < int64(len(magic)) || end > size {
|
||||
return nil, fmt.Errorf("row group %d has an invalid byte range [%d, %d)", i, start, end)
|
||||
}
|
||||
if i > 0 && start < lastEnd {
|
||||
return nil, fmt.Errorf("row group %d overlaps its predecessor", i)
|
||||
}
|
||||
starts = append(starts, start)
|
||||
lastEnd = end
|
||||
}
|
||||
if lastEnd >= size {
|
||||
return nil, fmt.Errorf("row groups leave no room for the footer")
|
||||
}
|
||||
|
||||
var sizes []int64
|
||||
var previous int64
|
||||
for _, start := range starts[1:] {
|
||||
sizes = append(sizes, start-previous)
|
||||
previous = start
|
||||
}
|
||||
sizes = append(sizes, lastEnd-previous)
|
||||
sizes = append(sizes, size-lastEnd)
|
||||
|
||||
layout := &format.Layout{Format: FormatName, ExtentSizes: sizes, Align: 1}
|
||||
if err := layout.Validate(size); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return layout, nil
|
||||
}
|
||||
@@ -0,0 +1,108 @@
|
||||
package parquet
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
parquetgo "github.com/parquet-go/parquet-go"
|
||||
"github.com/seaweedfs/seaweedfs/weed/format"
|
||||
"github.com/seaweedfs/seaweedfs/weed/format/formattest"
|
||||
)
|
||||
|
||||
type row struct {
|
||||
ID int64 `parquet:"id"`
|
||||
Name string `parquet:"name"`
|
||||
}
|
||||
|
||||
// buildParquet writes rowGroups row groups of rowsPerGroup rows each.
|
||||
func buildParquet(t *testing.T, rowGroups, rowsPerGroup int) []byte {
|
||||
t.Helper()
|
||||
var buf bytes.Buffer
|
||||
writer := parquetgo.NewGenericWriter[row](&buf)
|
||||
for g := 0; g < rowGroups; g++ {
|
||||
rows := make([]row, rowsPerGroup)
|
||||
for i := range rows {
|
||||
rows[i] = row{ID: int64(g*rowsPerGroup + i), Name: "some filler content to give row groups a little size"}
|
||||
}
|
||||
if _, err := writer.Write(rows); err != nil {
|
||||
t.Fatalf("Write() error = %v", err)
|
||||
}
|
||||
if err := writer.Flush(); err != nil {
|
||||
t.Fatalf("Flush() error = %v", err)
|
||||
}
|
||||
}
|
||||
if err := writer.Close(); err != nil {
|
||||
t.Fatalf("Close() error = %v", err)
|
||||
}
|
||||
return buf.Bytes()
|
||||
}
|
||||
|
||||
func TestIndex(t *testing.T) {
|
||||
data := buildParquet(t, 3, 100)
|
||||
size := int64(len(data))
|
||||
layout, err := Adapter{}.Index(context.Background(), bytes.NewReader(data), size)
|
||||
if err != nil {
|
||||
t.Fatalf("Index() error = %v", err)
|
||||
}
|
||||
// one extent per row group plus the trailing footer extent
|
||||
if len(layout.ExtentSizes) != 4 {
|
||||
t.Fatalf("extents = %v, want 4", layout.ExtentSizes)
|
||||
}
|
||||
if err := layout.Validate(size); err != nil {
|
||||
t.Fatalf("Validate() error = %v", err)
|
||||
}
|
||||
formattest.EncodeRoundTrip(t, layout)
|
||||
|
||||
// extent cuts must land on the row-group starts the footer declares
|
||||
file, err := parquetgo.OpenFile(bytes.NewReader(data), size)
|
||||
if err != nil {
|
||||
t.Fatalf("OpenFile() error = %v", err)
|
||||
}
|
||||
var offset int64
|
||||
for i, extentSize := range layout.ExtentSizes[:3] {
|
||||
if i > 0 {
|
||||
want := file.Metadata().RowGroups[i].Columns[0].MetaData.DataPageOffset
|
||||
if dictionary := file.Metadata().RowGroups[i].Columns[0].MetaData.DictionaryPageOffset; dictionary > 0 && dictionary < want {
|
||||
want = dictionary
|
||||
}
|
||||
if offset != want {
|
||||
t.Fatalf("extent %d starts at %d, row group starts at %d", i, offset, want)
|
||||
}
|
||||
}
|
||||
offset += extentSize
|
||||
}
|
||||
}
|
||||
|
||||
func TestIndexSingleRowGroup(t *testing.T) {
|
||||
data := buildParquet(t, 1, 10)
|
||||
layout, err := Adapter{}.Index(context.Background(), bytes.NewReader(data), int64(len(data)))
|
||||
if err != nil {
|
||||
t.Fatalf("Index() error = %v", err)
|
||||
}
|
||||
if len(layout.ExtentSizes) != 2 {
|
||||
t.Fatalf("extents = %v, want 2", layout.ExtentSizes)
|
||||
}
|
||||
}
|
||||
|
||||
func TestIndexRejectsNonParquet(t *testing.T) {
|
||||
data := []byte("this is not a parquet file, not even close, but long enough")
|
||||
if _, err := (Adapter{}).Index(context.Background(), bytes.NewReader(data), int64(len(data))); err == nil {
|
||||
t.Fatalf("Index() accepted junk")
|
||||
}
|
||||
}
|
||||
|
||||
func TestIndexTruncations(t *testing.T) {
|
||||
formattest.IndexTruncations(t, Adapter{}, buildParquet(t, 2, 50))
|
||||
}
|
||||
|
||||
func TestSniff(t *testing.T) {
|
||||
data := buildParquet(t, 1, 10)
|
||||
head, tail := data[:4], data[len(data)-4:]
|
||||
if !(Adapter{}).Sniff(format.Hint{Head: head, Tail: tail}) {
|
||||
t.Fatalf("Sniff() rejected parquet magic")
|
||||
}
|
||||
if (Adapter{}).Sniff(format.Hint{Head: []byte("PAR1"), Tail: []byte("nope")}) {
|
||||
t.Fatalf("Sniff() accepted missing tail magic")
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user