Skip to content

Commit 376b140

Browse files
authored
Merge branch 'main' into fix/csv/supported-types
2 parents 3f47c1f + 8cc722e commit 376b140

4 files changed

Lines changed: 29 additions & 25 deletions

File tree

go.mod

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@ require (
66
github.com/apache/arrow/go/v16 v16.0.0
77
github.com/bradleyjkemp/cupaloy/v2 v2.8.0
88
github.com/cloudquery/codegen v0.3.16
9-
github.com/cloudquery/plugin-sdk/v4 v4.40.2
9+
github.com/cloudquery/plugin-sdk/v4 v4.42.1
1010
github.com/goccy/go-json v0.10.2
1111
github.com/invopop/jsonschema v0.12.0
1212
github.com/stretchr/testify v1.9.0
@@ -29,7 +29,7 @@ require (
2929
github.com/bytedance/sonic v1.11.2 // indirect
3030
github.com/chenzhuoyu/base64x v0.0.0-20230717121745-296ad89f973d // indirect
3131
github.com/chenzhuoyu/iasm v0.9.1 // indirect
32-
github.com/cloudquery/cloudquery-api-go v1.9.1 // indirect
32+
github.com/cloudquery/cloudquery-api-go v1.11.1 // indirect
3333
github.com/davecgh/go-spew v1.1.1 // indirect
3434
github.com/deepmap/oapi-codegen v1.16.2 // indirect
3535
github.com/fatih/structs v1.1.0 // indirect
@@ -103,7 +103,7 @@ require (
103103
golang.org/x/xerrors v0.0.0-20231012003039-104605ab7028 // indirect
104104
google.golang.org/genproto/googleapis/rpc v0.0.0-20240401170217-c3f982113cda // indirect
105105
google.golang.org/grpc v1.63.2 // indirect
106-
google.golang.org/protobuf v1.34.0 // indirect
106+
google.golang.org/protobuf v1.34.1 // indirect
107107
gopkg.in/ini.v1 v1.67.0 // indirect
108108
gopkg.in/yaml.v3 v3.0.1 // indirect
109109
)

go.sum

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -43,14 +43,14 @@ github.com/chenzhuoyu/base64x v0.0.0-20230717121745-296ad89f973d/go.mod h1:8EPpV
4343
github.com/chenzhuoyu/iasm v0.9.0/go.mod h1:Xjy2NpN3h7aUqeqM+woSuuvxmIe6+DDsiNLIrkAmYog=
4444
github.com/chenzhuoyu/iasm v0.9.1 h1:tUHQJXo3NhBqw6s33wkGn9SP3bvrWLdlVIJ3hQBL7P0=
4545
github.com/chenzhuoyu/iasm v0.9.1/go.mod h1:Xjy2NpN3h7aUqeqM+woSuuvxmIe6+DDsiNLIrkAmYog=
46-
github.com/cloudquery/cloudquery-api-go v1.9.1 h1:Nq6SnE4V9A8YprLMXO/8QszWJMNloqW6c2G00ADYoI4=
47-
github.com/cloudquery/cloudquery-api-go v1.9.1/go.mod h1:F4kuaNBAVqsS9ZRHuX+tV2m6+Khoa2Rb9lROGhinGPk=
46+
github.com/cloudquery/cloudquery-api-go v1.11.1 h1:R7f+Lk16Exx0FAIx+0XuFC35e4UhJXctCxCubPOxitc=
47+
github.com/cloudquery/cloudquery-api-go v1.11.1/go.mod h1:F4kuaNBAVqsS9ZRHuX+tV2m6+Khoa2Rb9lROGhinGPk=
4848
github.com/cloudquery/codegen v0.3.16 h1:kZLOuVvEIHuk6QoFlRWQvUr5XRMxzYsDXpsNnF8SJHQ=
4949
github.com/cloudquery/codegen v0.3.16/go.mod h1:NOLLrXLTKpiJ3z7d11HiS4vIT+HkKaKe+q3USwuq4+E=
5050
github.com/cloudquery/jsonschema v0.0.0-20240220124159-92878faa2a66 h1:OZLPSIBYEfvkAUeOeM8CwTgVQy5zhayI99ishCrsFV0=
5151
github.com/cloudquery/jsonschema v0.0.0-20240220124159-92878faa2a66/go.mod h1:0SoZ/U7yJlNOR+fWsBSeTvTbGXB6DK01tzJ7m2Xfg34=
52-
github.com/cloudquery/plugin-sdk/v4 v4.40.2 h1:i9ZMIrgS6G5Sa0v26C/Y2rE/WUaZUaMRLYCRJgCRau0=
53-
github.com/cloudquery/plugin-sdk/v4 v4.40.2/go.mod h1:lBg32JZd7sVJy4IOrGEZrJPI+EnrRn9QEu2ovYL36iE=
52+
github.com/cloudquery/plugin-sdk/v4 v4.42.1 h1:m6UblpYDKDOyoFXB/cbq35rM0OLETYhS0tIGvWTFkB0=
53+
github.com/cloudquery/plugin-sdk/v4 v4.42.1/go.mod h1:NdMMpHZAsmNPpACaIjjlGZtu7uPdHA364oWmyT2deGI=
5454
github.com/coreos/go-systemd/v22 v22.5.0/go.mod h1:Y58oyj3AT4RCenI/lSvhwexgC+NSVTIJ3seZv2GcEnc=
5555
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
5656
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
@@ -313,8 +313,8 @@ google.golang.org/genproto/googleapis/rpc v0.0.0-20240401170217-c3f982113cda h1:
313313
google.golang.org/genproto/googleapis/rpc v0.0.0-20240401170217-c3f982113cda/go.mod h1:WtryC6hu0hhx87FDGxWCDptyssuo68sk10vYjF+T9fY=
314314
google.golang.org/grpc v1.63.2 h1:MUeiw1B2maTVZthpU5xvASfTh3LDbxHd6IJ6QQVU+xM=
315315
google.golang.org/grpc v1.63.2/go.mod h1:WAX/8DgncnokcFUldAxq7GeB5DXHDbMF+lLvDomNkRA=
316-
google.golang.org/protobuf v1.34.0 h1:Qo/qEd2RZPCf2nKuorzksSknv0d3ERwp1vFG38gSmH4=
317-
google.golang.org/protobuf v1.34.0/go.mod h1:c6P6GXX6sHbq/GpV6MGZEdwhWPcYBgnhAHhKbcUYpos=
316+
google.golang.org/protobuf v1.34.1 h1:9ddQBjfCyZPOHPUiPxpYESBLc+T8P3E+Vo4IbKZgFWg=
317+
google.golang.org/protobuf v1.34.1/go.mod h1:c6P6GXX6sHbq/GpV6MGZEdwhWPcYBgnhAHhKbcUYpos=
318318
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
319319
gopkg.in/check.v1 v1.0.0-20200902074654-038fdea0a05b h1:QRR6H1YWRnHb4Y/HeNFCTJLFVxaq6wH4YuVdsUOr75U=
320320
gopkg.in/check.v1 v1.0.0-20200902074654-038fdea0a05b/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=

parquet/read.go

Lines changed: 6 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -35,9 +35,7 @@ func (*Client) Read(f types.ReaderAtSeeker, table *schema.Table, res chan<- arro
3535

3636
sc := table.ToArrowSchema()
3737
for rr.Next() {
38-
for _, r := range slice(reverseTransformRecord(sc, rr.Record())) {
39-
res <- r
40-
}
38+
res <- reverseTransformRecord(sc, rr.Record())
4139
}
4240
if rr.Err() != nil && rr.Err() != io.EOF {
4341
return fmt.Errorf("failed to read parquet record: %w", rr.Err())
@@ -46,14 +44,6 @@ func (*Client) Read(f types.ReaderAtSeeker, table *schema.Table, res chan<- arro
4644
return nil
4745
}
4846

49-
func slice(r arrow.Record) []arrow.Record {
50-
res := make([]arrow.Record, r.NumRows())
51-
for i := int64(0); i < r.NumRows(); i++ {
52-
res[i] = r.NewSlice(i, i+1)
53-
}
54-
return res
55-
}
56-
5747
func reverseTransformRecord(sc *arrow.Schema, rec arrow.Record) arrow.Record {
5848
cols := make([]arrow.Array, rec.NumCols())
5949
for i := 0; i < int(rec.NumCols()); i++ {
@@ -63,6 +53,10 @@ func reverseTransformRecord(sc *arrow.Schema, rec arrow.Record) arrow.Record {
6353
}
6454

6555
func reverseTransformArray(dt arrow.DataType, arr arrow.Array) arrow.Array {
56+
if arrow.TypeEqual(dt, arr.DataType()) {
57+
return arr
58+
}
59+
6660
switch arr := arr.(type) {
6761
case *array.String:
6862
return reverseTransformFromString(dt, arr)
@@ -102,7 +96,7 @@ func reverseTransformArray(dt arrow.DataType, arr arrow.Array) arrow.Array {
10296
))
10397

10498
default:
105-
return arr
99+
panic(fmt.Errorf("unsupported conversion from %s to %s", arr.DataType(), dt))
106100
}
107101
}
108102

parquet/write_read_test.go

Lines changed: 14 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -54,13 +54,14 @@ func TestWriteRead(t *testing.T) {
5454
readErr = cl.Read(byteReader, table, ch)
5555
close(ch)
5656
}()
57-
received := make([]arrow.Record, 0, rows)
57+
received, total := make([]arrow.Record, 0, rows), 0
5858
for got := range ch {
5959
received = append(received, got)
60+
total += int(got.NumRows())
6061
}
6162
require.Empty(t, plugin.RecordsDiff(table.ToArrowSchema(), []arrow.Record{record}, received))
6263
require.NoError(t, readErr)
63-
require.Equalf(t, rows, len(received), "got %d row(s), want %d", len(received), rows)
64+
require.Equalf(t, rows, total, "got %d row(s), want %d", total, rows)
6465
}
6566
func TestWriteReadSliced(t *testing.T) {
6667
const rows = 10
@@ -102,13 +103,22 @@ func TestWriteReadSliced(t *testing.T) {
102103
readErr = cl.Read(byteReader, table, ch)
103104
close(ch)
104105
}()
105-
received := make([]arrow.Record, 0, rows)
106+
received, total := make([]arrow.Record, 0, rows), 0
106107
for got := range ch {
107108
received = append(received, got)
109+
total += int(got.NumRows())
108110
}
109111
require.Empty(t, plugin.RecordsDiff(table.ToArrowSchema(), []arrow.Record{record}, received))
110112
require.NoError(t, readErr)
111-
require.Equalf(t, rows, len(received), "got %d row(s), want %d", len(received), rows)
113+
require.Equalf(t, rows, total, "got %d row(s), want %d", total, rows)
114+
}
115+
116+
func slice(r arrow.Record) []arrow.Record {
117+
res := make([]arrow.Record, r.NumRows())
118+
for i := int64(0); i < r.NumRows(); i++ {
119+
res[i] = r.NewSlice(i, i+1)
120+
}
121+
return res
112122
}
113123

114124
func BenchmarkWrite(b *testing.B) {

0 commit comments

Comments
 (0)