#!/usr/bin/env python """ Upsert embeddings to Pinecone indexes. Supports dual-path ingestion to both dense and sparse indexes. Usage: python upsert_db.py # Dense only (default) python upsert_db.py --mode sparse # Sparse only python upsert_db.py --mode both # Both indexes """ import argparse import json import sys from pathlib import Path try: from pinecone_client import batch_upsert_from_file, get_sparse_client import config except ImportError: from .pinecone_client import batch_upsert_from_file, get_sparse_client from . import config def upsert_dense(embeddings_file: Path, namespace: str = "") -> dict: """Upsert to dense index (existing flow).""" print(f"Upserting to dense index: {config.PINECONE_INDEX_NAME}") result = batch_upsert_from_file( str(embeddings_file), batch_size=100, namespace=namespace, show_progress=True ) return result def upsert_sparse(embeddings_file: Path, namespace: str = "") -> dict: """Upsert to sparse index with streaming to avoid memory issues.""" print(f"Upserting to sparse index: {config.SPARSE_INDEX_NAME}") sparse_client = get_sparse_client() batch_size = 96 stored_count = 0 errors = [] batch = [] # Count total lines for progress total_lines = sum(1 for _ in open(embeddings_file, 'r', encoding='utf-8')) print(f"Total chunks: {total_lines}") with open(embeddings_file, 'r', encoding='utf-8') as f: for i, line in enumerate(f): chunk = json.loads(line.strip()) batch.append(chunk) if len(batch) >= batch_size: try: result = sparse_client.store_sparse(batch, batch_size=batch_size, namespace=namespace) stored_count += result.get("stored_count", 0) if result.get("errors"): errors.extend(result["errors"]) print(f" Progress: {stored_count}/{total_lines} ({100*stored_count/total_lines:.1f}%)") except Exception as e: errors.append(f"Batch {i//batch_size} error: {str(e)}") print(f" Error at batch {i//batch_size}: {e}") batch = [] # Final batch if batch: try: result = sparse_client.store_sparse(batch, batch_size=batch_size, namespace=namespace) stored_count += result.get("stored_count", 0) if result.get("errors"): errors.extend(result["errors"]) except Exception as e: errors.append(f"Final batch error: {str(e)}") return { "stored_count": stored_count, "total_chunks": total_lines, "errors": errors if errors else None } def main(): parser = argparse.ArgumentParser(description="Upsert embeddings to Pinecone") parser.add_argument( "--mode", choices=["dense", "sparse", "both"], default="dense", help="Which indexes to upsert to" ) parser.add_argument( "--namespace", default="", help="Namespace for vectors" ) parser.add_argument( "--file", type=str, default=None, help="Path to embeddings JSONL file" ) args = parser.parse_args() script_dir = Path(__file__).parent project_root = script_dir.parent if args.file: embeddings_file = Path(args.file) else: embeddings_file = project_root / "03_embedding" / "output" / "test_text" / "embeddings_mtsamples.jsonl" if not embeddings_file.exists(): embeddings_file = project_root / "03_embedding" / "output" / "embeddings_mtsamples.jsonl" if not embeddings_file.exists(): print(f"Error: Embeddings file not found: {embeddings_file}") sys.exit(1) print(f"Reading embeddings from: {embeddings_file}") print(f"Mode: {args.mode}") results = {} if args.mode in ["dense", "both"]: try: result = upsert_dense(embeddings_file, args.namespace) results["dense"] = result print(f"\n[Dense] Stored: {result['stored_count']}") if result.get('errors'): print(f"[Dense] Errors: {len(result['errors'])}") except Exception as e: print(f"[Dense] Error: {e}") results["dense"] = {"error": str(e)} if args.mode in ["sparse", "both"]: try: result = upsert_sparse(embeddings_file, args.namespace) results["sparse"] = result print(f"\n[Sparse] Stored: {result['stored_count']}") if result.get('errors'): print(f"[Sparse] Errors: {len(result['errors'])}") except Exception as e: print(f"[Sparse] Error: {e}") results["sparse"] = {"error": str(e)} print("\nUpsert complete.") return results if __name__ == "__main__": main()