File size: 4,588 Bytes
7034da5
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
"""Compact metadata storage shared by full-index building and merging."""
import json
from contextlib import closing
import sqlite3
import zlib
from pathlib import Path

from catalog import normalize, unversioned

PER_RECORD = {'source_key','record_name','aligned_bp_length','segment_start_bp','segment_end_bp','segment_bp_length','segment_index','segment_count'}
DDL = '''
CREATE TABLE IF NOT EXISTS files(id INTEGER PRIMARY KEY,path TEXT UNIQUE,hash TEXT,size INTEGER,rows INTEGER,row_groups INTEGER,max_group_bytes INTEGER,bytes_read INTEGER,range_reads INTEGER);
CREATE TABLE IF NOT EXISTS contexts(id INTEGER PRIMARY KEY,assembly_accession TEXT,assembly_base TEXT,metadata_json TEXT UNIQUE);
CREATE TABLE IF NOT EXISTS segment_data(id INTEGER PRIMARY KEY,record_name TEXT,record_base TEXT,segment_start_bp INTEGER,segment_end_bp INTEGER,aligned_bp_length INTEGER,segment_index INTEGER,segment_count INTEGER,file_id INTEGER,row_group INTEGER,row_in_group INTEGER,context_id INTEGER,source_key TEXT);
CREATE TABLE IF NOT EXISTS metadata(key TEXT PRIMARY KEY,value TEXT);
CREATE TABLE IF NOT EXISTS failures(path TEXT PRIMARY KEY,error TEXT);
CREATE VIEW IF NOT EXISTS segments AS SELECT s.*,c.assembly_accession,c.metadata_json AS context_json,f.path AS object_path,f.hash AS object_hash FROM segment_data s JOIN contexts c ON c.id=s.context_id JOIN files f ON f.id=s.file_id;
'''


def open_database(path):
    c=sqlite3.connect(path, timeout=60)
    c.execute('PRAGMA journal_mode=WAL'); c.execute('PRAGMA synchronous=NORMAL'); c.execute('PRAGMA cache_size=-65536')
    c.executescript(DDL)
    return c


class Writer:
    def __init__(self,c):
        self.c=c
        self.contexts={r[1]:r[0] for r in c.execute('SELECT id,metadata_json FROM contexts')}
    def record(self,record,fid,g,rowno,rid=None):
        common={k:v for k,v in record.items() if k not in PER_RECORD}
        encoded=json.dumps(common,separators=(',',':'),sort_keys=True)
        cid=self.contexts.get(encoded)
        if cid is None:
            a=record['assembly_accession']
            cid=self.c.execute('INSERT INTO contexts(assembly_accession,assembly_base,metadata_json) VALUES (?,?,?)',(normalize(a),unversioned(normalize(a)),encoded)).lastrowid
            self.contexts[encoded]=cid
        source=record['source_key']
        if source==record['assembly_accession']+'|'+record['record_name']: source=None
        return (rid,record['record_name'],unversioned(normalize(record['record_name'])),record['segment_start_bp'],record['segment_end_bp'],
                record.get('aligned_bp_length',record['segment_end_bp']),record.get('segment_index',0),record.get('segment_count',1),fid,g,rowno,cid,source)
    def insert(self,rows):
        self.c.executemany('INSERT INTO segment_data VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?)',rows)


def migrate(path):
    if not path.exists(): return
    with closing(sqlite3.connect(path, timeout=60)) as old:
        columns={r[1] for r in old.execute('PRAGMA table_info(segment_data)')}
        if 'context_id' in columns: return
        target=path.with_suffix('.migration.sqlite')
        if target.exists(): target.unlink()
        c=open_database(target); writer=Writer(c)
        c.executemany('INSERT INTO files VALUES (?,?,?,?,?,?,?,?,?)',old.execute('SELECT * FROM files'))
        cursor=old.execute('SELECT id,file_id,row_group,row_in_group,metadata_json FROM segment_data ORDER BY id')
        count=0
        while batch:=cursor.fetchmany(5000):
            writer.insert([writer.record(json.loads(zlib.decompress(r[4])),r[1],r[2],r[3],r[0]) for r in batch])
            count+=len(batch)
            if count%100000==0: c.commit(); print('Compacted',count,'records',flush=True)
        c.commit(); c.execute('PRAGMA wal_checkpoint(TRUNCATE)'); c.execute('PRAGMA journal_mode=DELETE'); c.close()
        old.execute('PRAGMA wal_checkpoint(TRUNCATE)')
    path.rename(path.with_suffix('.legacy.sqlite'))
    target.rename(path)


def finalize_indexes(c):
    c.executescript('''
CREATE INDEX IF NOT EXISTS segment_files ON segment_data(file_id);
CREATE INDEX IF NOT EXISTS record_exact ON segment_data(record_name COLLATE NOCASE);
CREATE INDEX IF NOT EXISTS record_versions ON segment_data(record_base);
CREATE INDEX IF NOT EXISTS segment_context ON segment_data(context_id);
CREATE INDEX IF NOT EXISTS context_assembly ON contexts(assembly_accession);
CREATE INDEX IF NOT EXISTS context_assembly_base ON contexts(assembly_base);
CREATE INDEX IF NOT EXISTS custom_source_key ON segment_data(source_key COLLATE NOCASE) WHERE source_key IS NOT NULL;
''')