This is an automated email from the ASF dual-hosted git repository.

zeroshade pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/arrow-go.git


The following commit(s) were added to refs/heads/main by this push:
     new 7b361b72 perf(parquet/pqarrow): avoid reflection in integer transfers 
(#1163)
7b361b72 is described below

commit 7b361b72f12f0b41f3c7c4f9d2c7b98ff4f433c3
Author: Minh Vu <[email protected]>
AuthorDate: Wed Aug 12 21:41:47 2026 +0200

    perf(parquet/pqarrow): avoid reflection in integer transfers (#1163)
    
    ### Rationale for this change
    
    Parquet to Arrow conversions for narrow and unsigned integers currently
    use reflection for every value. Integer-backed decimal conversions do
    the same when reading their source values. This per-value reflection is
    expensive for large columns.
    
    ### What changes are included in this PR?
    
    - Replace reflective integer assignments with typed conversion loops
    selected once per source and destination type.
    - Use typed INT32 and INT64 source loops for Decimal128 and Decimal256
    conversions.
    - Add boundary-value coverage for every existing integer and temporal
    destination type.
    - Add focused benchmarks for integer and decimal transfers.
    
    Results for 1 million values on an Apple M1 Pro:
    
    | Conversion | Before | After | Improvement |
    | --- | ---: | ---: | ---: |
    | INT32 to INT8 | 3.28 ms | 0.36 ms | 9.0x |
    | INT64 to UINT64 | 3.75 ms | 0.51 ms | 7.3x |
    | INT64 to Decimal128 | 3.97 ms | 1.23 ms | 3.2x |
    | INT64 to Decimal256 | 4.60 ms | 2.48 ms | 1.9x |
    
    Allocation counts are unchanged.
    
    ### Are these changes tested?
    
    Yes. The focused transfer tests, the full `parquet/pqarrow` suite, all
    `parquet/...` packages, and the focused race test pass.
    
    ### Are there any user-facing changes?
    
    No. This only improves Parquet read performance.
---
 parquet/pqarrow/column_readers.go                  | 121 ++++++++-------
 .../pqarrow/column_readers_transfer_bench_test.go  | 111 ++++++++++++++
 parquet/pqarrow/column_readers_transfer_test.go    | 170 +++++++++++++++++++++
 3 files changed, 341 insertions(+), 61 deletions(-)

diff --git a/parquet/pqarrow/column_readers.go 
b/parquet/pqarrow/column_readers.go
index 5aedeb47..5757e863 100644
--- a/parquet/pqarrow/column_readers.go
+++ b/parquet/pqarrow/column_readers.go
@@ -21,7 +21,6 @@ import (
        "encoding/binary"
        "errors"
        "fmt"
-       "reflect"
        "sync"
        "sync/atomic"
        "time"
@@ -776,64 +775,58 @@ func transferBinary(rdr file.RecordReader, dt 
arrow.DataType, mem memory.Allocat
        return arrow.NewChunked(dt, chunks), nil
 }
 
-func transferInt(rdr file.RecordReader, dt arrow.DataType) arrow.ArrayData {
-       var output reflect.Value
+type parquetInteger interface {
+       ~int32 | ~int64
+}
 
-       signed := true
-       // create buffer for proper type since parquet only has int32 and int64
-       // physical representations, but we want the correct type representation
-       // for Arrow's in memory buffer.
-       data := make([]byte, 
rdr.ValuesWritten()*int(bitutil.BytesForBits(int64(dt.(arrow.FixedWidthDataType).BitWidth()))))
-       switch dt.ID() {
+type arrowInteger interface {
+       ~int8 | ~uint8 | ~int16 | ~uint16 | ~int32 | ~uint32 | ~int64 | ~uint64
+}
+
+func convertIntegerValues[Out arrowInteger, In parquetInteger](out []Out, 
values []In) {
+       for i, value := range values {
+               out[i] = Out(value)
+       }
+}
+
+func transferIntegerValues[In parquetInteger](values []In, data []byte, dt 
arrow.Type) {
+       switch dt {
        case arrow.INT8:
-               output = reflect.ValueOf(arrow.Int8Traits.CastFromBytes(data))
+               convertIntegerValues(arrow.Int8Traits.CastFromBytes(data), 
values)
        case arrow.UINT8:
-               signed = false
-               output = reflect.ValueOf(arrow.Uint8Traits.CastFromBytes(data))
+               convertIntegerValues(arrow.Uint8Traits.CastFromBytes(data), 
values)
        case arrow.INT16:
-               output = reflect.ValueOf(arrow.Int16Traits.CastFromBytes(data))
+               convertIntegerValues(arrow.Int16Traits.CastFromBytes(data), 
values)
        case arrow.UINT16:
-               signed = false
-               output = reflect.ValueOf(arrow.Uint16Traits.CastFromBytes(data))
+               convertIntegerValues(arrow.Uint16Traits.CastFromBytes(data), 
values)
        case arrow.UINT32:
-               signed = false
-               output = reflect.ValueOf(arrow.Uint32Traits.CastFromBytes(data))
+               convertIntegerValues(arrow.Uint32Traits.CastFromBytes(data), 
values)
        case arrow.UINT64:
-               signed = false
-               output = reflect.ValueOf(arrow.Uint64Traits.CastFromBytes(data))
+               convertIntegerValues(arrow.Uint64Traits.CastFromBytes(data), 
values)
        case arrow.DATE32:
-               output = reflect.ValueOf(arrow.Date32Traits.CastFromBytes(data))
+               convertIntegerValues(arrow.Date32Traits.CastFromBytes(data), 
values)
        case arrow.TIME32:
-               output = reflect.ValueOf(arrow.Time32Traits.CastFromBytes(data))
+               convertIntegerValues(arrow.Time32Traits.CastFromBytes(data), 
values)
        case arrow.TIME64:
-               output = reflect.ValueOf(arrow.Time64Traits.CastFromBytes(data))
+               convertIntegerValues(arrow.Time64Traits.CastFromBytes(data), 
values)
        }
+}
+
+func transferInt(rdr file.RecordReader, dt arrow.DataType) arrow.ArrayData {
+       // create buffer for proper type since parquet only has int32 and int64
+       // physical representations, but we want the correct type representation
+       // for Arrow's in memory buffer.
+       data := make([]byte, 
rdr.ValuesWritten()*int(bitutil.BytesForBits(int64(dt.(arrow.FixedWidthDataType).BitWidth()))))
 
        length := rdr.ValuesWritten()
        // copy the values semantically with the correct types
        switch rdr.Type() {
        case parquet.Types.Int32:
-               values := arrow.Int32Traits.CastFromBytes(rdr.Values())
-               if signed {
-                       for idx, v := range values[:length] {
-                               output.Index(idx).SetInt(int64(v))
-                       }
-               } else {
-                       for idx, v := range values[:length] {
-                               output.Index(idx).SetUint(uint64(v))
-                       }
-               }
+               values := arrow.Int32Traits.CastFromBytes(rdr.Values())[:length]
+               transferIntegerValues(values, data, dt.ID())
        case parquet.Types.Int64:
-               values := arrow.Int64Traits.CastFromBytes(rdr.Values())
-               if signed {
-                       for idx, v := range values[:length] {
-                               output.Index(idx).SetInt(v)
-                       }
-               } else {
-                       for idx, v := range values[:length] {
-                               output.Index(idx).SetUint(uint64(v))
-                       }
-               }
+               values := arrow.Int64Traits.CastFromBytes(rdr.Values())[:length]
+               transferIntegerValues(values, data, dt.ID())
        }
 
        bitmap := rdr.ReleaseValidBits()
@@ -919,31 +912,37 @@ func transferInt96(rdr file.RecordReader, dt 
arrow.DataType) arrow.ArrayData {
 }
 
 // convert physical integer storage of a decimal logical type to a decimal128 
typed array
-func transferDecimalInteger(rdr file.RecordReader, dt arrow.DataType) 
arrow.ArrayData {
-       length := rdr.ValuesWritten()
-
-       var values reflect.Value
-       switch rdr.Type() {
-       case parquet.Types.Int32:
-               values = 
reflect.ValueOf(arrow.Int32Traits.CastFromBytes(rdr.Values())[:length])
-       case parquet.Types.Int64:
-               values = 
reflect.ValueOf(arrow.Int64Traits.CastFromBytes(rdr.Values())[:length])
-       }
-
-       var data []byte
-       switch dt.ID() {
+func transferDecimalIntegerValues[In parquetInteger](values []In, dt 
arrow.Type) []byte {
+       switch dt {
        case arrow.DECIMAL128:
-               data = make([]byte, 
arrow.Decimal128Traits.BytesRequired(length))
+               data := make([]byte, 
arrow.Decimal128Traits.BytesRequired(len(values)))
                out := arrow.Decimal128Traits.CastFromBytes(data)
-               for i := 0; i < values.Len(); i++ {
-                       out[i] = decimal128.FromI64(values.Index(i).Int())
+               for i, value := range values {
+                       out[i] = decimal128.FromI64(int64(value))
                }
+               return data
        case arrow.DECIMAL256:
-               data = make([]byte, 
arrow.Decimal256Traits.BytesRequired(length))
+               data := make([]byte, 
arrow.Decimal256Traits.BytesRequired(len(values)))
                out := arrow.Decimal256Traits.CastFromBytes(data)
-               for i := 0; i < values.Len(); i++ {
-                       out[i] = decimal256.FromI64(values.Index(i).Int())
+               for i, value := range values {
+                       out[i] = decimal256.FromI64(int64(value))
                }
+               return data
+       }
+       return nil
+}
+
+func transferDecimalInteger(rdr file.RecordReader, dt arrow.DataType) 
arrow.ArrayData {
+       length := rdr.ValuesWritten()
+
+       var data []byte
+       switch rdr.Type() {
+       case parquet.Types.Int32:
+               values := arrow.Int32Traits.CastFromBytes(rdr.Values())[:length]
+               data = transferDecimalIntegerValues(values, dt.ID())
+       case parquet.Types.Int64:
+               values := arrow.Int64Traits.CastFromBytes(rdr.Values())[:length]
+               data = transferDecimalIntegerValues(values, dt.ID())
        }
 
        var nullmap *memory.Buffer
diff --git a/parquet/pqarrow/column_readers_transfer_bench_test.go 
b/parquet/pqarrow/column_readers_transfer_bench_test.go
new file mode 100644
index 00000000..85a4d0b6
--- /dev/null
+++ b/parquet/pqarrow/column_readers_transfer_bench_test.go
@@ -0,0 +1,111 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+package pqarrow
+
+import (
+       "testing"
+
+       "github.com/apache/arrow-go/v18/arrow"
+       "github.com/apache/arrow-go/v18/parquet"
+)
+
+func BenchmarkTransferInteger(b *testing.B) {
+       for _, tc := range []struct {
+               name string
+               size int
+       }{{"1K", 1_000}, {"64K", 64_000}, {"1M", 1_000_000}} {
+               b.Run("Int32ToInt8/"+tc.name, func(b *testing.B) {
+                       size := tc.size
+                       values := make([]int32, size)
+                       for i := range values {
+                               values[i] = int32(i*31 - 1_000_000)
+                       }
+                       rdr := &integerTransferRecordReader{
+                               physicalType: parquet.Types.Int32,
+                               values:       
arrow.Int32Traits.CastToBytes(values),
+                               length:       len(values),
+                       }
+
+                       b.ReportAllocs()
+                       b.SetBytes(int64(len(values) * arrow.Int32SizeBytes))
+                       b.ResetTimer()
+                       for range b.N {
+                               data := transferInt(rdr, 
arrow.PrimitiveTypes.Int8)
+                               data.Release()
+                       }
+               })
+
+               b.Run("Int64ToUint64/"+tc.name, func(b *testing.B) {
+                       size := tc.size
+                       values := make([]int64, size)
+                       for i := range values {
+                               values[i] = int64(i)*6364136223846793005 - 1
+                       }
+                       rdr := &integerTransferRecordReader{
+                               physicalType: parquet.Types.Int64,
+                               values:       
arrow.Int64Traits.CastToBytes(values),
+                               length:       len(values),
+                       }
+
+                       b.ReportAllocs()
+                       b.SetBytes(int64(len(values) * arrow.Int64SizeBytes))
+                       b.ResetTimer()
+                       for range b.N {
+                               data := transferInt(rdr, 
arrow.PrimitiveTypes.Uint64)
+                               data.Release()
+                       }
+               })
+       }
+}
+
+func BenchmarkTransferDecimalInteger(b *testing.B) {
+       for _, tc := range []struct {
+               name string
+               size int
+       }{{"1K", 1_000}, {"64K", 64_000}, {"1M", 1_000_000}} {
+               size := tc.size
+               values := make([]int64, size)
+               for i := range values {
+                       values[i] = int64(i%1_000_000) - 500_000
+               }
+               rdr := &integerTransferRecordReader{
+                       physicalType: parquet.Types.Int64,
+                       values:       arrow.Int64Traits.CastToBytes(values),
+                       length:       len(values),
+               }
+
+               b.Run("Int64ToDecimal128/"+tc.name, func(b *testing.B) {
+                       b.ReportAllocs()
+                       b.SetBytes(int64(len(values) * arrow.Int64SizeBytes))
+                       b.ResetTimer()
+                       for range b.N {
+                               data := transferDecimalInteger(rdr, 
&arrow.Decimal128Type{Precision: 18})
+                               data.Release()
+                       }
+               })
+
+               b.Run("Int64ToDecimal256/"+tc.name, func(b *testing.B) {
+                       b.ReportAllocs()
+                       b.SetBytes(int64(len(values) * arrow.Int64SizeBytes))
+                       b.ResetTimer()
+                       for range b.N {
+                               data := transferDecimalInteger(rdr, 
&arrow.Decimal256Type{Precision: 18})
+                               data.Release()
+                       }
+               })
+       }
+}
diff --git a/parquet/pqarrow/column_readers_transfer_test.go 
b/parquet/pqarrow/column_readers_transfer_test.go
new file mode 100644
index 00000000..5d063d5d
--- /dev/null
+++ b/parquet/pqarrow/column_readers_transfer_test.go
@@ -0,0 +1,170 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+package pqarrow
+
+import (
+       "testing"
+
+       "github.com/apache/arrow-go/v18/arrow"
+       "github.com/apache/arrow-go/v18/arrow/decimal128"
+       "github.com/apache/arrow-go/v18/arrow/decimal256"
+       "github.com/apache/arrow-go/v18/arrow/memory"
+       "github.com/apache/arrow-go/v18/parquet"
+       "github.com/apache/arrow-go/v18/parquet/file"
+       "github.com/stretchr/testify/require"
+)
+
+type integerTransferRecordReader struct {
+       file.RecordReader
+       physicalType parquet.Type
+       values       []byte
+       length       int
+}
+
+func (r *integerTransferRecordReader) ValuesWritten() int { return r.length }
+func (r *integerTransferRecordReader) Type() parquet.Type { return 
r.physicalType }
+func (r *integerTransferRecordReader) Values() []byte     { return r.values }
+func (r *integerTransferRecordReader) NullCount() int64   { return 0 }
+func (r *integerTransferRecordReader) ReleaseValidBits() *memory.Buffer {
+       return nil
+}
+
+func requireConvertedValues[Out arrowInteger, In parquetInteger](t *testing.T, 
got []Out, values []In) {
+       t.Helper()
+       require.Len(t, got, len(values))
+       for i, value := range values {
+               require.Equal(t, Out(value), got[i])
+       }
+}
+
+func requireIntegerTransfer[In parquetInteger](t *testing.T, rdr 
file.RecordReader, values []In, dt arrow.DataType) {
+       t.Helper()
+       data := transferInt(rdr, dt)
+       defer data.Release()
+
+       buf := data.Buffers()[1].Bytes()
+       switch dt.ID() {
+       case arrow.INT8:
+               requireConvertedValues(t, arrow.Int8Traits.CastFromBytes(buf), 
values)
+       case arrow.UINT8:
+               requireConvertedValues(t, arrow.Uint8Traits.CastFromBytes(buf), 
values)
+       case arrow.INT16:
+               requireConvertedValues(t, arrow.Int16Traits.CastFromBytes(buf), 
values)
+       case arrow.UINT16:
+               requireConvertedValues(t, 
arrow.Uint16Traits.CastFromBytes(buf), values)
+       case arrow.UINT32:
+               requireConvertedValues(t, 
arrow.Uint32Traits.CastFromBytes(buf), values)
+       case arrow.UINT64:
+               requireConvertedValues(t, 
arrow.Uint64Traits.CastFromBytes(buf), values)
+       case arrow.DATE32:
+               requireConvertedValues(t, 
arrow.Date32Traits.CastFromBytes(buf), values)
+       case arrow.TIME32:
+               requireConvertedValues(t, 
arrow.Time32Traits.CastFromBytes(buf), values)
+       case arrow.TIME64:
+               requireConvertedValues(t, 
arrow.Time64Traits.CastFromBytes(buf), values)
+       }
+}
+
+func TestTransferIntegerValues(t *testing.T) {
+       types := []arrow.DataType{
+               arrow.PrimitiveTypes.Int8,
+               arrow.PrimitiveTypes.Uint8,
+               arrow.PrimitiveTypes.Int16,
+               arrow.PrimitiveTypes.Uint16,
+               arrow.PrimitiveTypes.Uint32,
+               arrow.PrimitiveTypes.Uint64,
+               arrow.FixedWidthTypes.Date32,
+               arrow.FixedWidthTypes.Time32ms,
+               arrow.FixedWidthTypes.Time64us,
+       }
+
+       t.Run("int32 physical values", func(t *testing.T) {
+               values := []int32{-1 << 31, -32769, -1, 0, 1, 255, 1<<31 - 1}
+               rdr := &integerTransferRecordReader{
+                       physicalType: parquet.Types.Int32,
+                       values:       arrow.Int32Traits.CastToBytes(values),
+                       length:       len(values),
+               }
+               for _, dt := range types {
+                       t.Run(dt.Name(), func(t *testing.T) {
+                               requireIntegerTransfer(t, rdr, values, dt)
+                       })
+               }
+       })
+
+       t.Run("int64 physical values", func(t *testing.T) {
+               values := []int64{-1 << 63, -1<<32 - 1, -1, 0, 1, 1<<32 - 1, 
1<<63 - 1}
+               rdr := &integerTransferRecordReader{
+                       physicalType: parquet.Types.Int64,
+                       values:       arrow.Int64Traits.CastToBytes(values),
+                       length:       len(values),
+               }
+               for _, dt := range types {
+                       t.Run(dt.Name(), func(t *testing.T) {
+                               requireIntegerTransfer(t, rdr, values, dt)
+                       })
+               }
+       })
+}
+
+func TestTransferDecimalIntegerValues(t *testing.T) {
+       t.Run("int32 physical values", func(t *testing.T) {
+               values := []int32{-999_999_999, -1, 0, 1, 999_999_999}
+               rdr := &integerTransferRecordReader{
+                       physicalType: parquet.Types.Int32,
+                       values:       arrow.Int32Traits.CastToBytes(values),
+                       length:       len(values),
+               }
+
+               data128 := transferDecimalInteger(rdr, 
&arrow.Decimal128Type{Precision: 9})
+               defer data128.Release()
+               got128 := 
arrow.Decimal128Traits.CastFromBytes(data128.Buffers()[1].Bytes())
+               for i, value := range values {
+                       require.Equal(t, decimal128.FromI64(int64(value)), 
got128[i])
+               }
+
+               data256 := transferDecimalInteger(rdr, 
&arrow.Decimal256Type{Precision: 9})
+               defer data256.Release()
+               got256 := 
arrow.Decimal256Traits.CastFromBytes(data256.Buffers()[1].Bytes())
+               for i, value := range values {
+                       require.Equal(t, decimal256.FromI64(int64(value)), 
got256[i])
+               }
+       })
+
+       t.Run("int64 physical values", func(t *testing.T) {
+               values := []int64{-999_999_999_999_999_999, -1, 0, 1, 
999_999_999_999_999_999}
+               rdr := &integerTransferRecordReader{
+                       physicalType: parquet.Types.Int64,
+                       values:       arrow.Int64Traits.CastToBytes(values),
+                       length:       len(values),
+               }
+
+               data128 := transferDecimalInteger(rdr, 
&arrow.Decimal128Type{Precision: 18})
+               defer data128.Release()
+               got128 := 
arrow.Decimal128Traits.CastFromBytes(data128.Buffers()[1].Bytes())
+               for i, value := range values {
+                       require.Equal(t, decimal128.FromI64(value), got128[i])
+               }
+
+               data256 := transferDecimalInteger(rdr, 
&arrow.Decimal256Type{Precision: 18})
+               defer data256.Release()
+               got256 := 
arrow.Decimal256Traits.CastFromBytes(data256.Buffers()[1].Bytes())
+               for i, value := range values {
+                       require.Equal(t, decimal256.FromI64(value), got256[i])
+               }
+       })
+}

Reply via email to