diff --git a/sdks/go/pkg/beam/core/runtime/xlangx/expand.go b/sdks/go/pkg/beam/core/runtime/xlangx/expand.go index 0ce3fa495420..56b2c0b6a814 100644 --- a/sdks/go/pkg/beam/core/runtime/xlangx/expand.go +++ b/sdks/go/pkg/beam/core/runtime/xlangx/expand.go @@ -90,6 +90,10 @@ func expand( ext *graph.ExternalTransform) (*jobpb.ExpansionResponse, error) { h, config := defaultReg.getHandlerFunc(transform.GetSpec().GetUrn(), ext.ExpansionAddr) + // Overwrite expansion address if changed due to override for service or URN. + if config != ext.ExpansionAddr { + ext.ExpansionAddr = config + } return h(ctx, &HandlerParams{ Config: config, Req: &jobpb.ExpansionRequest{ diff --git a/sdks/go/pkg/beam/transforms/sql/sql.go b/sdks/go/pkg/beam/transforms/sql/sql.go index cff87133250f..86f1eb13e5dc 100644 --- a/sdks/go/pkg/beam/transforms/sql/sql.go +++ b/sdks/go/pkg/beam/transforms/sql/sql.go @@ -108,7 +108,6 @@ func Transform(s beam.Scope, query string, opts ...Option) beam.PCollection { payload := beam.CrossLanguagePayload(&sqlx.ExpansionPayload{ Query: query, Dialect: o.dialect, - Options: o.customs, }) expansionAddr := sqlx.DefaultExpansionAddr diff --git a/sdks/go/pkg/beam/transforms/sql/sqlx/sqlx.go b/sdks/go/pkg/beam/transforms/sql/sqlx/sqlx.go index d39820c37dc4..e83b92db437b 100644 --- a/sdks/go/pkg/beam/transforms/sql/sqlx/sqlx.go +++ b/sdks/go/pkg/beam/transforms/sql/sqlx/sqlx.go @@ -48,5 +48,5 @@ type Option struct { type ExpansionPayload struct { Query string `beam:"query"` Dialect string `beam:"dialect"` - Options []Option `beam:"options"` + options []Option `beam:"options"` }