laskoviymishka commented on code in PR #2125:
URL: https://github.com/apache/iceberg-go/pull/2125#discussion_r4231823641
##########
go.mod:
##########
@@ -38,7 +38,7 @@ require (
github.com/beltran/gohive v1.8.1
github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc
github.com/docker/docker v28.5.2+incompatible
- github.com/geoarrow/geoarrow-go v0.0.0-20260403143023-f54751c3e3a1
+ github.com/geoarrow/geoarrow-go v0.0.0-20261005150217-fc2b33c3141d
Review Comment:
Still the untagged pseudo-version from round 1. Since this is the dep that
sets the on-disk Parquet schema and the read-side type mapping, I'd rather pin
a tagged geoarrow-go release before we merge than an untagged commit. If we
have a reason to go in on this commit, a one-line note on why plus a follow-up
to pin a tag works for me.
##########
table/internal/parquet_files.go:
##########
@@ -1188,25 +1196,31 @@ func normalizeWKBArrayReachable(ext
array.ExtensionArray, active []bool, mem mem
}
// accumulateGeoBounds extends the per-field bounding boxes with the WKB values
-// in this batch. Null rows are skipped; a malformed WKB value fails the write.
+// in this batch and tallies their null counts. Null rows are skipped; a
+// malformed WKB value fails the write.
+//
+// Every geoCols entry comes from a *geoarrow.WKBType field of the writer's
+// schema, so a batch whose column does not match is an error rather than a
+// skip: skipping would leave a partial null count that under-reports nulls.
func (w *ParquetFileWriter) accumulateGeoBounds(batch arrow.RecordBatch) error
{
for _, gc := range w.geoCols {
Review Comment:
This guard change is the actual fix for last round's false 0, but nothing
exercises the three new error returns. Swapping them back to `continue` would
resurface the under-counted-nulls regression and every test would still pass. A
table test feeding `Write` a mismatched batch (missing column, non-extension
array, non-`wkbStorage`) and asserting the returned error would pin it.
While we're here, the message only prints the field ID. These are
unreachable for well-formed input, so the error is the only diagnostic when one
does fire, which makes adding the Arrow field name worth it.
##########
table/geo_write_test.go:
##########
@@ -296,3 +296,70 @@ func TestWriteGeometryColumnCheckedAllocator(t *testing.T)
{
require.Contains(t, df.LowerBoundValues(), geoTestGeomFieldID,
"geometry column must record a lower bound")
require.Contains(t, df.UpperBoundValues(), geoTestGeomFieldID,
"geometry column must record an upper bound")
}
+
+// TestWriteGeoColumnMultiRowGroupStats writes geo columns across several row
+// groups and batches. Parquet GEOMETRY/GEOGRAPHY column chunks omit the
+// standard Statistics block (min, max, null count), so DataFileStatsFromMeta
+// invalidates the column in the first row group; value counts and column sizes
+// (which live outside that block) must still sum over every row group, and the
+// null count (tallied from the Arrow data) must sum over every batch. A value
+// count from row group 0 only, paired with the whole-file null count, would
make
+// the file look all-null and let NotNull/IsNull evaluators drop live rows.
+func TestWriteGeoColumnMultiRowGroupStats(t *testing.T) {
+ t.Parallel()
+
+ writer, schema, arrowSchema := newGeoTestWriter(t, t.TempDir(),
iceberg.Properties{
+ tblutils.ParquetRowGroupLimitKey: "2",
+ })
+
+ pt := wktToWKB(t, "POINT (1 2)").String()
+ // Nulls fill row group 0 only. The second batch has none, so a null
tally
+ // that overwrote rather than summed across batches would record 0.
+ first, _, err := array.RecordFromJSON(memory.DefaultAllocator,
arrowSchema, strings.NewReader(`[
+ {"id": 1, "geom": null, "geog": null},
+ {"id": 2, "geom": null, "geog": null},
+ {"id": 3, "geom": "`+pt+`", "geog": "`+pt+`"}
+ ]`))
+ require.NoError(t, err)
+ defer first.Release()
+ second, _, err := array.RecordFromJSON(memory.DefaultAllocator,
arrowSchema, strings.NewReader(`[
+ {"id": 4, "geom": "`+pt+`", "geog": "`+pt+`"},
+ {"id": 5, "geom": "`+pt+`", "geog": "`+pt+`"},
+ {"id": 6, "geom": "`+pt+`", "geog": "`+pt+`"}
+ ]`))
+ require.NoError(t, err)
+ defer second.Release()
+
+ df, err := writer.writeFile(t.Context(), nil, WriteTask{
+ Uuid: uuid.New(),
+ ID: 0,
+ FileCount: 1,
+ Schema: schema,
+ Batches: []arrow.RecordBatch{first, second},
+ })
+ require.NoError(t, err)
+ require.EqualValues(t, 6, df.Count())
+ require.Len(t, df.SplitOffsets(), 3, "row group limit of 2 must produce
3 row groups")
+
+ for _, fieldID := range []int{geoTestGeomFieldID, geoTestGeogFieldID} {
+ require.Contains(t, df.ValueCounts(), fieldID)
+ assert.EqualValues(t, 6, df.ValueCounts()[fieldID], "value
count must cover every row group")
+ require.Contains(t, df.NullValueCounts(), fieldID)
+ assert.EqualValues(t, 2, df.NullValueCounts()[fieldID], "null
count must cover every batch")
+ require.Contains(t, df.ColumnSizes(), fieldID)
Review Comment:
This asserts the value-count and null-count sums but only `Contains` for
`ColumnSizes`, so a regression recording just row group 0's size would still
pass here. Worth asserting `ColumnSizes[fieldID] > 0` at least, or comparing
against the summed per-row-group compressed size, so the multi-row-group
guarantee is actually pinned.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]