eludius18/json-processor
0
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 