From 2df015d4459481279dfae23f1ebaf81a2d2cb0ff Mon Sep 17 00:00:00 2001 From: Asad Ullah Date: Fri, 29 Dec 2023 15:33:58 +0500 Subject: [PATCH 1/4] Update oplog_manager.py --- mongo_connector/oplog_manager.py | 12 +++++++++--- 1 file changed, 9 insertions(+), 3 deletions(-) diff --git a/mongo_connector/oplog_manager.py b/mongo_connector/oplog_manager.py index cfd6d48a..3dac46e9 100755 --- a/mongo_connector/oplog_manager.py +++ b/mongo_connector/oplog_manager.py @@ -193,6 +193,12 @@ def _should_skip_entry(self, entry): ) return True, False + # Within mongo-6 in case of update the op-log does not contain the + # complete document, so we must re-fetch it from database + if entry["op"] == "u" and 'diff' in entry["o"]: + from_coll = self.get_collection(entry["ns"]) + entry["o"] = from_coll.find_one({'_id': entry["o2"]["_id"]}) + # Update the namespace. entry["ns"] = namespace.dest_name @@ -545,7 +551,7 @@ def get_all_ns(): if database == "config" or database == "local": continue coll_list = retry_until_ok( - self.primary_client[database].collection_names + self.primary_client[database].list_collection_names ) for coll in coll_list: # ignore system collections @@ -607,7 +613,7 @@ def upsert_each(dm): for namespace in dump_set: from_coll = self.get_collection(namespace) mapped_ns = self.namespace_config.map_namespace(namespace) - total_docs = retry_until_ok(from_coll.count) + total_docs = retry_until_ok(from_coll.aggregate, [{'$match': {}}]).next() num = None for num, doc in enumerate(docs_to_dump(from_coll)): try: @@ -641,7 +647,7 @@ def upsert_all(dm): try: for namespace in dump_set: from_coll = self.get_collection(namespace) - total_docs = retry_until_ok(from_coll.count) + total_docs = retry_until_ok(from_coll.aggregate, [{'$match': {}}]).next() mapped_ns = self.namespace_config.map_namespace(namespace) LOG.info( "Bulk upserting approximately %d docs from " "collection '%s'", From 19c24c91144d7fd14e8cf29cc495f709a6f77647 Mon Sep 17 00:00:00 2001 From: Asad Ullah Date: Fri, 29 Dec 2023 16:57:50 +0500 Subject: [PATCH 2/4] Update oplog_manager.py --- mongo_connector/oplog_manager.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/mongo_connector/oplog_manager.py b/mongo_connector/oplog_manager.py index 3dac46e9..cefd195c 100755 --- a/mongo_connector/oplog_manager.py +++ b/mongo_connector/oplog_manager.py @@ -196,6 +196,8 @@ def _should_skip_entry(self, entry): # Within mongo-6 in case of update the op-log does not contain the # complete document, so we must re-fetch it from database if entry["op"] == "u" and 'diff' in entry["o"]: + LOG.debug("OplogThread: updating entry '%s' from " + "collection '%s'" % (entry["o2"]["_id"], entry["ns"])) from_coll = self.get_collection(entry["ns"]) entry["o"] = from_coll.find_one({'_id': entry["o2"]["_id"]}) From 04316a7a52dee929ea03114c595632291aa22ced Mon Sep 17 00:00:00 2001 From: Asad Ullah Date: Mon, 1 Jan 2024 18:06:43 +0500 Subject: [PATCH 3/4] Update oplog_manager.py change total count funtion --- mongo_connector/oplog_manager.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/mongo_connector/oplog_manager.py b/mongo_connector/oplog_manager.py index cefd195c..25d11bd0 100755 --- a/mongo_connector/oplog_manager.py +++ b/mongo_connector/oplog_manager.py @@ -615,7 +615,7 @@ def upsert_each(dm): for namespace in dump_set: from_coll = self.get_collection(namespace) mapped_ns = self.namespace_config.map_namespace(namespace) - total_docs = retry_until_ok(from_coll.aggregate, [{'$match': {}}]).next() + total_docs = retry_until_ok(from_coll.count_documents, {}) num = None for num, doc in enumerate(docs_to_dump(from_coll)): try: @@ -649,7 +649,7 @@ def upsert_all(dm): try: for namespace in dump_set: from_coll = self.get_collection(namespace) - total_docs = retry_until_ok(from_coll.aggregate, [{'$match': {}}]).next() + total_docs = retry_until_ok(from_coll.count_documents, {}) mapped_ns = self.namespace_config.map_namespace(namespace) LOG.info( "Bulk upserting approximately %d docs from " "collection '%s'", From b2f96d3ed5c7fc5cc527567cff60a1f5a76fb648 Mon Sep 17 00:00:00 2001 From: Asad Ullah Date: Wed, 27 Mar 2024 20:59:33 +0500 Subject: [PATCH 4/4] Update connector.py --- mongo_connector/connector.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/mongo_connector/connector.py b/mongo_connector/connector.py index 5103f8d4..3e9e7803 100644 --- a/mongo_connector/connector.py +++ b/mongo_connector/connector.py @@ -39,6 +39,7 @@ from mongo_connector.util import log_fatal_exceptions, retry_until_ok from mongo_connector.namespace_config import NamespaceConfig, validate_namespace_options from mongo_connector.version import Version +from bson.codec_options import DatetimeConversion # Monkey patch logging to add Logger.always @@ -332,7 +333,8 @@ def create_authed_client(self, hosts=None, **kwargs): new_uri = self.address else: new_uri = self.copy_uri_options(hosts, self.address) - client = MongoClient(new_uri, tz_aware=self.tz_aware, **kwargs) + client = MongoClient(new_uri, tz_aware=self.tz_aware, datetime_conversion=DatetimeConversion.DATETIME_CLAMP, **kwargs) + if self.auth_key is not None: client["admin"].authenticate(self.auth_username, self.auth_key) return client