Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@ require (
github.com/testcontainers/testcontainers-go/modules/compose v0.42.0
github.com/twmb/avro v1.7.2
github.com/twmb/murmur3 v1.1.8
github.com/twpayne/go-geom v1.6.1
github.com/uptrace/bun v1.2.18
github.com/uptrace/bun/dialect/mssqldialect v1.2.18
github.com/uptrace/bun/dialect/mysqldialect v1.2.18
Expand Down
10 changes: 10 additions & 0 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,8 @@ github.com/AzureAD/microsoft-authentication-extensions-for-go/cache v0.1.1 h1:WJ
github.com/AzureAD/microsoft-authentication-extensions-for-go/cache v0.1.1/go.mod h1:tCcJZ0uHAmvjsVYzEFivsRTN00oz5BEsRgQHu5JZ9WE=
github.com/AzureAD/microsoft-authentication-library-for-go v1.7.2 h1:RHK7bS+HQMslb1sZpAokUt+zTVmue0hKSs2C791hhzU=
github.com/AzureAD/microsoft-authentication-library-for-go v1.7.2/go.mod h1:HKpQxkWaGLJ+D/5H8QRpyQXA1eKjxkFlOMwck5+33Jk=
github.com/DATA-DOG/go-sqlmock v1.5.2 h1:OcvFkGmslmlZibjAjaHm3L//6LiuBgolP7OputlJIzU=
github.com/DATA-DOG/go-sqlmock v1.5.2/go.mod h1:88MAG/4G7SMwSE3CeA0ZKzrT5CiOU3OJ+JlNzwDqpNU=
github.com/DefangLabs/secret-detector v0.0.0-20250403165618-22662109213e h1:rd4bOvKmDIx0WeTv9Qz+hghsgyjikFiPrseXHlKepO0=
github.com/DefangLabs/secret-detector v0.0.0-20250403165618-22662109213e/go.mod h1:blbwPQh4DTlCZEfk1BLU4oMIhLda2U+A840Uag9DsZw=
github.com/GoogleCloudPlatform/opentelemetry-operations-go/detectors/gcp v1.31.0 h1:DHa2U07rk8syqvCge0QIGMCE1WxGj9njT44GH7zNJLQ=
Expand Down Expand Up @@ -85,6 +87,10 @@ github.com/RoaringBitmap/roaring/v2 v2.18.2 h1:oPq3Cgx//iDuJQVp6xSInAKW34J9CEwE5
github.com/RoaringBitmap/roaring/v2 v2.18.2/go.mod h1:eq4wdNXxtJIS/oikeCzdX1rBzek7ANzbth041hrU8Q4=
github.com/acarl005/stripansi v0.0.0-20180116102854-5a71ef0e047d h1:licZJFw2RwpHMqeKTCYkitsPqHNxTmd4SNR5r94FGM8=
github.com/acarl005/stripansi v0.0.0-20180116102854-5a71ef0e047d/go.mod h1:asat636LX7Bqt5lYEZ27JNDcqxfjdBQuJ/MM4CN/Lzo=
github.com/alecthomas/assert/v2 v2.10.0 h1:jjRCHsj6hBJhkmhznrCzoNpbA3zqy0fYiUcYZP/GkPY=
github.com/alecthomas/assert/v2 v2.10.0/go.mod h1:Bze95FyfUr7x34QZrjL+XP+0qgp/zg8yS+TtBj1WA3k=
github.com/alecthomas/repr v0.4.0 h1:GhI2A8MACjfegCPVq9f1FLvIBS+DrQ2KQBFZP1iFzXc=
github.com/alecthomas/repr v0.4.0/go.mod h1:Fr0507jx4eOXV7AlPV6AVZLYrLIuIeSOWtW57eE/O/4=
github.com/alexflint/go-arg v1.6.1 h1:uZogJ6VDBjcuosydKgvYYRhh9sRCusjOvoOLZopBlnA=
github.com/alexflint/go-arg v1.6.1/go.mod h1:nQ0LFYftLJ6njcaee0sU+G0iS2+2XJQfA8I062D0LGc=
github.com/alexflint/go-scalar v1.2.0 h1:WR7JPKkeNpnYIOfHRa7ivM21aWAdHD0gEWHCx+WQBRw=
Expand Down Expand Up @@ -392,6 +398,8 @@ github.com/hashicorp/go-version v1.9.0 h1:CeOIz6k+LoN3qX9Z0tyQrPtiB1DFYRPfCIBtaX
github.com/hashicorp/go-version v1.9.0/go.mod h1:fltr4n8CU8Ke44wwGCBoEymUuxUHl09ZGVZPK5anwXA=
github.com/hashicorp/golang-lru/v2 v2.0.7 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs4luLUK2k=
github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM=
github.com/hexops/gotextdiff v1.0.3 h1:gitA9+qJrrTCsiCl7+kh75nPqQt1cx4ZkudSTLoUqJM=
github.com/hexops/gotextdiff v1.0.3/go.mod h1:pSWU5MAI3yDq+fZBTazCSJysOMbxWL1BSow5/V2vxeg=
github.com/in-toto/attestation v1.1.2 h1:MBFn6lsMq6dptQZJBhalXTcWMb/aJy3V+GX3VYj/V1E=
github.com/in-toto/attestation v1.1.2/go.mod h1:gYFddHMZj3DiQ0b62ltNi1Vj5rC879bTmBbrv9CRHpM=
github.com/in-toto/in-toto-golang v0.11.0 h1:nfidMYBFx+E0lnmX5KUnN2Pdm8zdNKal1ayjJuzzRoA=
Expand Down Expand Up @@ -624,6 +632,8 @@ github.com/twmb/avro v1.7.2 h1:cmrEBRSbELRqsg/dRkQvVWuOaR2EfGifHIt/2iJ9lfI=
github.com/twmb/avro v1.7.2/go.mod h1:X0fT1dY2xcbV4YuCE4mYro+qljHl4kUF5uA/2z1rgSk=
github.com/twmb/murmur3 v1.1.8 h1:8Yt9taO/WN3l08xErzjeschgZU2QSrwm1kclYq+0aRg=
github.com/twmb/murmur3 v1.1.8/go.mod h1:Qq/R7NUyOfr65zD+6Q5IHKsJLwP7exErjN6lyyq3OSQ=
github.com/twpayne/go-geom v1.6.1 h1:iLE+Opv0Ihm/ABIcvQFGIiFBXd76oBIar9drAwHFhR4=
github.com/twpayne/go-geom v1.6.1/go.mod h1:Kr+Nly6BswFsKM5sd31YaoWS5PeDDH2NftJTK7Gd028=
github.com/uptrace/bun v1.2.18 h1:3HnRcMfS6OBPMG1eSOzlbFJ/X/AyMEJb7rMxE6VQvDU=
github.com/uptrace/bun v1.2.18/go.mod h1:wNltaKJk4JtOt4SG5I5zmA7v0/Mzjh1+/S906Rayd3Y=
github.com/uptrace/bun/dialect/mssqldialect v1.2.18 h1:nYzHoyJKJlIyl5i95Exi8ZTK8ooKWG+o3z3f404d/yQ=
Expand Down
179 changes: 173 additions & 6 deletions table/arrow_utils.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,12 +18,15 @@
package table

import (
"bytes"
"context"
"encoding/base64"
"encoding/json"
"errors"
"fmt"
"iter"
"log/slog"
"regexp"
"slices"
"strconv"
"strings"
Expand Down Expand Up @@ -375,6 +378,17 @@ func (c convertToIceberg) Primitive(dt arrow.DataType) (result iceberg.NestedFie
result.Type = iceberg.PrimitiveTypes.UUID
case "parquet.variant":
result.Type = iceberg.VariantType{}
case "geoarrow.wkb":
wkb, ok := dt.(*geoarrow.WKBType)
if !ok {
panic(fmt.Errorf("%w: unsupported arrow type for conversion - %s", iceberg.ErrInvalidSchema, dt))
}

iceType, err := geoArrowMetadataToIcebergType(wkb.Metadata())
if err != nil {
panic(fmt.Errorf("%w: converting geoarrow metadata: %w", iceberg.ErrInvalidSchema, err))
}
result.Type = iceType
default:
panic(fmt.Errorf("%w: unsupported arrow type for conversion - %s", iceberg.ErrInvalidSchema, dt))
}
Expand Down Expand Up @@ -658,20 +672,26 @@ func (c convertToArrow) VisitVariant() arrow.Field {
return arrow.Field{Type: extensions.NewDefaultVariantType()}
}

func (c convertToArrow) VisitGeometry(iceberg.GeometryType) arrow.Field {
func (c convertToArrow) VisitGeometry(g iceberg.GeometryType) arrow.Field {
meta := icebergCRSToGeoArrowMetadata(g.CRS())
if c.useLargeTypes {
return arrow.Field{Type: geoarrow.NewWKBType(geoarrow.WKBWithLargeBinaryStorage())}
return arrow.Field{Type: geoarrow.NewWKBType(geoarrow.WKBWithLargeBinaryStorage(), geoarrow.WKBWithMetadata(meta))}
}

return arrow.Field{Type: geoarrow.NewWKBType(geoarrow.WKBWithBinaryStorage())}
return arrow.Field{Type: geoarrow.NewWKBType(geoarrow.WKBWithBinaryStorage(), geoarrow.WKBWithMetadata(meta))}
}

Comment thread
happydave1 marked this conversation as resolved.
func (c convertToArrow) VisitGeography(iceberg.GeographyType) arrow.Field {
func (c convertToArrow) VisitGeography(g iceberg.GeographyType) arrow.Field {
meta := icebergCRSToGeoArrowMetadata(g.CRS())
// Always add an edge to differentiate between Geography and Geometry arrow fields.
// Note that the edge convention is a best-effort hint and planar geography from other clients won't round-trip through Arrow alone.
Comment on lines +686 to +687

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I think this comment is misleading (planar geography is not a thing, and this is not a hint, it affects correctness).

Suggested change
// Always add an edge to differentiate between Geography and Geometry arrow fields.
// Note that the edge convention is a best-effort hint and planar geography from other clients won't round-trip through Arrow alone.
// Always add an edge to differentiate between Geography and Geometry arrow fields.

meta.Edges = geoarrow.EdgeInterpolation(g.Algorithm())
Comment thread
happydave1 marked this conversation as resolved.

if c.useLargeTypes {
return arrow.Field{Type: geoarrow.NewWKBType(geoarrow.WKBWithLargeBinaryStorage())}
return arrow.Field{Type: geoarrow.NewWKBType(geoarrow.WKBWithLargeBinaryStorage(), geoarrow.WKBWithMetadata(meta))}
}

return arrow.Field{Type: geoarrow.NewWKBType(geoarrow.WKBWithBinaryStorage())}
return arrow.Field{Type: geoarrow.NewWKBType(geoarrow.WKBWithBinaryStorage(), geoarrow.WKBWithMetadata(meta))}
}

var _ iceberg.SchemaVisitorPerPrimitiveType[arrow.Field] = convertToArrow{}
Expand Down Expand Up @@ -1826,3 +1846,150 @@ func positionDeleteRecordsToDataFilesDV(ctx context.Context, rootLocation string
}
}
}

func checkCRSString(rawCrs json.RawMessage) bool {
b := bytes.TrimSpace(rawCrs)

return len(b) > 0 && b[0] == '"'
}

func checkCRSJSON(rawCrs json.RawMessage) bool {
b := bytes.TrimSpace(rawCrs)

return len(b) > 0 && b[0] == '{'
}

func geoArrowCRSToIcebergCRS(meta geoarrow.Metadata) (string, error) {
if len(meta.CRS) == 0 {
return "srid:0", nil

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A bare ga.wkb() field with no CRS — the PyIceberg default and what older clients emit — reads back as geometry("srid:0") rather than geometry(), which isn't Equals to a default-CRS geometry.

Separately, an empty meta.CRS with a non-empty CRSType skips this early return and falls into the switch, so CRSTypeSRID yields "srid:" with an empty id and GeometryTypeOf accepts it silently. I'd return the OGC:CRS84 default for the bare case and guard the empty-CRS-with-CRSType combination, with a test for bare geoarrow.wkb.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

bare ga.wkb() field with no CRS — the PyIceberg default and what older clients emit — reads back as geometry("srid:0") rather than geometry(), which isn't Equals to a default-CRS geometry.

This is the correct behaviour: the GeoArrow default does not equal the Iceberg default. PyIceberg is probably wrong here.

Separately, an empty meta.CRS with a non-empty CRSType skips this early return

This should be fixed...the CRSType can actually just be ignored for the purposes of this function (it's purely a hint)

}

switch {
case checkCRSString(meta.CRS):

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A small cleanup cluster in this function:

The if len(meta.CRS) > 0 wrappers inside both this checkCRSString arm and the checkCRSSJSON arm below are dead — the guard at the top of the function already returns on empty, and the check* funcs only return true for non-empty input (staticcheck SA9003). I'd drop both wrappers.

On the write path (line 1984), strings.EqualFold(lowerCRS, "EPSG:4326") runs a case-fold against an already-lowercased string — lowerCRS == "epsg:4326" is enough.

And in the JSON-object code parsing (line 1923), || codeNum.String() == "" is unreachable — a successful json.Number unmarshal is never empty. I'd drop just that clause; the trailing if code == "" is reachable for {"code":""}, so leave that one.

var crs string

if err := json.Unmarshal(meta.CRS, &crs); err != nil {
return "", fmt.Errorf("invalid geoarrow CRS metadata: %w", err)
}

if strings.EqualFold(crs, "OGC:CRS84") || strings.EqualFold(crs, "EPSG:4326") {
return "OGC:CRS84", nil
}

switch meta.CRSType {
case geoarrow.CRSTypeSRID:
return "srid:" + crs, nil
case geoarrow.CRSTypeWKT22019:
return "", errors.New("CRS type wkt2:2019 not supported")
default:
if len(crs) <= 32 {
return crs, nil
}

return "", errors.New("crs length too long")
}
case checkCRSJSON(meta.CRS):
var crs map[string]json.RawMessage

if err := json.Unmarshal(meta.CRS, &crs); err != nil {
return "", fmt.Errorf("invalid geoarrow CRS metadata: %w", err)
}

idRaw, ok := crs["id"]
if !ok || len(idRaw) == 0 {
return "", errors.New("unsupported CRS")
}

var id map[string]json.RawMessage
if err := json.Unmarshal(idRaw, &id); err != nil {
return "", errors.New("unsupported CRS")
}

var authority string
if err := json.Unmarshal(id["authority"], &authority); err != nil || authority == "" {
return "", errors.New("unsupported CRS")
}

codeRaw := id["code"]
if len(codeRaw) == 0 {
return "", errors.New("unsupported CRS")
}

var code string
if err := json.Unmarshal(codeRaw, &code); err != nil {
var codeNum json.Number
if err := json.Unmarshal(codeRaw, &codeNum); err != nil {
return "", errors.New("unsupported CRS")
}
code = codeNum.String()
}
if code == "" {
return "", errors.New("unsupported CRS")
}

authorityCode := authority + ":" + code
if strings.EqualFold(authorityCode, "OGC:CRS84") || strings.EqualFold(authorityCode, "EPSG:4326") {
return "OGC:CRS84", nil
}

return authorityCode, nil
default:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

CRSTypeAuthorityCode and CRSTypeWKT22019 both fall into default and get returned as-is. For authority_code that's almost right but loses the type tag; for WKT2:2019 it hands a multi-KB WKT2 blob straight back as the Iceberg CRS string, which isn't a valid CRS identifier.

I'd add an explicit CRSTypeAuthorityCode case, and for CRSTypeWKT22019 either error like projjson does or reduce it to the authority code — whichever we pick, a test for each so the behavior is pinned.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The CRSType, if it's coming from JSON, is just a hint, is not required, and is often absent. Here you probably want to:

  • Check if crs is a JSON object. If it is, check the id member and paste together crs["id"]["authority"], :, and crs["id"]["code"]. If any of those are missing, error for an unsupported CRS>
  • Check if crs is a string. If it's shorter than 32 characters, let it through verbatim. There's no official restriction on allowed characters in authorities or codes but the length check should reject anything questionable.

return "", errors.New("unsupported CRS: CRS must either be omitted, a string, or a JSON object")
}
}

func geoArrowMetadataToIcebergType(meta geoarrow.Metadata) (iceberg.Type, error) {
crs, err := geoArrowCRSToIcebergCRS(meta)
if err != nil {
return nil, err
}

switch meta.Edges {
case geoarrow.EdgePlanar:
return iceberg.GeometryTypeOf(crs)
case geoarrow.EdgeVincenty, geoarrow.EdgeKarney, geoarrow.EdgeThomas, geoarrow.EdgeAndoyer, geoarrow.EdgeSpherical:
return iceberg.GeographyTypeOf(crs, string(meta.Edges))
default:
return nil, fmt.Errorf("unsupported geoarrow edges %q", meta.Edges)
}
}

var authorityCodeCRS = regexp.MustCompile(`^[A-Za-z][A-Za-z0-9_-]*:[A-Za-z0-9_.-]+$`)

func icebergCRSToGeoArrowMetadata(crs string) geoarrow.Metadata {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Small thing while we're here: strings.ToLower(crs) runs twice (here and for the projjson check) and the slice crs[len("srid:"):] indexes the original-case string after a lowercased prefix check. I'd hoist lower := strings.ToLower(crs) once and reuse it. Also worth noting the EPSG:4326 match in the read function is case-sensitive, so epsg:4326 falls through — strings.EqualFold would close that.

lowerCRS := strings.ToLower(crs)
if strings.HasPrefix(lowerCRS, "srid:") {
id := crs[len("srid:"):]

if id == "0" {
return geoarrow.NewMetadata() // srid:0 maps to omitted GeoArrow CRS
}

raw, _ := json.Marshal(id) //nolint:errcheck // Marshalling a string can't fail

return geoarrow.Metadata{
CRS: raw,
CRSType: geoarrow.CRSTypeSRID,
}
}

if strings.HasPrefix(lowerCRS, "projjson:") {
panic(fmt.Errorf("%w: projjson CRS not supported yet", iceberg.ErrInvalidSchema))
}

var raw []byte

if lowerCRS == "epsg:4326" {
// collapse EPSG:4326 to OGC:CRS84
raw, _ = json.Marshal("OGC:CRS84") //nolint:errcheck // Marshalling a string can't fail
} else {
raw, _ = json.Marshal(crs) //nolint:errcheck // Marshalling a string can't fail
}

meta := geoarrow.Metadata{CRS: raw}
if authorityCodeCRS.MatchString(crs) {
meta.CRSType = geoarrow.CRSTypeAuthorityCode
}

return meta
}
Loading
Loading