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

jrmccluskey pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/master by this push:
     new 22abffca3a4 [Go SDK] Encode INT16 row fields as big endian to match 
Java and Python (#40152)
22abffca3a4 is described below

commit 22abffca3a40c81d9ee4161cd5e7657f9f85ad56
Author: Elia Liu <[email protected]>
AuthorDate: Fri Oct 9 02:40:22 2026 +1100

    [Go SDK] Encode INT16 row fields as big endian to match Java and Python 
(#40152)
    
    * [Go SDK] Encode INT16 row fields as big endian to match Java and Python
    
    The reflection based row coder used the varint encoding for int16 and
    uint16 fields. Java and Python encode the INT16 schema type as 2 big
    endian bytes, so rows with such fields could not be exchanged with the
    other SDKs.
    
    Fixes #40151. Part of #39684.
    
    * [Go SDK] Describe the INT16 change as update incompatible in CHANGES.md
---
 CHANGES.md                                       |  1 +
 sdks/go/pkg/beam/core/graph/coder/int.go         | 31 ++++++++++++++++++
 sdks/go/pkg/beam/core/graph/coder/int_test.go    | 32 +++++++++++++++++++
 sdks/go/pkg/beam/core/graph/coder/row_decoder.go | 26 +++++++++++++--
 sdks/go/pkg/beam/core/graph/coder/row_encoder.go | 14 +++++++--
 sdks/go/pkg/beam/core/graph/coder/row_test.go    | 40 ++++++++++++++++++++++++
 6 files changed, 140 insertions(+), 4 deletions(-)

diff --git a/CHANGES.md b/CHANGES.md
index 997ee9adb49..d87eb6181c5 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -77,6 +77,7 @@
 ## Breaking Changes
 
 * X behavior was changed ([#X](https://github.com/apache/beam/issues/X)).
+* (Go) The row coder now encodes `int16` and `uint16` struct fields as 2 byte 
big endian INT16 values, matching the Java and Python SDKs. This is an update 
incompatible change for streaming pipelines that use rows with `int16` or 
`uint16` fields ([#40151](https://github.com/apache/beam/issues/40151)).
 
 ## Deprecations
 
diff --git a/sdks/go/pkg/beam/core/graph/coder/int.go 
b/sdks/go/pkg/beam/core/graph/coder/int.go
index 6eda98ccc38..8ef56b015ac 100644
--- a/sdks/go/pkg/beam/core/graph/coder/int.go
+++ b/sdks/go/pkg/beam/core/graph/coder/int.go
@@ -69,3 +69,34 @@ func DecodeInt32(r io.Reader) (int32, error) {
        }
        return int32(ret), nil
 }
+
+// EncodeUint16 encodes an uint16 in big endian format.
+func EncodeUint16(value uint16, w io.Writer) error {
+       var data [2]byte
+       binary.BigEndian.PutUint16(data[:], value)
+       _, err := ioutilx.WriteUnsafe(w, data[:])
+       return err
+}
+
+// DecodeUint16 decodes an uint16 in big endian format.
+func DecodeUint16(r io.Reader) (uint16, error) {
+       var data [2]byte
+       if err := ioutilx.ReadNBufUnsafe(r, data[:]); err != nil {
+               return 0, err
+       }
+       return binary.BigEndian.Uint16(data[:]), nil
+}
+
+// EncodeInt16 encodes an int16 in big endian format.
+func EncodeInt16(value int16, w io.Writer) error {
+       return EncodeUint16(uint16(value), w)
+}
+
+// DecodeInt16 decodes an int16 in big endian format.
+func DecodeInt16(r io.Reader) (int16, error) {
+       ret, err := DecodeUint16(r)
+       if err != nil {
+               return 0, err
+       }
+       return int16(ret), nil
+}
diff --git a/sdks/go/pkg/beam/core/graph/coder/int_test.go 
b/sdks/go/pkg/beam/core/graph/coder/int_test.go
index 25dc681710d..737340f1663 100644
--- a/sdks/go/pkg/beam/core/graph/coder/int_test.go
+++ b/sdks/go/pkg/beam/core/graph/coder/int_test.go
@@ -86,3 +86,35 @@ func TestEncodeDecodeInt32(t *testing.T) {
                }
        }
 }
+
+func TestEncodeDecodeInt16(t *testing.T) {
+       tests := []struct {
+               value int16
+               want  []byte
+       }{
+               {-32768, []byte{0x80, 0x00}},
+               {-2, []byte{0xff, 0xfe}},
+               {-1, []byte{0xff, 0xff}},
+               {0, []byte{0x00, 0x00}},
+               {1, []byte{0x00, 0x01}},
+               {999, []byte{0x03, 0xe7}},
+               {32767, []byte{0x7f, 0xff}},
+       }
+
+       for _, test := range tests {
+               var buf bytes.Buffer
+               if err := EncodeInt16(test.value, &buf); err != nil {
+                       t.Fatalf("EncodeInt16(%v) failed: %v", test.value, err)
+               }
+               if got := buf.Bytes(); !bytes.Equal(got, test.want) {
+                       t.Errorf("EncodeInt16(%v) = %v, want %v", test.value, 
got, test.want)
+               }
+               actual, err := DecodeInt16(&buf)
+               if err != nil {
+                       t.Fatalf("DecodeInt16(<%v>) failed: %v", test.value, 
err)
+               }
+               if actual != test.value {
+                       t.Errorf("DecodeInt16(<%v>) = %v, want %v", test.value, 
actual, test.value)
+               }
+       }
+}
diff --git a/sdks/go/pkg/beam/core/graph/coder/row_decoder.go 
b/sdks/go/pkg/beam/core/graph/coder/row_decoder.go
index 8dabc2be9ff..90938c237fe 100644
--- a/sdks/go/pkg/beam/core/graph/coder/row_decoder.go
+++ b/sdks/go/pkg/beam/core/graph/coder/row_decoder.go
@@ -241,6 +241,24 @@ func reflectDecodeUint(rv reflect.Value, r io.Reader) 
error {
        return nil
 }
 
+func reflectDecodeInt16(rv reflect.Value, r io.Reader) error {
+       v, err := DecodeInt16(r)
+       if err != nil {
+               return errors.Wrap(err, "error decoding int16 field")
+       }
+       rv.SetInt(int64(v))
+       return nil
+}
+
+func reflectDecodeUint16(rv reflect.Value, r io.Reader) error {
+       v, err := DecodeUint16(r)
+       if err != nil {
+               return errors.Wrap(err, "error decoding uint16 field")
+       }
+       rv.SetUint(uint64(v))
+       return nil
+}
+
 func reflectDecodeSinglePrecisionFloat(rv reflect.Value, r io.Reader) error {
        v, err := DecodeSinglePrecisionFloat(r)
        if err != nil {
@@ -341,10 +359,14 @@ func (b *RowDecoderBuilder) decoderForSingleTypeReflect(t 
reflect.Type) (typeDec
                return typeDecoderFieldReflect{decode: reflectDecodeByte}, nil
        case reflect.String:
                return typeDecoderFieldReflect{decode: reflectDecodeString}, nil
-       case reflect.Int, reflect.Int8, reflect.Int16, reflect.Int32, 
reflect.Int64:
+       case reflect.Int, reflect.Int8, reflect.Int32, reflect.Int64:
                return typeDecoderFieldReflect{decode: reflectDecodeInt}, nil
-       case reflect.Uint, reflect.Uint64, reflect.Uint32, reflect.Uint16:
+       case reflect.Int16:
+               return typeDecoderFieldReflect{decode: reflectDecodeInt16}, nil
+       case reflect.Uint, reflect.Uint64, reflect.Uint32:
                return typeDecoderFieldReflect{decode: reflectDecodeUint}, nil
+       case reflect.Uint16:
+               return typeDecoderFieldReflect{decode: reflectDecodeUint16}, nil
        case reflect.Float32:
                return typeDecoderFieldReflect{decode: 
reflectDecodeSinglePrecisionFloat}, nil
        case reflect.Float64:
diff --git a/sdks/go/pkg/beam/core/graph/coder/row_encoder.go 
b/sdks/go/pkg/beam/core/graph/coder/row_encoder.go
index 2756c9b8e70..6f5c4c04e81 100644
--- a/sdks/go/pkg/beam/core/graph/coder/row_encoder.go
+++ b/sdks/go/pkg/beam/core/graph/coder/row_encoder.go
@@ -200,14 +200,24 @@ func (b *RowEncoderBuilder) encoderForSingleTypeReflect(t 
reflect.Type) (typeEnc
                return typeEncoderFieldReflect{encode: func(rv reflect.Value, w 
io.Writer) error {
                        return EncodeStringUTF8(rv.String(), w)
                }}, nil
-       case reflect.Int, reflect.Int64, reflect.Int16, reflect.Int32, 
reflect.Int8:
+       case reflect.Int, reflect.Int64, reflect.Int32, reflect.Int8:
                return typeEncoderFieldReflect{encode: func(rv reflect.Value, w 
io.Writer) error {
                        return EncodeVarInt(rv.Int(), w)
                }}, nil
-       case reflect.Uint, reflect.Uint64, reflect.Uint32, reflect.Uint16:
+       case reflect.Int16:
+               // INT16 schema fields are 2 big endian bytes.
+               return typeEncoderFieldReflect{encode: func(rv reflect.Value, w 
io.Writer) error {
+                       return EncodeInt16(int16(rv.Int()), w)
+               }}, nil
+       case reflect.Uint, reflect.Uint64, reflect.Uint32:
                return typeEncoderFieldReflect{encode: func(rv reflect.Value, w 
io.Writer) error {
                        return EncodeVarUint64(rv.Uint(), w)
                }}, nil
+       case reflect.Uint16:
+               // uint16 is stored as an INT16 schema field.
+               return typeEncoderFieldReflect{encode: func(rv reflect.Value, w 
io.Writer) error {
+                       return EncodeUint16(uint16(rv.Uint()), w)
+               }}, nil
        case reflect.Float32:
                return typeEncoderFieldReflect{encode: func(rv reflect.Value, w 
io.Writer) error {
                        return EncodeSinglePrecisionFloat(float32(rv.Float()), 
w)
diff --git a/sdks/go/pkg/beam/core/graph/coder/row_test.go 
b/sdks/go/pkg/beam/core/graph/coder/row_test.go
index 1549123558b..c2a925efd30 100644
--- a/sdks/go/pkg/beam/core/graph/coder/row_test.go
+++ b/sdks/go/pkg/beam/core/graph/coder/row_test.go
@@ -905,3 +905,43 @@ func TestRowHeader_TrailingZeroBytes(t *testing.T) {
        }
 
 }
+
+// TestReflectionRowCoderGeneration_Int16BigEndian checks that int16 and uint16
+// fields use the 2 byte big endian encoding of the INT16 schema type.
+func TestReflectionRowCoderGeneration_Int16BigEndian(t *testing.T) {
+       type row struct {
+               A int16
+               B uint16
+               C int16
+       }
+       want := row{A: 999, B: 65535, C: -2}
+       wantBytes := []byte{
+               0x03, 0x00, // 3 fields, no nils.
+               0x03, 0xe7, // A
+               0xff, 0xff, // B
+               0xff, 0xfe, // C
+       }
+       rt := reflect.TypeOf(want)
+       enc, err := RowEncoderForStruct(rt)
+       if err != nil {
+               t.Fatalf("RowEncoderForStruct(%v) = %v, want nil error", rt, 
err)
+       }
+       var buf bytes.Buffer
+       if err := enc(want, &buf); err != nil {
+               t.Fatalf("enc(%v) = %v, want nil error", want, err)
+       }
+       if got := buf.Bytes(); !bytes.Equal(got, wantBytes) {
+               t.Fatalf("enc(%v) = %v, want %v", want, got, wantBytes)
+       }
+       dec, err := RowDecoderForStruct(rt)
+       if err != nil {
+               t.Fatalf("RowDecoderForStruct(%v) = %v, want nil error", rt, 
err)
+       }
+       got, err := dec(bytes.NewBuffer(wantBytes))
+       if err != nil {
+               t.Fatalf("dec(%v) = %v, want nil error", wantBytes, err)
+       }
+       if d := cmp.Diff(want, got); d != "" {
+               t.Fatalf("dec(enc(%v)) = %v\ndiff (-want, +got): %v", want, 
got, d)
+       }
+}

Reply via email to