diff options
| author | Wisdom Omuya <deafgoat@gmail.com> | 2014-09-29 15:56:43 -0400 |
|---|---|---|
| committer | Wisdom Omuya <deafgoat@gmail.com> | 2014-09-29 16:37:49 -0400 |
| commit | 1d39e81053332f6e01ae1168f5ddf8ac21cf008a (patch) | |
| tree | f0c6160c87fbcb62db41fc9707fc8c053c88da32 | |
| parent | ebc9fb27ad9697fa0fb784c40fc29413ff40f135 (diff) | |
add --upsertFields and --directoryperdb for mongoimport2.7.7
Former-commit-id: 7ff74594390549b9a06f92500c8d374a2dcda027
5 files changed, 78 insertions, 80 deletions
diff --git a/src/github.com/mongodb/mongo-tools/common/db/shim.go b/src/github.com/mongodb/mongo-tools/common/db/shim.go index 02f313b0aa6..c843b18853f 100644 --- a/src/github.com/mongodb/mongo-tools/common/db/shim.go +++ b/src/github.com/mongodb/mongo-tools/common/db/shim.go @@ -20,13 +20,33 @@ const MaxBSONSize = 16 * 1024 * 1024 type ShimMode int +var ErrShimNotFound = errors.New("Shim not found") + const ( Dump ShimMode = iota Insert + Upsert Drop Remove ) +type StorageShim struct { + DBPath string + Database string + Collection string + Skip int + Limit int + ShimPath string + Query string + UpsertFields string + Sort []string + DirectoryPerDB bool + Journal bool + Mode ShimMode + shimProcess *exec.Cmd + stdin io.WriteCloser +} + type Shim struct { DBPath string ShimPath string @@ -256,22 +276,6 @@ func (shim *Shim) Run(command interface{}, out interface{}, database string) (er return err } -type StorageShim struct { - DBPath string - Database string - Collection string - Skip int - Limit int - ShimPath string - Query string - Sort []string - DirectoryPerDB bool - Journal bool - Mode ShimMode - shimProcess *exec.Cmd - stdin io.WriteCloser -} - func makeSort(fields []string) bson.D { val := bson.D{} for _, field := range fields { @@ -318,25 +322,21 @@ func buildArgs(shim StorageShim) ([]string, error) { returnVal = append(returnVal, "--sort", string(sortObjJson)) } - if shim.Mode != Drop && shim.Query != "" { returnVal = append(returnVal, "--query", shim.Query) } - /* - if shim.Mode == Dump { - returnVal = append(returnVal, "--query", shim.Sort) - } - */ - + returnVal = append(returnVal, "--mode") switch shim.Mode { case Dump: case Insert: - returnVal = append(returnVal, "--load") + returnVal = append(returnVal, "insert") + case Upsert: + returnVal = append(returnVal, "upsert", "--upsertFields", shim.UpsertFields) case Drop: - returnVal = append(returnVal, "--drop") + returnVal = append(returnVal, "drop") case Remove: - returnVal = append(returnVal, "--remove") + returnVal = append(returnVal, "remove") } return returnVal, nil } @@ -351,8 +351,6 @@ func checkExists(path string) (bool, error) { return true, nil } -var ErrShimNotFound = errors.New("Shim not found") - func LocateShim() (string, error) { shimLoc := os.Getenv("MONGOSHIM") if shimLoc == "" { @@ -407,14 +405,12 @@ func (shim *StorageShim) Open() (*BSONSource, *BSONSink, error) { return nil, nil, err } shim.shimProcess = cmd - return &BSONSource{stdOut, nil}, &BSONSink{stdin}, nil } func (shim *StorageShim) WaitResult() error { if shim.shimProcess != nil { - err := shim.shimProcess.Wait() - return err + return shim.shimProcess.Wait() } return nil } diff --git a/src/github.com/mongodb/mongo-tools/mongofiles/mongofiles.go b/src/github.com/mongodb/mongo-tools/mongofiles/mongofiles.go index e2dc7f595b3..bb80c4c5feb 100644 --- a/src/github.com/mongodb/mongo-tools/mongofiles/mongofiles.go +++ b/src/github.com/mongodb/mongo-tools/mongofiles/mongofiles.go @@ -2,8 +2,8 @@ package mongofiles import ( "fmt" - "github.com/mongodb/mongo-tools/common/db" "github.com/mongodb/mongo-tools/common/bsonutil" + "github.com/mongodb/mongo-tools/common/db" commonopts "github.com/mongodb/mongo-tools/common/options" "github.com/mongodb/mongo-tools/mongofiles/options" "gopkg.in/mgo.v2" @@ -38,7 +38,6 @@ const ( DefaultChunkSize = 255 * 1024 ) - type MongoFiles struct { // generic mongo tool options ToolOptions *commonopts.ToolOptions @@ -297,7 +296,7 @@ func (self *MongoFiles) ensureIndex(collection string, indexDoc bsonutil.Marshal var result bson.M indexExists := docSource.Next(&result) - + if err = docSource.Err(); err != nil { return fmt.Errorf("error retrieving indexes on the '%v' collection: %v", collection, err) } @@ -307,7 +306,7 @@ func (self *MongoFiles) ensureIndex(collection string, indexDoc bsonutil.Marshal if err != nil { return fmt.Errorf("error closing shim: %v", err) } - + if !indexExists { // if an index doesn't exist, create one err = self.createIndex(collection, indexDoc, indexName, isUnique) diff --git a/src/github.com/mongodb/mongo-tools/mongoimport/import_writers.go b/src/github.com/mongodb/mongo-tools/mongoimport/import_writers.go index 474165d9266..75298a2521c 100644 --- a/src/github.com/mongodb/mongo-tools/mongoimport/import_writers.go +++ b/src/github.com/mongodb/mongo-tools/mongoimport/import_writers.go @@ -28,13 +28,14 @@ type ShimImportWriter struct { importShim *db.StorageShim docSink *db.EncodedBSONSink dbPath string - dbName string + dirPerDB bool + db string collection string shimPath string } func (siw *ShimImportWriter) Open(dbName, collection string) error { - siw.dbName = dbName + siw.db = dbName siw.collection = collection shimPath, err := db.LocateShim() if err != nil { @@ -46,59 +47,57 @@ func (siw *ShimImportWriter) Open(dbName, collection string) error { func (siw *ShimImportWriter) Drop() error { dropShim := db.StorageShim{ - DBPath: siw.dbPath, - Database: siw.dbName, - Collection: siw.collection, - ShimPath: siw.shimPath, - Query: "{}", - Mode: db.Drop, + DBPath: siw.dbPath, + DirectoryPerDB: siw.dirPerDB, + Database: siw.db, + Collection: siw.collection, + ShimPath: siw.shimPath, + Query: "{}", + Mode: db.Drop, } _, _, err := dropShim.Open() if err != nil { return err } - - /* - decodedResult := db.NewDecodedBSONSource(out) - resultDoc := bson.M{} - decodedResult.Next(&resultDoc) - */ - defer dropShim.Close() return dropShim.WaitResult() } -func (siw *ShimImportWriter) initImportShim() error { +func (siw *ShimImportWriter) initImportShim(upsert bool) error { + mode := db.Insert + if upsert { + mode = db.Upsert + } importShim := &db.StorageShim{ - DBPath: siw.dbPath, - Database: siw.dbName, - Collection: siw.collection, - ShimPath: siw.shimPath, - Query: "", - Mode: db.Insert, + DBPath: siw.dbPath, + DirectoryPerDB: siw.dirPerDB, + Database: siw.db, + Collection: siw.collection, + ShimPath: siw.shimPath, + Query: "", + Mode: mode, + UpsertFields: strings.Join(siw.upsertFields, ","), } _, inStream, err := importShim.Open() if err != nil { return err } siw.importShim = importShim - siw.docSink = &db.EncodedBSONSink{inStream, importShim} + siw.docSink = &db.EncodedBSONSink{ + BSONIn: inStream, + WriterShim: importShim, + } return nil } func (siw *ShimImportWriter) Import(doc bson.M) error { if siw.importShim == nil { - //lazily initialize import shim - err := siw.initImportShim() - if err != nil { + // lazily initialize import shim + if err := siw.initImportShim(siw.upsertMode); err != nil { return err } } - err := siw.docSink.WriteDoc(doc) - if err != nil { - return err - } - return nil + return siw.docSink.WriteDoc(doc) } func (siw *ShimImportWriter) Close() error { diff --git a/src/github.com/mongodb/mongo-tools/mongoimport/main/mongoimport.go b/src/github.com/mongodb/mongo-tools/mongoimport/main/mongoimport.go index 857c2edadcd..b00547ca6ea 100644 --- a/src/github.com/mongodb/mongo-tools/mongoimport/main/mongoimport.go +++ b/src/github.com/mongodb/mongo-tools/mongoimport/main/mongoimport.go @@ -24,7 +24,7 @@ func main() { _, err := opts.Parse() if err != nil { - fmt.Fprintf(os.Stderr, "Error parsing command line options: %v", err) + fmt.Fprintf(os.Stderr, "Error parsing command line options: %v\n", err) util.ExitFail() } diff --git a/src/github.com/mongodb/mongo-tools/mongoimport/mongoimport.go b/src/github.com/mongodb/mongo-tools/mongoimport/mongoimport.go index 0544de8143b..da23e6e927b 100644 --- a/src/github.com/mongodb/mongo-tools/mongoimport/mongoimport.go +++ b/src/github.com/mongodb/mongo-tools/mongoimport/mongoimport.go @@ -73,14 +73,12 @@ func (mongoImport *MongoImport) getImportWriter() ImportWriter { session: nil, } } - if mongoImport.IngestOptions.Upsert { - panic("not implemented! see SERVER-15309") - } return &ShimImportWriter{ upsertMode: mongoImport.IngestOptions.Upsert, upsertFields: upsertFields, dbPath: mongoImport.ToolOptions.DBPath, - dbName: mongoImport.ToolOptions.Namespace.DB, + dirPerDB: mongoImport.ToolOptions.DirectoryPerDB, + db: mongoImport.ToolOptions.Namespace.DB, collection: mongoImport.ToolOptions.Namespace.Collection, } } @@ -183,22 +181,27 @@ func (mongoImport *MongoImport) importDocuments(importInput ImportInput) (docsCo if mongoImport.ToolOptions.Port != "" { connUrl = connUrl + ":" + mongoImport.ToolOptions.Port } - fmt.Fprintf(os.Stdout, "connected to: %v\n", connUrl) - err = importWriter.Open(mongoImport.ToolOptions.Namespace.DB, mongoImport.ToolOptions.Namespace.Collection) + util.PrintfTimeStamped("connected to: %v\n", connUrl) + + err = importWriter.Open( + mongoImport.ToolOptions.Namespace.DB, + mongoImport.ToolOptions.Namespace.Collection, + ) if err != nil { return } defer func() { - err2 := importWriter.Close() + closeErr := importWriter.Close() if err == nil { - err = err2 + err = closeErr } }() // drop the database if necessary if mongoImport.IngestOptions.Drop { - util.PrintfTimeStamped("dropping: %v.%v\n", mongoImport.ToolOptions.DB, + util.PrintfTimeStamped("dropping: %v.%v\n", + mongoImport.ToolOptions.DB, mongoImport.ToolOptions.Collection) if err := importWriter.Drop(); err != nil && @@ -207,6 +210,9 @@ func (mongoImport *MongoImport) importDocuments(importInput ImportInput) (docsCo } } + ignoreBlanks := mongoImport.IngestOptions.IgnoreBlanks && + mongoImport.InputOptions.Type != JSON + for { document, err := importInput.ImportDocument() if err != nil { @@ -223,12 +229,10 @@ func (mongoImport *MongoImport) importDocuments(importInput ImportInput) (docsCo } // ignore blank fields if specified - if mongoImport.IngestOptions.IgnoreBlanks && - mongoImport.InputOptions.Type != JSON { + if ignoreBlanks { document = removeBlankFields(document) } - err = importWriter.Import(document) - if err != nil { + if err = importWriter.Import(document); err != nil { if mongoImport.IngestOptions.StopOnError { return docsCount, err } |
