Download compact_index.py from HuggingFaceBio/carbon-a-database-explorer: direct link, hf CLI and curl.
- Browser
- Download file 4.59 kB
-
https://huggingface.co/spaces/HuggingFaceBio/carbon-a-database-explorer/resolve/main/compact_index.py
- Command line
-
hf download hf://spaces/HuggingFaceBio/carbon-a-database-explorer/compact_index.py
-
curl -L -o compact_index.py https://huggingface.co/spaces/HuggingFaceBio/carbon-a-database-explorer/resolve/main/compact_index.py
4.59 kB
| """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; | |
| ''') | |