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 ab2efb08 fix(arrow/util): return protobuf record conversion errors
(#1088)
ab2efb08 is described below
commit ab2efb08ff9e892d7751571da457672555306aff
Author: Minh Vu <[email protected]>
AuthorDate: Mon Aug 10 18:45:23 2026 +0200
fix(arrow/util): return protobuf record conversion errors (#1088)
### Rationale for this change
Record discards errors from AppendValueOrNull. An unknown protobuf enum
value can then panic or leave builders with mismatched lengths instead
of reporting the conversion failure.
### What changes are included in this PR?
Add RecordWithError with field context, reject undefined enum numbers
safely, and clean up temporary builders and dictionary values on both
success and failure. Keep Record with its existing panic-on-error
behavior for compatibility.
### Are these changes tested?
- `go test ./arrow/util`
- Added checked-allocator coverage for undefined enum values and
conversion cleanup.
### Are there any user-facing changes?
Yes. RecordWithError provides an error-returning path while the existing
panic-on-error Record API remains available.
---
arrow/util/protobuf_reflect.go | 32 +++++++++++++++++++++++------
arrow/util/protobuf_reflect_test.go | 40 +++++++++++++++++++++++++++++++++++++
2 files changed, 66 insertions(+), 6 deletions(-)
diff --git a/arrow/util/protobuf_reflect.go b/arrow/util/protobuf_reflect.go
index ca17c838..e6e01109 100644
--- a/arrow/util/protobuf_reflect.go
+++ b/arrow/util/protobuf_reflect.go
@@ -661,23 +661,35 @@ func (msg ProtobufMessageReflection) Schema()
*arrow.Schema {
// Record returns an arrow.RecordBatch for a protobuf message
func (msg ProtobufMessageReflection) Record(mem memory.Allocator)
arrow.RecordBatch {
+ record, err := msg.RecordWithError(mem)
+ if err != nil {
+ panic(err)
+ }
+ return record
+}
+
+// RecordWithError returns an arrow.RecordBatch for a protobuf message and
+// reports conversion failures.
+func (msg ProtobufMessageReflection) RecordWithError(mem memory.Allocator)
(arrow.RecordBatch, error) {
if mem == nil {
mem = memory.NewGoAllocator()
}
schema := msg.Schema()
if schema.NumFields() == 0 {
- return array.NewRecordBatch(schema, nil, 1)
+ return array.NewRecordBatch(schema, nil, 1), nil
}
recordBuilder := array.NewRecordBuilder(mem, schema)
defer recordBuilder.Release()
for i, f := range msg.fields {
- f.AppendValueOrNull(recordBuilder.Field(i), mem)
+ if err := f.AppendValueOrNull(recordBuilder.Field(i), mem); err
!= nil {
+ return nil, fmt.Errorf("failed to append protobuf field
%q: %w", f.name(), err)
+ }
}
- return recordBuilder.NewRecordBatch()
+ return recordBuilder.NewRecordBatch(), nil
}
// NewProtobufMessageReflection initialises a ProtobufMessageReflection
@@ -776,7 +788,11 @@ func (f ProtobufMessageFieldReflection)
AppendValueOrNull(b array.Builder, mem m
switch b.Type().ID() {
case arrow.STRING:
if f.isEnum() {
-
b.(*array.StringBuilder).Append(string(fd.Enum().Values().ByNumber(pv.Enum()).Name()))
+ enumValue := fd.Enum().Values().ByNumber(pv.Enum())
+ if enumValue == nil {
+ return fmt.Errorf("enum value %d is not defined
for %s", pv.Enum(), fd.Enum().FullName())
+ }
+
b.(*array.StringBuilder).Append(string(enumValue.Name()))
} else {
b.(*array.StringBuilder).Append(pv.String())
}
@@ -816,6 +832,11 @@ func (f ProtobufMessageFieldReflection)
AppendValueOrNull(b array.Builder, mem m
return err
}
case arrow.DICTIONARY:
+ enumNum := int(f.reflectValue().Int())
+ enumValue :=
fd.Enum().Values().ByNumber(protoreflect.EnumNumber(enumNum))
+ if enumValue == nil {
+ return fmt.Errorf("enum value %d is not defined for
%s", enumNum, fd.Enum().FullName())
+ }
pdr := f.asDictionary()
db := b.(*array.BinaryDictionaryBuilder)
dictValues := pdr.getDictValues(mem).(*array.String)
@@ -824,8 +845,7 @@ func (f ProtobufMessageFieldReflection) AppendValueOrNull(b
array.Builder, mem m
if err != nil {
return err
}
- enumNum := int(f.reflectValue().Int())
- enumVal :=
fd.Enum().Values().ByNumber(protoreflect.EnumNumber(enumNum)).Name()
+ enumVal := enumValue.Name()
err = db.AppendValueFromString(string(enumVal))
if err != nil {
return err
diff --git a/arrow/util/protobuf_reflect_test.go
b/arrow/util/protobuf_reflect_test.go
index c93331cd..79cc6e28 100644
--- a/arrow/util/protobuf_reflect_test.go
+++ b/arrow/util/protobuf_reflect_test.go
@@ -390,6 +390,19 @@ func TestRecordReleasesConstructionBuffers(t *testing.T) {
mem.AssertSize(t, 0)
}
+func TestRecordWithErrorReportsFieldConversion(t *testing.T) {
+ msg := AllTheTypesNoAnyFixture().msg.(*util_message.AllTheTypesNoAny)
+ msg.Enum = util_message.AllTheTypesNoAny_ExampleEnum(999)
+ pmr := NewProtobufMessageReflection(msg)
+ mem := memory.NewCheckedAllocator(memory.NewGoAllocator())
+ rec, err := pmr.RecordWithError(mem)
+ assert.Nil(t, rec)
+ require.Error(t, err)
+ assert.ErrorContains(t, err, `failed to append protobuf field "enum"`)
+ assert.ErrorContains(t, err, "enum value 999 is not defined")
+ mem.AssertSize(t, 0)
+}
+
func TestRecordWithAllFieldsExcluded(t *testing.T) {
pmr := NewProtobufMessageReflection(AllTheTypesNoAnyFixture().msg,
WithExclusionPolicy(func(*ProtobufFieldReflection) bool {
return true }))
@@ -400,6 +413,33 @@ func TestRecordWithAllFieldsExcluded(t *testing.T) {
assert.Zero(t, rec.NumCols())
}
+func TestRecordWithErrorReportsUnknownEnumValueForEnumValueHandler(t
*testing.T) {
+ msg := AllTheTypesNoAnyFixture().msg.(*util_message.AllTheTypesNoAny)
+ msg.Enum = util_message.AllTheTypesNoAny_ExampleEnum(999)
+ onlyEnum := func(pfr *ProtobufFieldReflection) bool { return
!pfr.isEnum() }
+ pmr := NewProtobufMessageReflection(msg, WithExclusionPolicy(onlyEnum),
WithEnumHandler(EnumValue))
+ mem := memory.NewCheckedAllocator(memory.NewGoAllocator())
+
+ rec, err := pmr.RecordWithError(mem)
+ assert.Nil(t, rec)
+ require.Error(t, err)
+ assert.ErrorContains(t, err, `failed to append protobuf field "enum"`)
+ assert.ErrorContains(t, err, "enum value 999 is not defined")
+ mem.AssertSize(t, 0)
+}
+
+func TestRecordWithErrorReturnsOneRowForZeroFieldSchema(t *testing.T) {
+ excludeAll := func(*ProtobufFieldReflection) bool { return true }
+ pmr := NewProtobufMessageReflection(&util_message.AllTheTypesNoAny{},
WithExclusionPolicy(excludeAll))
+
+ rec, err := pmr.RecordWithError(memory.DefaultAllocator)
+ require.NoError(t, err)
+ require.NotNil(t, rec)
+ defer rec.Release()
+ require.EqualValues(t, 1, rec.NumRows())
+ require.EqualValues(t, 0, rec.NumCols())
+}
+
func TestRecordFromEmptyMessage(t *testing.T) {
pmr := NewProtobufMessageReflection(&emptypb.Empty{})
rec := pmr.Record(nil)