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
26 changes: 22 additions & 4 deletions connectors/mongo/conn.go
Original file line number Diff line number Diff line change
Expand Up @@ -37,13 +37,20 @@ type ConnectorSettings struct {
TargetDocCountPerPartition int64 //target number of documents per partition (256k docs is 256MB with 1KB average doc size)
SampleFactor int // a factor to determine how many extra samples per partition are used
MaxPageSize int
NamespaceFanout int
DocumentDBSamplingFanout int
PerNamespaceStreams bool
SkipBatchOverwrite bool
FullDocumentKey bool

Query string // query filter, as a v2 Extended JSON string, e.g., '{\"x\":{\"$gt\":1}}'"
}

const (
defaultNamespaceFanoutLimit = 100
defaultDocumentDBSamplingFanoutLimit = 100
)

func setDefault[T comparable](field *T, defaultValue T) {
if *field == *new(T) {
*field = defaultValue
Expand Down Expand Up @@ -239,6 +246,7 @@ func (c *conn) GeneratePlan(ctx context.Context, r *connect.Request[adiomv1.Gene

done := make(chan struct{})
eg, ctx := errgroup.WithContext(ctx)
eg.SetLimit(c.planningFanoutLimit())
var finalPartitions []*adiomv1.Partition
ch := make(chan *adiomv1.Partition)

Expand Down Expand Up @@ -288,6 +296,7 @@ func (c *conn) GeneratePlan(ctx context.Context, r *connect.Request[adiomv1.Gene
if numSamples > 1000000 {
slog.Warn("More than 1000000 samples requested", "samples", numSamples)
}
slog.Debug("Getting samples for namespace", "namespace", partition.GetNamespace(), "samples", numSamples)
ids, err := c.sampleIDs(ctx, col, numSamples)
if err != nil {
return fmt.Errorf("error getting %v samples for partition: %w", numSamples, err)
Expand Down Expand Up @@ -330,6 +339,13 @@ func (c *conn) GeneratePlan(ctx context.Context, r *connect.Request[adiomv1.Gene
}), nil
}

func (c *conn) planningFanoutLimit() int {
if c.settings.NamespaceFanout > 0 {
return c.settings.NamespaceFanout
}
return defaultNamespaceFanoutLimit
}

// GetInfo implements adiomv1connect.ConnectorServiceHandler.
func (c *conn) GetInfo(ctx context.Context, r *connect.Request[adiomv1.GetInfoRequest]) (*connect.Response[adiomv1.GetInfoResponse], error) {
// Get version of the MongoDB server
Expand Down Expand Up @@ -691,11 +707,11 @@ func toTimestampPB(t bson.Timestamp) *timestamppb.Timestamp {
}

type MongoUpdate struct {
NS bson.M `bson:"ns"`
NS bson.M `bson:"ns"`
ClusterTime *bson.Timestamp `bson:"clusterTime"`
DocumentKey bson.D `bson:"documentKey"`
FullDocument bson.Raw `bson:"fullDocument"`
OperationType string `bson:"operationType"`
DocumentKey bson.D `bson:"documentKey"`
FullDocument bson.Raw `bson:"fullDocument"`
OperationType string `bson:"operationType"`
}

// StreamUpdates implements adiomv1connect.ConnectorServiceHandler.
Expand Down Expand Up @@ -1122,6 +1138,8 @@ func (c *conn) Teardown() {
func NewConn(connSettings ConnectorSettings) (adiomv1connect.ConnectorServiceHandler, error) {
setDefault(&connSettings.TargetDocCountPerPartition, 512*1000)
setDefault(&connSettings.SampleFactor, 1)
setDefault(&connSettings.NamespaceFanout, defaultNamespaceFanoutLimit)
setDefault(&connSettings.DocumentDBSamplingFanout, defaultDocumentDBSamplingFanoutLimit)
client, err := MongoClient(context.Background(), connSettings)
if err != nil {
slog.Error(fmt.Sprintf("unable to connect to mongo client: %v", err))
Expand Down
17 changes: 14 additions & 3 deletions connectors/mongo/conn_unit_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,17 @@ func bsonID(t *testing.T, name string, v interface{}) *adiomv1.BsonValue {
return &adiomv1.BsonValue{Name: name, Type: uint32(typ), Data: data}
}

func TestPlanningFanoutLimit(t *testing.T) {
assert.Equal(t, defaultNamespaceFanoutLimit, (&conn{flavor: FlavorMongoDB}).planningFanoutLimit())
assert.Equal(t, defaultNamespaceFanoutLimit, (&conn{flavor: FlavorDocumentDB}).planningFanoutLimit())
assert.Equal(t, 12, (&conn{settings: ConnectorSettings{NamespaceFanout: 12}, flavor: FlavorDocumentDB}).planningFanoutLimit())
}

func TestDocumentDBSamplingFanout(t *testing.T) {
assert.Equal(t, defaultDocumentDBSamplingFanoutLimit, (&conn{}).documentDBSamplingFanout())
assert.Equal(t, 17, (&conn{settings: ConnectorSettings{DocumentDBSamplingFanout: 17}}).documentDBSamplingFanout())
}

// buildIdFilter

func TestBuildIdFilter_EmptyId(t *testing.T) {
Expand Down Expand Up @@ -173,9 +184,9 @@ func deleteU(t *testing.T, id string) *adiomv1.Update {

func partialU(t *testing.T, id string, data bson.M, unset ...string) *adiomv1.Update {
u := &adiomv1.Update{
Id: []*adiomv1.BsonValue{bsonID(t, "_id", id)},
Type: adiomv1.UpdateType_UPDATE_TYPE_PARTIAL_UPDATE,
PartialUpdateUnset: unset,
Id: []*adiomv1.BsonValue{bsonID(t, "_id", id)},
Type: adiomv1.UpdateType_UPDATE_TYPE_PARTIAL_UPDATE,
PartialUpdateUnset: unset,
}
if data != nil {
u.Data = mustMarshal(t, data)
Expand Down
8 changes: 8 additions & 0 deletions connectors/mongo/docdb.go
Original file line number Diff line number Diff line change
Expand Up @@ -71,6 +71,7 @@ func (c *conn) sampleIDs(ctx context.Context, col *mongo.Collection, numSamples

results := make([]sampleResult, numSamples)
eg, ctx := errgroup.WithContext(ctx)
eg.SetLimit(c.documentDBSamplingFanout())
for i := int64(0); i < numSamples; i++ {
i := i
eg.Go(func() error {
Expand Down Expand Up @@ -122,3 +123,10 @@ func (c *conn) sampleIDs(ctx context.Context, col *mongo.Collection, numSamples

return ids, nil
}

func (c *conn) documentDBSamplingFanout() int {
if c.settings.DocumentDBSamplingFanout > 0 {
return c.settings.DocumentDBSamplingFanout
}
return defaultDocumentDBSamplingFanoutLimit
}
6 changes: 6 additions & 0 deletions gen/adiom/commands/connectors/v1/cosmos_commandargs.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

14 changes: 12 additions & 2 deletions gen/adiom/commands/connectors/v1/mongo.pb.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

12 changes: 12 additions & 0 deletions gen/adiom/commands/connectors/v1/mongo_commandargs.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

42 changes: 27 additions & 15 deletions internal/app/options/connectorflags.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,8 +23,8 @@ import (
"github.com/adiom-data/dsync/connectors/sqlbatch"
"github.com/adiom-data/dsync/connectors/testconn"
"github.com/adiom-data/dsync/connectors/vector"
adiomv1 "github.com/adiom-data/dsync/gen/adiom/v1"
connectorsv1 "github.com/adiom-data/dsync/gen/adiom/commands/connectors/v1"
adiomv1 "github.com/adiom-data/dsync/gen/adiom/v1"
"github.com/adiom-data/dsync/gen/adiom/v1/adiomv1connect"
"github.com/urfave/cli/v2"
"github.com/urfave/cli/v2/altsrc"
Expand Down Expand Up @@ -359,12 +359,14 @@ func GetRegisteredConnectors() []RegisteredConnector {
return false
},
Create: func(args []string, as AdditionalSettings) (adiomv1connect.ConnectorServiceHandler, []string, error) {
return CreateHelper("MongoDB", "mongodb://connection-string [options]", connectorsv1.MongoFlagsCommandFlags(), func(c *cli.Context, args []string, _ AdditionalSettings) (adiomv1connect.ConnectorServiceHandler, error) {
return CreateHelper("MongoDB", "mongodb://connection-string [options]", mongoFlagsCommandFlags(), func(c *cli.Context, args []string, _ AdditionalSettings) (adiomv1connect.ConnectorServiceHandler, error) {
flags, err := connectorsv1.ParseMongoFlags(c)
if err != nil {
return nil, err
}
return mongo.NewConn(mongoSettingsFromFlags(flags, args[0]))
settings := mongoSettingsFromFlags(flags, args[0])
settings.DocumentDBSamplingFanout = c.Int("documentdb-sampling-fanout")
return mongo.NewConn(settings)
})(args, as)
},
},
Expand Down Expand Up @@ -696,18 +698,18 @@ func fileSettingsFromFlags(f *connectorsv1.FileFlags, uri string) (fileconnector

func s3SettingsFromFlags(f *connectorsv1.S3Flags, uri string) s3connector.ConnectorSettings {
return s3connector.ConnectorSettings{
Uri: uri,
PrettyJSON: f.PrettyJson,
Region: f.Region,
Prefix: f.Prefix,
OutputFormat: f.OutputFormat,
Profile: f.Profile,
Endpoint: f.Endpoint,
AccessKeyID: f.AccessKeyId,
SecretAccessKey: f.SecretAccessKey,
SessionToken: f.SessionToken,
UsePathStyle: f.UsePathStyle,
MaxFileSizeMB: f.MaxFileSize,
Uri: uri,
PrettyJSON: f.PrettyJson,
Region: f.Region,
Prefix: f.Prefix,
OutputFormat: f.OutputFormat,
Profile: f.Profile,
Endpoint: f.Endpoint,
AccessKeyID: f.AccessKeyId,
SecretAccessKey: f.SecretAccessKey,
SessionToken: f.SessionToken,
UsePathStyle: f.UsePathStyle,
MaxFileSizeMB: f.MaxFileSize,
MaxTotalMemoryMB: f.MaxTotalMemory,
}
}
Expand All @@ -724,10 +726,20 @@ func mongoBaseSettingsFromFlags(f *connectorsv1.MongoBaseFlags, connString strin
s.WriterMaxBatchSize = int(f.WriterBatchSize)
s.TargetDocCountPerPartition = f.DocPartition
s.MaxPageSize = int(f.MaxPageSize)
s.NamespaceFanout = int(f.NamespaceFanout)
s.Query = f.InitialSyncQuery
return s
}

func mongoFlagsCommandFlags() []cli.Flag {
flags := connectorsv1.MongoFlagsCommandFlags()
return append(flags, &cli.IntFlag{
Name: "documentdb-sampling-fanout",
Value: 100,
Hidden: true,
})
}

func mongoSettingsFromFlags(f *connectorsv1.MongoFlags, connString string) mongo.ConnectorSettings {
s := mongoBaseSettingsFromFlags(f.Base, connString)
s.SampleFactor = int(f.SampleFactor)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,11 @@ message MongoBaseFlags {
usage: "query filter for the initial data copy (v2 Extended JSON)"
aliases: ["q"]
}];
int32 namespace_fanout = 9 [(commandargs.v1.flag) = {
name: "namespace-fanout"
usage: "maximum number of namespaces to plan concurrently"
default: "100"
}];
}

// Full MongoDB flags = base + MongoDB-specific extras
Expand Down
Loading