From 02d09df9af1ff29a481ea01e3c51e8efd675479d Mon Sep 17 00:00:00 2001 From: Minh Vu Date: Thu, 6 Aug 2026 00:12:40 +0200 Subject: [PATCH] fix(arrow/array): preserve JSON rows with empty schemas --- arrow/array/json_reader.go | 22 ++++++++++++++++++---- arrow/array/json_reader_test.go | 28 ++++++++++++++++++++++++++++ 2 files changed, 46 insertions(+), 4 deletions(-) diff --git a/arrow/array/json_reader.go b/arrow/array/json_reader.go index 86f3383f0..52938af0b 100644 --- a/arrow/array/json_reader.go +++ b/arrow/array/json_reader.go @@ -184,11 +184,17 @@ func (r *JSONReader) readNext() bool { } func (r *JSONReader) nextall() bool { + n := 0 for r.readNext() { + n++ } - r.cur = r.bldr.NewRecordBatch() - return r.cur.NumRows() > 0 + if r.schema.NumFields() == 0 { + r.cur = NewRecordBatch(r.schema, nil, int64(n)) + } else { + r.cur = r.bldr.NewRecordBatch() + } + return n > 0 } func (r *JSONReader) next1() bool { @@ -196,7 +202,11 @@ func (r *JSONReader) next1() bool { return false } - r.cur = r.bldr.NewRecordBatch() + if r.schema.NumFields() == 0 { + r.cur = NewRecordBatch(r.schema, nil, 1) + } else { + r.cur = r.bldr.NewRecordBatch() + } return true } @@ -210,7 +220,11 @@ func (r *JSONReader) nextn() bool { } if n > 0 { - r.cur = r.bldr.NewRecordBatch() + if r.schema.NumFields() == 0 { + r.cur = NewRecordBatch(r.schema, nil, int64(n)) + } else { + r.cur = r.bldr.NewRecordBatch() + } } return n > 0 } diff --git a/arrow/array/json_reader_test.go b/arrow/array/json_reader_test.go index 3d0def659..4254347e7 100644 --- a/arrow/array/json_reader_test.go +++ b/arrow/array/json_reader_test.go @@ -95,6 +95,34 @@ func TestJSONReaderAll(t *testing.T) { assert.False(t, rdr.Next()) } +func TestJSONReaderPreservesRowsForEmptySchema(t *testing.T) { + schema := arrow.NewSchema(nil, nil) + mem := memory.NewCheckedAllocator(memory.NewGoAllocator()) + defer mem.AssertSize(t, 0) + + t.Run("one row per batch", func(t *testing.T) { + rdr := array.NewJSONReader(strings.NewReader("{} {}"), schema, array.WithAllocator(mem)) + defer rdr.Release() + + assert.True(t, rdr.Next()) + assert.EqualValues(t, 1, rdr.RecordBatch().NumRows()) + assert.True(t, rdr.Next()) + assert.EqualValues(t, 1, rdr.RecordBatch().NumRows()) + assert.False(t, rdr.Next()) + assert.NoError(t, rdr.Err()) + }) + + t.Run("all rows in one batch", func(t *testing.T) { + rdr := array.NewJSONReader(strings.NewReader("{} {}"), schema, array.WithAllocator(mem), array.WithChunk(-1)) + defer rdr.Release() + + assert.True(t, rdr.Next()) + assert.EqualValues(t, 2, rdr.RecordBatch().NumRows()) + assert.False(t, rdr.Next()) + assert.NoError(t, rdr.Err()) + }) +} + func TestJSONReaderChunked(t *testing.T) { schema := arrow.NewSchema([]arrow.Field{ {Name: "region", Type: arrow.BinaryTypes.String, Nullable: true},