Team Ai
Apppublic

eludius18/json-processor

sourceHugging Faceupdated 1y agoView on Hugging Face
0likes
main.py612 linesDownload Raw Back to root
1"""2JSON Massive Processor - Hugging Face Space Backend3Procesamiento asíncrono de archivos JSON masivos (1M+ registros)4 5@author Eladio Robles Casas6@version 1.07"""8 9from fastapi import FastAPI, HTTPException, BackgroundTasks10from fastapi.middleware.cors import CORSMiddleware11from fastapi.staticfiles import StaticFiles12from fastapi.responses import FileResponse13from pydantic import BaseModel14from typing import List, Optional, Dict, Any15import json16import os17import time18from datetime import datetime19import csv20import logging21from pathlib import Path22 23# ==================== CONFIGURACIÓN ====================24 25logging.basicConfig(26    level=logging.INFO,27    format='%(asctime)s - %(name)s - %(levelname)s - %(message)s'28)29logger = logging.getLogger(__name__)30 31app = FastAPI(32    title="JSON Massive Processor",33    description="Procesamiento asíncrono de archivos JSON masivos",34    version="1.0.0"35)36 37# CORS - permite acceso desde n8n Cloud38app.add_middleware(39    CORSMiddleware,40    allow_origins=["*"],  # En producción: especificar dominios41    allow_credentials=True,42    allow_methods=["*"],43    allow_headers=["*"],44)45 46# Crear directorios necesarios47Path("static/csv").mkdir(parents=True, exist_ok=True)48Path("jobs").mkdir(parents=True, exist_ok=True)49Path("data").mkdir(parents=True, exist_ok=True)50 51# Montar archivos estáticos52app.mount("/static", StaticFiles(directory="static"), name="static")53 54# ==================== MODELOS ====================55 56class StartAsyncRequest(BaseModel):57    """Request para iniciar procesamiento asíncrono"""58    mode: str = "local_jsonl"  # local_jsonl, local_json_array, url_json59    chunk_size: int = 2500060    local_path: str = "/home/user/app/data/massive_500k_records.json"61    json_pointer: str = ""62    header_order: List[str] = [63        "id", "name", "firstName", "street", "countryCode", "zipCode"64    ]65 66class JobStatus(BaseModel):67    """Estado de un job"""68    ok: bool69    job_id: str70    status: str  # pending, processing, completed, failed71    progress: float72    chunks_processed: int73    total_chunks: int74    rows_processed: int75    estimated_total: Optional[int] = None76    started_at: str77    updated_at: str78    completed_at: Optional[str] = None79    error: Optional[str] = None80    error_type: Optional[str] = None81 82# ==================== HELPERS ====================83 84def generate_job_id() -> str:85    """Genera un job_id único"""86    timestamp = int(time.time())87    return f"job_{timestamp}_{os.urandom(4).hex()}"88 89def get_state_path(job_id: str) -> str:90    """Retorna el path del archivo de estado"""91    return f"jobs/{job_id}_state.json"92 93def read_state(job_id: str) -> Dict[str, Any]:94    """Lee el estado de un job"""95    state_path = get_state_path(job_id)96    if not os.path.exists(state_path):97        raise FileNotFoundError(f"Job not found: {job_id}")98    99    with open(state_path, 'r') as f:100        return json.load(f)101 102def write_state(job_id: str, state: Dict[str, Any]):103    """Escribe el estado de un job"""104    state_path = get_state_path(job_id)105    state["updated_at"] = datetime.utcnow().isoformat()106    107    with open(state_path, 'w') as f:108        json.dump(state, f, indent=2)109    110    logger.info(f"Job {job_id}: State updated - {state.get('status')}")111 112def extract_contract_data(contract: Dict[str, Any]) -> Dict[str, str]:113    """114    Extrae datos de un contrato según la estructura de massive_500k_records.json115    116    Estructura esperada:117    {118      "data": {119        "contracts": [120          {121            "contractHolder": {122              "name": "...",123              "firstName": "..."124            },125            "insuredProperties": [126              {127                "property": {128                  "address": {129                    "street": "...",130                    "countryCode": "...",131                    "zipCode": "..."132                  }133                }134              }135            ]136          }137        ]138      }139    }140    """141    holder = contract.get('contractHolder', {})142    143    # Si hay insuredProperties, extraer del primero144    insured_props = contract.get('insuredProperties', [])145    address = {}146    147    if insured_props:148        first_prop = insured_props[0]149        property_data = first_prop.get('property', {})150        address = property_data.get('address', {})151    152    return {153        'id': contract.get('id', ''),154        'name': holder.get('name', ''),155        'firstName': holder.get('firstName', ''),156        'street': address.get('street', ''),157        'countryCode': address.get('countryCode', ''),158        'zipCode': address.get('zipCode', '')159    }160 161# ==================== BACKGROUND TASK ====================162 163def process_jsonl_to_csv_background(164    job_id: str,165    local_path: str,166    chunk_size: int,167    header_order: List[str]168):169    """170    Background task que procesa JSONL → CSV directamente171    Lee el archivo línea a línea (JSONL/NDJSON format)172    """173    try:174        logger.info(f"Job {job_id}: Starting background processing")175        176        # Inicializar estado177        state = {178            "job_id": job_id,179            "status": "processing",180            "progress": 0.0,181            "chunks_processed": 0,182            "total_chunks": 0,183            "rows_processed": 0,184            "started_at": datetime.utcnow().isoformat(),185            "updated_at": datetime.utcnow().isoformat()186        }187        write_state(job_id, state)188        189        # Preparar CSV190        csv_path = f"static/csv/{job_id}_final.csv"191        192        # Contar total de líneas para progress (opcional, puede ser lento)193        logger.info(f"Job {job_id}: Counting total lines...")194        total_lines = 0195        try:196            with open(local_path, 'r', encoding='utf-8') as f:197                for _ in f:198                    total_lines += 1199        except Exception as e:200            logger.warning(f"Job {job_id}: Could not count lines: {e}")201            total_lines = 0  # No sabemos el total202        203        logger.info(f"Job {job_id}: Total lines to process: {total_lines}")204        state["estimated_total"] = total_lines205        write_state(job_id, state)206        207        # Procesar JSONL streaming208        chunk_index = 0209        rows_processed = 0210        rows_buffer = []211        212        with open(csv_path, 'w', newline='', encoding='utf-8') as csv_file:213            writer = csv.DictWriter(csv_file, fieldnames=header_order)214            writer.writeheader()215            216            with open(local_path, 'r', encoding='utf-8') as jsonl_file:217                for line_num, line in enumerate(jsonl_file, 1):218                    line = line.strip()219                    if not line:220                        continue221                    222                    try:223                        # Parsear línea JSON224                        obj = json.loads(line)225                        226                        # Extraer contratos227                        contracts = obj.get('data', {}).get('contracts', [])228                        229                        # Procesar cada contrato230                        for contract in contracts:231                            extracted = extract_contract_data(contract)232                            233                            # Asegurar que tiene todos los campos234                            row = {key: extracted.get(key, '') for key in header_order}235                            rows_buffer.append(row)236                            rows_processed += 1237                            238                            # Escribir cuando llegamos a chunk_size239                            if len(rows_buffer) >= chunk_size:240                                for buffered_row in rows_buffer:241                                    writer.writerow(buffered_row)242                                243                                chunk_index += 1244                                245                                # Actualizar progreso246                                if total_lines > 0:247                                    progress = (line_num / total_lines) * 100248                                else:249                                    progress = 0.0250                                251                                state.update({252                                    "status": "processing",253                                    "progress": round(progress, 2),254                                    "chunks_processed": chunk_index,255                                    "rows_processed": rows_processed256                                })257                                write_state(job_id, state)258                                259                                logger.info(260                                    f"Job {job_id}: Chunk {chunk_index} - "261                                    f"{rows_processed} rows ({progress:.2f}%)"262                                )263                                264                                rows_buffer = []265                        266                    except json.JSONDecodeError as e:267                        logger.error(f"Job {job_id}: JSON parse error at line {line_num}: {e}")268                        continue269                    except Exception as e:270                        logger.error(f"Job {job_id}: Error processing line {line_num}: {e}")271                        continue272                273                # Escribir último buffer274                if rows_buffer:275                    for buffered_row in rows_buffer:276                        writer.writerow(buffered_row)277                    chunk_index += 1278        279        # Calcular tamaño del CSV280        csv_size_mb = os.path.getsize(csv_path) / 1024 / 1024281        282        # Marcar como completado283        state.update({284            "status": "completed",285            "progress": 100.0,286            "chunks_processed": chunk_index,287            "total_chunks": chunk_index,288            "rows_processed": rows_processed,289            "completed_at": datetime.utcnow().isoformat(),290            "csv_path": csv_path,291            "csv_size_mb": round(csv_size_mb, 2)292        })293        write_state(job_id, state)294        295        logger.info(f"Job {job_id}: ✅ Completed - {rows_processed} rows processed")296        297    except FileNotFoundError as e:298        logger.error(f"Job {job_id}: File not found: {e}")299        state = read_state(job_id)300        state.update({301            "status": "failed",302            "error": f"File not found: {str(e)}",303            "error_type": "FileNotFoundError"304        })305        write_state(job_id, state)306        307    except Exception as e:308        logger.error(f"Job {job_id}: ❌ Failed: {e}", exc_info=True)309        try:310            state = read_state(job_id)311            state.update({312                "status": "failed",313                "error": str(e),314                "error_type": type(e).__name__315            })316            write_state(job_id, state)317        except:318            pass319 320# ==================== ENDPOINTS ====================321 322@app.get("/")323async def root():324    """Endpoint raíz"""325    return {326        "service": "JSON Massive Processor",327        "version": "1.0.0",328        "status": "running",329        "endpoints": {330            "health": "/health",331            "start_job": "POST /start_async_job",332            "check_status": "GET /status/{job_id}",333            "download": "GET /download/{job_id}"334        }335    }336 337@app.get("/health")338async def health_check():339    """Health check endpoint"""340    return {341        "status": "healthy",342        "timestamp": datetime.utcnow().isoformat(),343        "service": "JSON Massive Processor"344    }345 346@app.post("/start_async_job")347async def start_async_job(req: StartAsyncRequest, background_tasks: BackgroundTasks):348    """349    Inicia procesamiento asíncrono de JSON/JSONL masivo350    351    Devuelve job_id inmediatamente y procesa en background352    """353    try:354        # Validar que existe el archivo355        if not os.path.exists(req.local_path):356            raise HTTPException(357                status_code=404,358                detail=f"File not found: {req.local_path}"359            )360        361        # Generar job_id362        job_id = generate_job_id()363        364        logger.info(f"Job {job_id}: Starting async job for {req.local_path}")365        366        # Crear estado inicial367        initial_state = {368            "job_id": job_id,369            "status": "pending",370            "progress": 0.0,371            "chunks_processed": 0,372            "total_chunks": 0,373            "rows_processed": 0,374            "started_at": datetime.utcnow().isoformat(),375            "updated_at": datetime.utcnow().isoformat(),376            "config": {377                "mode": req.mode,378                "chunk_size": req.chunk_size,379                "local_path": req.local_path,380                "header_order": req.header_order381            }382        }383        write_state(job_id, initial_state)384        385        # Lanzar background task386        background_tasks.add_task(387            process_jsonl_to_csv_background,388            job_id=job_id,389            local_path=req.local_path,390            chunk_size=req.chunk_size,391            header_order=req.header_order392        )393        394        return {395            "ok": True,396            "job_id": job_id,397            "status": "pending",398            "message": "Job queued for processing"399        }400        401    except HTTPException:402        raise403    except Exception as e:404        logger.error(f"Error starting async job: {e}", exc_info=True)405        raise HTTPException(status_code=500, detail=str(e))406 407@app.get("/status/{job_id}")408async def get_status(job_id: str):409    """410    Consulta el estado actual de un job411    """412    try:413        state = read_state(job_id)414        415        return {416            "ok": True,417            "job_id": state.get("job_id"),418            "status": state.get("status"),419            "progress": state.get("progress", 0.0),420            "chunks_processed": state.get("chunks_processed", 0),421            "total_chunks": state.get("total_chunks", 0),422            "rows_processed": state.get("rows_processed", 0),423            "estimated_total": state.get("estimated_total"),424            "started_at": state.get("started_at"),425            "updated_at": state.get("updated_at"),426            "completed_at": state.get("completed_at"),427            "error": state.get("error"),428            "error_type": state.get("error_type")429        }430        431    except FileNotFoundError:432        raise HTTPException(433            status_code=404,434            detail=f"Job not found: {job_id}"435        )436    except Exception as e:437        logger.error(f"Error getting status for {job_id}: {e}")438        raise HTTPException(status_code=500, detail=str(e))439 440@app.get("/download/{job_id}")441async def download_csv(job_id: str):442    """443    Descarga el CSV final (solo si status=completed)444    """445    try:446        state = read_state(job_id)447        448        # Verificar que está completado449        if state.get("status") != "completed":450            return {451                "ok": False,452                "error": f"Job not completed yet (status: {state.get('status')})",453                "current_status": state.get("status"),454                "progress": state.get("progress", 0)455            }456        457        csv_path = state.get("csv_path")458        if not csv_path or not os.path.exists(csv_path):459            raise HTTPException(460                status_code=404,461                detail="CSV file not found"462            )463        464        # Generar nombre de archivo descriptivo465        filename = f"massive_data_{job_id}_{state.get('rows_processed', 0)}_rows.csv"466        467        return FileResponse(468            csv_path,469            media_type="text/csv",470            headers={471                "Content-Disposition": f'attachment; filename="{filename}"'472            }473        )474        475    except HTTPException:476        raise477    except FileNotFoundError:478        raise HTTPException(479            status_code=404,480            detail=f"Job not found: {job_id}"481        )482    except Exception as e:483        logger.error(f"Error downloading CSV for {job_id}: {e}")484        raise HTTPException(status_code=500, detail=str(e))485 486@app.delete("/job/{job_id}")487async def delete_job(job_id: str):488    """489    Elimina un job y sus archivos asociados (limpieza)490    """491    try:492        state = read_state(job_id)493        494        # Eliminar CSV si existe495        csv_path = state.get("csv_path")496        if csv_path and os.path.exists(csv_path):497            os.remove(csv_path)498            logger.info(f"Job {job_id}: CSV deleted")499        500        # Eliminar estado501        state_path = get_state_path(job_id)502        if os.path.exists(state_path):503            os.remove(state_path)504            logger.info(f"Job {job_id}: State deleted")505        506        return {507            "ok": True,508            "message": f"Job {job_id} deleted successfully"509        }510        511    except FileNotFoundError:512        raise HTTPException(513            status_code=404,514            detail=f"Job not found: {job_id}"515        )516    except Exception as e:517        logger.error(f"Error deleting job {job_id}: {e}")518        raise HTTPException(status_code=500, detail=str(e))519 520@app.get("/jobs")521async def list_jobs():522    """523    Lista todos los jobs existentes524    """525    try:526        jobs = []527        jobs_dir = Path("jobs")528        529        for state_file in jobs_dir.glob("*_state.json"):530            try:531                with open(state_file, 'r') as f:532                    state = json.load(f)533                    jobs.append({534                        "job_id": state.get("job_id"),535                        "status": state.get("status"),536                        "progress": state.get("progress"),537                        "rows_processed": state.get("rows_processed"),538                        "started_at": state.get("started_at"),539                        "completed_at": state.get("completed_at")540                    })541            except Exception as e:542                logger.error(f"Error reading {state_file}: {e}")543                continue544        545        return {546            "ok": True,547            "total_jobs": len(jobs),548            "jobs": jobs549        }550        551    except Exception as e:552        logger.error(f"Error listing jobs: {e}")553        raise HTTPException(status_code=500, detail=str(e))554 555# ==================== DEBUG ENDPOINTS ====================556 557@app.get("/debug/files")558async def debug_list_files():559    """Lista archivos disponibles en data/"""560    try:561        data_dir = Path("data")562        files = []563        564        if data_dir.exists():565            for file_path in data_dir.iterdir():566                if file_path.is_file():567                    size_mb = file_path.stat().st_size / 1024 / 1024568                    files.append({569                        "name": file_path.name,570                        "size_mb": round(size_mb, 2),571                        "path": str(file_path)572                    })573        574        return {575            "ok": True,576            "data_directory": str(data_dir.absolute()),577            "files": files,578            "total_files": len(files)579        }580    except Exception as e:581        return {582            "ok": False,583            "error": str(e)584        }585 586# ==================== STARTUP ====================587 588@app.on_event("startup")589async def startup_event():590    logger.info("🚀 JSON Massive Processor started")591    logger.info(f"📁 Data directory: {os.path.abspath('data')}")592    logger.info(f"📁 CSV output directory: {os.path.abspath('static/csv')}")593    logger.info(f"📁 Jobs directory: {os.path.abspath('jobs')}")594    595    # List available files596    try:597        data_dir = Path("data")598        if data_dir.exists():599            files = list(data_dir.iterdir())600            logger.info(f"📂 Files in data/: {[f.name for f in files if f.is_file()]}")601    except Exception as e:602        logger.error(f"Error listing data files: {e}")603 604@app.on_event("shutdown")605async def shutdown_event():606    logger.info("👋 JSON Massive Processor shutting down")607 608if __name__ == "__main__":609    import uvicorn610    uvicorn.run(app, host="0.0.0.0", port=7860)611 612