package main import ( "context" "fmt" "go.mongodb.org/mongo-driver/mongo" "go.mongodb.org/mongo-driver/mongo/options" "go.mongodb.org/mongo-driver/mongo/readpref" "github.com/globalsign/mgo/bson" "time" ) type InsertChangeStream struct{ OperationType string `json:"operationType"` Document interface{} `json:"fullDocument"` DocumentKey DocumentKeyStruct `json:"documentKey"` } type DocumentKeyStruct struct{ Id map[string]string `json:"_id"` } // TODO: Get below structs from models // TODO: Add es tags type DemoDocument struct { A string `bson:"a" json:"a"` B string `bson:"b" json:"b"` } func main() { mongoDBURI := "mongodb://localhost:27017" client, err := mongo.NewClient(options.Client().ApplyURI(mongoDBURI)) ctx, _ := context.WithTimeout(context.Background(), 30*time.Second) _ = client _ = ctx _ = err err = client.Connect(ctx) if err != nil { fmt.Println("Err while context.WithTimeout : ", err.Error()) } err = client.Ping(ctx, readpref.Primary()) if err != nil { fmt.Println("error while ping : ", err.Error()) } fmt.Println("connection established") db := client.Database("testdb") coll := db.Collection("collection2") cs, err := coll.Watch(ctx, mongo.Pipeline{}) defer cs.Close(ctx) if err != nil { fmt.Println("Err while watching collection : ", err.Error()) } ok := cs.Next(ctx) fmt.Println(" ok fasdf", ok ) for ok { next := cs.Current //fmt.Println(ok, next) fmt.Println(next.String()) var insertData interface{} //err = json.Unmarshal([]byte(next.String()), &insertData) //if err != nil { // fmt.Println("Error while json marshalling: ", err.Error()) //} cs.Decode(insertData) fmt.Println("insert data", insertData) ok = cs.Next(ctx) //document := insertData.Document //objectId := insertData.DocumentKey.Id["$oid"] } //for cs.Next(ctx) { // //elem := bson.NewDocument() // if err := cs.Decode(elem); err != nil { // log.Fatal(err) // } // //} }