summaryrefslogtreecommitdiff
path: root/src/github.com
diff options
context:
space:
mode:
authorWisdom Omuya <deafgoat@gmail.com>2014-09-29 15:56:43 -0400
committerWisdom Omuya <deafgoat@gmail.com>2014-09-29 16:37:49 -0400
commit1d39e81053332f6e01ae1168f5ddf8ac21cf008a (patch)
treef0c6160c87fbcb62db41fc9707fc8c053c88da32 /src/github.com
parentebc9fb27ad9697fa0fb784c40fc29413ff40f135 (diff)
add --upsertFields and --directoryperdb for mongoimport2.7.7
Former-commit-id: 7ff74594390549b9a06f92500c8d374a2dcda027
Diffstat (limited to 'src/github.com')
-rw-r--r--src/github.com/mongodb/mongo-tools/common/db/shim.go58
-rw-r--r--src/github.com/mongodb/mongo-tools/mongofiles/mongofiles.go7
-rw-r--r--src/github.com/mongodb/mongo-tools/mongoimport/import_writers.go61
-rw-r--r--src/github.com/mongodb/mongo-tools/mongoimport/main/mongoimport.go2
-rw-r--r--src/github.com/mongodb/mongo-tools/mongoimport/mongoimport.go30
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
}