Skip to content

Commit e06570c

Browse files
committed
[IMP] dms_import: improve performance, reduce batch size, and bypass heavy compute fields
1 parent f755c12 commit e06570c

1 file changed

Lines changed: 148 additions & 111 deletions

File tree

dms_import/hooks.py

Lines changed: 148 additions & 111 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@
22
# License AGPL-3.0 or later (https://www.gnu.org/licenses/agpl).
33

44
import logging
5+
from collections import defaultdict
56
from random import randint
67

78
from psycopg2.sql import SQL, Identifier
@@ -37,68 +38,48 @@ def _get_or_create_default_storage(env):
3738
return db_storage
3839

3940

40-
def _assign_dms_groups(cr, env, folder_id, new_dir):
41-
"""Migrate read/write groups from a documents.folder to dms.access.groups."""
42-
DmsAccessGroups = env["dms.access.group"]
43-
all_new_groups = DmsAccessGroups.browse()
41+
def _batch_fetch_folder_relations(cr, folder_ids):
42+
write_groups = defaultdict(list)
43+
read_groups = defaultdict(list)
44+
tags_by_folder = defaultdict(list)
4445

45-
# Handle Write Groups
4646
if table_exists(cr, "documents_folder_res_groups_rel"):
4747
cr.execute(
4848
SQL(
49-
"""
50-
SELECT res_groups_id FROM {}
51-
WHERE documents_folder_id = %s
52-
"""
53-
).format(Identifier("documents_folder_res_groups_rel")),
54-
(folder_id,),
49+
"""SELECT documents_folder_id, res_groups_id
50+
FROM documents_folder_res_groups_rel
51+
WHERE documents_folder_id = ANY(%s)"""
52+
),
53+
(list(folder_ids),),
5554
)
56-
group_ids = [r[0] for r in cr.fetchall()]
57-
if group_ids:
58-
write_group = DmsAccessGroups.create(
59-
{
60-
"name": f"{new_dir.name} Write Group",
61-
"perm_create": True,
62-
"perm_write": True,
63-
"perm_unlink": True,
64-
"group_ids": [Command.set(group_ids)],
65-
}
66-
)
67-
all_new_groups |= write_group
55+
for folder_id, group_id in cr.fetchall():
56+
write_groups[folder_id].append(group_id)
6857

69-
# Handle Read Groups
7058
if table_exists(cr, "documents_folder_read_groups"):
7159
cr.execute(
7260
SQL(
73-
"""
74-
SELECT res_groups_id FROM {}
75-
WHERE documents_folder_id = %s
76-
"""
77-
).format(Identifier("documents_folder_read_groups")),
78-
(folder_id,),
61+
"""SELECT documents_folder_id, res_groups_id
62+
FROM documents_folder_read_groups
63+
WHERE documents_folder_id = ANY(%s)"""
64+
),
65+
(list(folder_ids),),
7966
)
80-
read_group_ids = [r[0] for r in cr.fetchall()]
81-
if read_group_ids:
82-
read_group = DmsAccessGroups.create(
83-
{
84-
"name": f"{new_dir.name} Read Group",
85-
"group_ids": [Command.set(read_group_ids)],
86-
}
87-
)
88-
all_new_groups |= read_group
67+
for folder_id, group_id in cr.fetchall():
68+
read_groups[folder_id].append(group_id)
8969

90-
# If dont have any group, we need to assign a group to visible folder
91-
if not all_new_groups:
92-
all_new_groups = DmsAccessGroups.create(
93-
{
94-
"name": f"{new_dir.name} Default Group",
95-
"perm_create": True,
96-
"perm_write": True,
97-
"perm_unlink": True,
98-
"group_ids": [Command.set(env.ref("base.group_user").ids)],
99-
}
70+
if table_exists(cr, "documents_facet") and table_exists(cr, "documents_tag"):
71+
cr.execute(
72+
SQL(
73+
"""SELECT f.folder_id, t.id
74+
FROM {} t JOIN {} f ON f.id = t.facet_id
75+
WHERE f.folder_id = ANY(%s)"""
76+
).format(Identifier("documents_tag"), Identifier("documents_facet")),
77+
(list(folder_ids),),
10078
)
101-
return all_new_groups
79+
for folder_id, tag_id in cr.fetchall():
80+
tags_by_folder[folder_id].append(tag_id)
81+
82+
return write_groups, read_groups, tags_by_folder
10283

10384

10485
def _create_unique_directory(DmsDirectory, vals):
@@ -131,7 +112,7 @@ def migrate_documents_tags(cr, env, lang):
131112
tag_mapping = {}
132113
category_mapping = {}
133114

134-
# 1. FACETSCATEGORIES
115+
# 1. Facetscategories
135116
if table_exists(cr, "documents_facet"):
136117
cr.execute(
137118
SQL("SELECT id, name->>%s AS name FROM {}").format(
@@ -166,7 +147,7 @@ def migrate_documents_tags(cr, env, lang):
166147
if norm and norm in existing_categories:
167148
category_mapping[f["id"]] = existing_categories[norm].id
168149

169-
# 2. TAGS
150+
# 2. Tags
170151
existing_tags = {
171152
(_normalize(tag.name), tag.category_id.id or False): tag
172153
for tag in DmsTag.search([])
@@ -227,7 +208,7 @@ def migrate_documents_tags(cr, env, lang):
227208
return tag_mapping, category_mapping
228209

229210

230-
def migrate_documents_folders(cr, env, lang, tag_mapping, category_mapping):
211+
def migrate_documents_folders(cr, env, lang, tag_mapping):
231212
"""Migrate documents.folder to dms.directory, optimized for batch processing."""
232213
if not table_exists(cr, "documents_folder"):
233214
_logger.warning("Skipping folder migration: 'documents_folder' not found.")
@@ -242,11 +223,22 @@ def migrate_documents_folders(cr, env, lang, tag_mapping, category_mapping):
242223
(lang,),
243224
)
244225
folders = cr.dictfetchall()
226+
if not folders:
227+
_logger.info("No folders to migrate.")
228+
return {}
245229

246230
DmsDirectory = env["dms.directory"]
231+
DmsAccessGroups = env["dms.access.group"]
247232
folder_mapping = {}
248233
storage = _get_or_create_default_storage(env)
234+
default_user_group_id = env.ref("base.group_user").id
249235

236+
folder_ids = [f["id"] for f in folders]
237+
write_groups, read_groups, tags_by_folder = _batch_fetch_folder_relations(
238+
cr, folder_ids
239+
)
240+
access_groups_to_create_vals = []
241+
dirs_to_update_tags = defaultdict(list)
250242
for folder in folders:
251243
try:
252244
parent_id = folder_mapping.get(folder["parent_folder_id"])
@@ -263,49 +255,83 @@ def migrate_documents_folders(cr, env, lang, tag_mapping, category_mapping):
263255
"Created dms.directory: '%s' (ID: %s)", new_dir.name, new_dir.id
264256
)
265257

266-
# Assign groups
267-
dms_groups = _assign_dms_groups(cr, env, folder["id"], new_dir)
268-
dms_groups._compute_users() # trigger to recompute count_users
269-
new_dir.group_ids = [Command.set(dms_groups.ids)]
270-
_logger.info(
271-
"Assigned %s groups to directory '%s'",
272-
dms_groups.mapped("name"),
273-
new_dir.name,
274-
)
258+
# Assign groups using pre-fetched data
259+
has_groups = False
260+
write_group_ids = write_groups.get(folder["id"], [])
261+
if write_group_ids:
262+
has_groups = True
263+
vals = {
264+
"name": f"{new_dir.name} Write Group",
265+
"perm_create": True,
266+
"perm_write": True,
267+
"perm_unlink": True,
268+
"group_ids": [Command.set(write_group_ids)],
269+
"_dir_id": new_dir.id,
270+
}
271+
access_groups_to_create_vals.append(vals)
272+
read_group_ids = read_groups.get(folder["id"], [])
273+
if read_group_ids:
274+
has_groups = True
275+
vals = {
276+
"name": f"{new_dir.name} Read Group",
277+
"group_ids": [Command.set(read_group_ids)],
278+
"_dir_id": new_dir.id,
279+
}
280+
access_groups_to_create_vals.append(vals)
281+
# In case no group, directory will be invisible, so I decide to assign a new group
282+
# this group has full permision like no group in document folder
283+
if not has_groups:
284+
vals = {
285+
"name": f"{new_dir.name} Default Group",
286+
"perm_create": True,
287+
"perm_write": True,
288+
"perm_unlink": True,
289+
"group_ids": [Command.set([default_user_group_id])],
290+
"_dir_id": new_dir.id,
291+
}
292+
access_groups_to_create_vals.append(vals)
275293

276294
# Assign Tags
277-
if table_exists(cr, "documents_facet") and table_exists(
278-
cr, "documents_tag"
279-
):
280-
cr.execute(
281-
SQL(
282-
"""SELECT t.id FROM {} t
283-
JOIN {} f ON f.id = t.facet_id
284-
WHERE f.folder_id = %s"""
285-
).format(
286-
Identifier("documents_tag"),
287-
Identifier("documents_facet"),
288-
),
289-
(folder["id"],),
290-
)
291-
dms_tag_ids = [
292-
tag_mapping.get(r[0])
293-
for r in cr.fetchall()
294-
if tag_mapping.get(r[0])
295-
]
296-
if dms_tag_ids:
297-
new_dir.tag_ids = [Command.set(dms_tag_ids)]
295+
folder_tag_ids = [
296+
tag_mapping[tag_id]
297+
for tag_id in tags_by_folder.get(folder["id"], [])
298+
if tag_id in tag_mapping
299+
]
300+
if folder_tag_ids:
301+
dirs_to_update_tags[tuple(sorted(folder_tag_ids))].append(new_dir.id)
298302

299303
except Exception:
300304
_logger.exception(
301305
"Error migrating folder ID %s (%s)", folder["id"], folder["name"]
302306
)
303307

308+
for tag_ids_tuple, dir_ids in dirs_to_update_tags.items():
309+
DmsDirectory.browse(dir_ids).write(
310+
{"tag_ids": [Command.set(list(tag_ids_tuple))]}
311+
)
312+
313+
if access_groups_to_create_vals:
314+
new_groups = DmsAccessGroups.create(
315+
[
316+
{k: v for k, v in vals.items() if k != "_dir_id"}
317+
for vals in access_groups_to_create_vals
318+
]
319+
)
320+
# Recompute all user counts at once
321+
new_groups.invalidate_model(fnames=["users"])
322+
323+
dirs_to_update_groups = defaultdict(list)
324+
for vals, group in zip(access_groups_to_create_vals, new_groups):
325+
dirs_to_update_groups[vals["_dir_id"]].append(group.id)
326+
327+
for dir_id, group_ids in dirs_to_update_groups.items():
328+
DmsDirectory.browse(dir_id).write({"group_ids": [Command.set(group_ids)]})
329+
304330
_logger.info("Successfully migrated %d folders.", len(folder_mapping))
305331
return folder_mapping
306332

307333

308-
def migrate_documents_files(cr, env, folder_mapping, tag_mapping, category_mapping):
334+
def migrate_documents_files(cr, env, folder_mapping, tag_mapping):
309335
"""Migrate documents.document to dms.file, optimized for batch processing."""
310336
if not table_exists(cr, "documents_document") or not folder_mapping:
311337
_logger.warning("Skipping file migration: table or folder mapping is missing.")
@@ -323,9 +349,15 @@ def migrate_documents_files(cr, env, folder_mapping, tag_mapping, category_mappi
323349
return
324350

325351
DmsFile = env["dms.file"]
326-
327-
for batch_of_ids in split_every(models.INSERT_BATCH_SIZE * 10, all_doc_ids):
328-
_logger.info("Migrating files batch with %d documents", len(batch_of_ids))
352+
BATCH_SIZE = models.INSERT_BATCH_SIZE
353+
total_docs = len(all_doc_ids)
354+
_logger.info("Processing %d documents in total.", total_docs)
355+
for i, batch_of_ids in enumerate(split_every(BATCH_SIZE, all_doc_ids)):
356+
_logger.info(
357+
"Processing batch %d/%d...",
358+
i + 1,
359+
(total_docs + BATCH_SIZE - 1) // BATCH_SIZE,
360+
)
329361
cr.execute(
330362
SQL(
331363
"""SELECT id, name, folder_id, attachment_id, active
@@ -367,29 +399,38 @@ def migrate_documents_files(cr, env, folder_mapping, tag_mapping, category_mappi
367399
"active": doc["active"],
368400
"attachment_id": doc["attachment_id"],
369401
"tag_ids": [Command.set(tags_by_doc.get(doc["id"], []))],
370-
"_old_attachment_id": doc["attachment_id"],
371402
}
372403
)
373404

374405
if files_to_create_vals:
375-
new_files = DmsFile.create(
376-
[
377-
{k: v for k, v in vals.items() if k != "_old_attachment_id"}
378-
for vals in files_to_create_vals
379-
]
380-
)
406+
# avoid _compute_content in dms.file
407+
with env.norecompute():
408+
new_files = DmsFile.create(files_to_create_vals)
381409

382-
attachments_to_update = {}
410+
# Update attachments in bulk
411+
attachment_update_map = {}
383412
for vals, new_file in zip(files_to_create_vals, new_files):
384-
if vals["_old_attachment_id"]:
385-
attachments_to_update[vals["_old_attachment_id"]] = new_file.id
386-
387-
if attachments_to_update:
388-
for att_id, file_id in attachments_to_update.items():
389-
env["ir.attachment"].browse(att_id).write(
390-
{"res_model": "dms.file", "res_id": file_id}
391-
)
413+
if vals.get("attachment_id"):
414+
attachment_update_map[vals["attachment_id"]] = new_file.id
415+
416+
if attachment_update_map:
417+
attachment_ids = list(attachment_update_map.keys())
418+
case_clauses = []
419+
params = []
420+
for att_id, file_id in attachment_update_map.items():
421+
case_clauses.append(SQL("WHEN %s THEN %s"))
422+
params.extend([att_id, file_id])
423+
params.append(tuple(attachment_ids))
424+
query = SQL(
425+
"""
426+
UPDATE ir_attachment
427+
SET res_model = 'dms.file',
428+
res_id = CASE id {cases} END
429+
WHERE id IN %s
430+
"""
431+
).format(cases=SQL(" ").join(case_clauses))
392432

433+
cr.execute(query, params)
393434
_logger.info("Files migration finished.")
394435

395436

@@ -403,15 +444,11 @@ def post_init_hook(cr, registry):
403444
try:
404445
lang = get_lang(env).code or env.lang
405446
_logger.info(
406-
"Starting migration from EE 'documents' to OCA 'dms' using lang '%s'", lang
447+
"Starting migration from 'documents' to 'dms' using lang '%s'", lang
407448
)
408-
409-
tag_mapping, category_mapping = migrate_documents_tags(cr, env, lang)
410-
folder_mapping = migrate_documents_folders(
411-
cr, env, lang, tag_mapping, category_mapping
412-
)
413-
migrate_documents_files(cr, env, folder_mapping, tag_mapping, category_mapping)
414-
449+
tag_mapping, _ = migrate_documents_tags(cr, env, lang)
450+
folder_mapping = migrate_documents_folders(cr, env, lang, tag_mapping)
451+
migrate_documents_files(cr, env, folder_mapping, tag_mapping)
415452
_logger.info("Migration from 'documents' to 'dms' completed successfully.")
416453
except Exception:
417454
_logger.exception("Migration from 'documents' to 'dms' failed!")

0 commit comments

Comments
 (0)