Skip to content

Commit 33ae9f5

Browse files
authored
Merge pull request #34 from fanyang89/fix-parquet-scan-panic
fix parquet rows panic on EOF-with-rows boundary
2 parents d2b03d1 + 4c144b1 commit 33ae9f5

3 files changed

Lines changed: 107 additions & 0 deletions

File tree

‎chdb/driver/parquet.go‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -65,6 +65,9 @@ func (r *parquetRows) readNextChunk() error {
6565
return err // no records read, should exit the loop
6666
}
6767
if err == io.EOF && readAmount > 0 {
68+
r.buffer = r.buffer[:readAmount]
69+
r.bufferIndex = 0
70+
r.needNewBuffer = false
6871
return nil //here we are at EOF, but since we read at least 1 record, we should consume it
6972
}
7073
if readAmount == 0 {
Lines changed: 99 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,99 @@
1+
package chdbdriver
2+
3+
import (
4+
"database/sql/driver"
5+
"io"
6+
"testing"
7+
8+
chdbpurego "github.com/chdb-io/chdb-go/chdb-purego"
9+
"github.com/parquet-go/parquet-go"
10+
)
11+
12+
type fakeResult struct{}
13+
14+
func (fakeResult) Buf() []byte { return nil }
15+
func (fakeResult) String() string { return "" }
16+
func (fakeResult) Len() int { return 0 }
17+
func (fakeResult) Elapsed() float64 { return 0 }
18+
func (fakeResult) RowsRead() uint64 { return 1 }
19+
func (fakeResult) BytesRead() uint64 { return 0 }
20+
func (fakeResult) Error() error { return nil }
21+
func (fakeResult) Free() {}
22+
23+
var _ chdbpurego.ChdbResult = (*fakeResult)(nil)
24+
25+
type eofRowGroup struct {
26+
schema *parquet.Schema
27+
}
28+
29+
func (g *eofRowGroup) NumRows() int64 { return 4 }
30+
func (g *eofRowGroup) ColumnChunks() []parquet.ColumnChunk { return nil }
31+
func (g *eofRowGroup) Schema() *parquet.Schema { return g.schema }
32+
func (g *eofRowGroup) SortingColumns() []parquet.SortingColumn { return nil }
33+
func (g *eofRowGroup) Rows() parquet.Rows { return &eofRows{schema: g.schema} }
34+
35+
type eofRows struct {
36+
schema *parquet.Schema
37+
phase int
38+
next int64
39+
}
40+
41+
func (r *eofRows) ReadRows(rows []parquet.Row) (int, error) {
42+
if len(rows) == 0 {
43+
return 0, nil
44+
}
45+
46+
fill := func(n int, err error) (int, error) {
47+
if n > len(rows) {
48+
n = len(rows)
49+
}
50+
for i := 0; i < n; i++ {
51+
rows[i] = parquet.Row{parquet.ValueOf(r.next).Level(0, 0, 0)}
52+
r.next++
53+
}
54+
return n, err
55+
}
56+
57+
switch r.phase {
58+
case 0:
59+
r.phase = 1
60+
return fill(2, nil)
61+
case 1:
62+
r.phase = 2
63+
return fill(2, io.EOF)
64+
default:
65+
return 0, io.EOF
66+
}
67+
}
68+
69+
func (r *eofRows) SeekToRow(int64) error { return nil }
70+
71+
func (r *eofRows) Close() error { return nil }
72+
73+
func (r *eofRows) Schema() *parquet.Schema { return r.schema }
74+
75+
func TestParquetNextHandlesEOFWithRemainingRows(t *testing.T) {
76+
schema := parquet.SchemaOf(struct {
77+
Number int64 `parquet:"number"`
78+
}{})
79+
80+
reader := parquet.NewGenericRowGroupReader[any](&eofRowGroup{schema: schema})
81+
rows := &parquetRows{
82+
reader: reader,
83+
schemaFields: schema.Fields(),
84+
bufferSize: 2,
85+
needNewBuffer: true,
86+
localResult: fakeResult{},
87+
}
88+
89+
dest := make([]driver.Value, 1)
90+
for i := 0; i < 4; i++ {
91+
if err := rows.Next(dest); err != nil {
92+
t.Fatalf("rows.Next failed at row %d, err: %v", i, err)
93+
}
94+
}
95+
96+
if err := rows.Next(dest); err != io.EOF {
97+
t.Fatalf("expected io.EOF after 4 rows, actual: %v", err)
98+
}
99+
}

‎chdb/driver/parquet_streaming.go‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -58,6 +58,11 @@ func (r *parquetStreamingRows) readNextChunkFromBuf() error {
5858
return err // no records read, should exit the loop
5959
}
6060
if err == io.EOF && readAmount > 0 {
61+
if readAmount < r.bufferSize {
62+
r.buffer = r.buffer[:readAmount]
63+
}
64+
r.bufferIndex = 0
65+
r.needNewBuffer = false
6166
return nil //here we are at EOF, but since we read at least 1 record, we should consume it
6267
}
6368
if readAmount == 0 {

0 commit comments

Comments
 (0)