404 Not Found

404 Not Found


nginx

رفع وتنزيل الملفات — الملفات الكبيرة والاستيراد الدفعي

رفع الملفات كاستلام طرد — الحزم الصغيرة تُوقَّع فورًا (في الذاكرة)، بينما الشحنات الكبيرة تتطلب تفريغًا على دفعات (كتابة متدفقة)؛ التنزيل كشحن طرد — طرود مفردة (FileResponse)، والشحنات الدفعية تسير على سير متحرك (StreamingResponse).

1. ما ستتعلمه


2. القصة الحقيقية لـ Alice

(1) المشكلة: خطأ نفاد الذاكرة عند استيراد ملف CSV بمليون صف

يرفع Bob ملف CSV يحتوي على مليون صف من بيانات الأسعار إلى PriceTracker كل شهر. التنفيذ السابق كان يقرأ الملف بالكامل في الذاكرة قبل تحليله، وتسبب ملف CSV بحجم 1GB في تعطّل عملية Python بسبب خطأ OOM (نفاد الذاكرة). وللأسوأ، كان CSV يحتوي على بيانات غير صالحة (أسعار سالبة، عملات غير صالحة)، وبعد فشل الاستيراد في منتصف الطريق، بقيت قاعدة البيانات في حالة عدم تناسق.

(2) حل باستخدام الرفع القائم على التدفق والمعالجة غير المتزامنة

يحفظ UploadFile في FastAPI الملفات كملفات مؤقتة افتراضيًا (لا تستهلك الذاكرة). مدمجًا مع Celery للمعالجة غير المتزامنة لملايين الأسطر من البيانات، وPydantic للتحقق سطرًا بسطر لتجاوز البيانات غير الصالحة، وإدراج دفعي تعاملي لضمان التناسق.

(3) العائد

انتقل استيراد CSV بحجم 1GB من التعطل بسبب خطأ OOM إلى العمل باستقرار، بانخفاض استخدام الذاكرة من 2GB إلى 50MB. يعالج Celery Worker مليون صف في حوالي 10 دقائق، ويمكن لـ Bob عرض التقدم في الوقت الفعلي عبر نقطة نهاية حالة المهمة.


3. أساسيات UploadFile

(1) UploadFile مقابل bytes

(1) ▶ مثال: رفع ملف بسيط

PYTHON
from fastapi import FastAPI, UploadFile, File, HTTPException

app = FastAPI()

@app.post("/upload/single")
async def upload_single(file: UploadFile = File(...)):
    # UploadFile: ملف محفوظ كملف مؤقت، وليس في الذاكرة
    content = await file.read()
    return {
        "filename": file.filename,
        "size": len(content),
        "content_type": file.content_type,
    }

الناتج:

TEXT
# تم تعريف الدالة بنجاح
الطريقة استخدام الذاكرة حجم الملف المناسب API
bytes تحميل الكل في الذاكرة < 2MB file: bytes = File()
UploadFile ملفات مؤقتة (تدفق) غير محدود file: UploadFile = File()

(2) ▶ مثال: رفع ملفات متعددة

PYTHON
@app.post("/upload/multiple")
async def upload_multiple(files: list[UploadFile] = File(...)):
    results = []
    for file in files:
        content = await file.read()
        results.append({
            "filename": file.filename,
            "size": len(content),
        })
    return {"uploaded": len(results), "files": results}

الناتج:

TEXT
# تم تعريف الدالة بنجاح

4. المعالجة المتدفقة للملفات الكبيرة

(1) القراءة كتلة بكتلة والكتابة القائمة على التدفق

(1) ▶ مثال: رفع ملفات كبيرة متدفقة

PYTHON
import shutil
from pathlib import Path
from fastapi import FastAPI, UploadFile, File

app = FastAPI()
UPLOAD_DIR = Path("uploads")
UPLOAD_DIR.mkdir(exist_ok=True)

@app.post("/upload/large")
async def upload_large_file(file: UploadFile = File(...)):
    # تدفق الملف إلى القرص - لا يُحمّل الملف بالكامل في الذاكرة أبدًا
    dest = UPLOAD_DIR / file.filename
    with open(dest, "wb") as buffer:
        # نسخ بكتل (مخزن مؤقت 64KB افتراضي)
        shutil.copyfileobj(file.file, buffer)
    
    file_size = dest.stat().st_size
    return {
        "filename": file.filename,
        "size_bytes": file_size,
        "saved_to": str(dest),
    }

الناتج:

TEXT
# تم تعريف الدالة بنجاح

(2) ▶ مثال: تخصيص حجم الكتلة

PYTHON
CHUNK_SIZE = 1024 * 1024  # كتل 1MB

@app.post("/upload/chunked")
async def upload_chunked(file: UploadFile = File(...)):
    dest = UPLOAD_DIR / file.filename
    bytes_written = 0
    with open(dest, "wb") as buffer:
        while chunk := await file.read(CHUNK_SIZE):
            buffer.write(chunk)
            bytes_written += len(chunk)
    return {"filename": file.filename, "bytes_written": bytes_written}

الناتج:

TEXT
# تم تعريف الدالة بنجاح

5. الاستيراد الدفعي CSV/Excel

(1) مسار رفع الملفات الكبيرة + المعالجة غير المتزامنة

100%
sequenceDiagram
    participant Bob as Bob الواجهة الأمامية
    participant API as FastAPI
    participant Temp as ملف مؤقت
    participant Celery as Celery Worker
    participant DB as PostgreSQL

    Bob->>API: POST /import/csv (UploadFile)
    API->>Temp: حفظ في ملف مؤقت
    API-->>Bob: 202 Accepted + task_id
    API->>Celery: تشغيل مهمة التحليل
    Celery->>Temp: قراءة CSV بكتل
    Celery->>Celery: التحقق من كل صف (Pydantic)
    Celery->>DB: إدراج دفعي للصفوف الصالحة
    Celery-->>Celery: الإبلاغ عن التقدم
    Bob->>API: GET /tasks/{task_id}
    API-->>Bob: التقدم: 75%
    Celery-->>Celery: اكتملت المهمة
    Bob->>API: GET /tasks/{task_id}
    API-->>Bob: الحالة: SUCCESS، تم استيراد: 950000

(1) ▶ مثال: نقطة نهاية رفع CSV + تشغيل مهمة Celery

PYTHON
import csv
import io
from fastapi import FastAPI, UploadFile, File, Depends
from app.tasks import import_prices_from_csv

app = FastAPI()
UPLOAD_DIR = Path("uploads")

@app.post("/api/v1/import/csv")
async def import_csv(
    file: UploadFile = File(..., description="ملف CSV ببيانات الأسعار"),
    user=Depends(require_subscription("pro")),
):
    if not file.filename.endswith(".csv"):
        raise HTTPException(status_code=400, detail="Only CSV files accepted")
    
    # حفظ الملف المرفوع
    dest = UPLOAD_DIR / f"{uuid4()}.csv"
    with open(dest, "wb") as buffer:
        shutil.copyfileobj(file.file, buffer)
    
    # تشغيل مهمة Celery للمعالجة غير المتزامنة
    task = import_prices_from_csv.delay(str(dest), user_id=user.id)
    
    return {
        "task_id": task.id,
        "filename": file.filename,
        "status": "processing",
        "message": "File uploaded. Check task status for progress.",
    }

الناتج:

TEXT
# تم تعريف الدالة بنجاح

(2) ▶ مثال: مهمة Celery تحلل ملف CSV سطرًا بسطر

PYTHON
# app/tasks/import_tasks.py
from app.core.celery_app import celery_app
from pydantic import BaseModel, Field, field_validator
import csv

class PriceRow(BaseModel):
    product_id: int = Field(gt=0)
    price: float = Field(gt=0)
    currency: str = Field(default="USD", pattern=r"^[A-Z]{3}$")
    source: str = Field(max_length=100, default="csv_import")

    @field_validator("price")
    @classmethod
    def round_price(cls, v: float) -> float:
        return round(v, 2)

@celery_app.task(bind=True)
def import_prices_from_csv(self, file_path: str, user_id: int):
    valid_rows = []
    invalid_rows = []
    total_rows = 0
    
    with open(file_path, "r") as f:
        reader = csv.DictReader(f)
        for row in reader:
            total_rows += 1
            try:
                validated = PriceRow(**row)
                valid_rows.append(validated.model_dump())
            except Exception as e:
                invalid_rows.append({"row": total_rows, "error": str(e)})
            
            # الإبلاغ عن التقدم كل 10000 صف
            if total_rows % 10000 == 0:
                self.update_state(
                    state="PROGRESS",
                    meta={"current": total_rows, "valid": len(valid_rows), "invalid": len(invalid_rows)},
                )
    
    # إدراج دفعي للصفوف الصالحة
    batch_insert_prices(valid_rows, batch_size=5000)
    
    # تنظيف الملف المؤقت
    Path(file_path).unlink(missing_ok=True)
    
    return {
        "total": total_rows,
        "imported": len(valid_rows),
        "skipped": len(invalid_rows),
    }

الناتج:

TEXT
# تم تعريف الدالة بنجاح

(2) مصفوفة معالجة تنسيقات الملفات

التنسيق مكتبة التحليل المميزات العيوب
CSV CSV / pandas خفيف، قائم على التدفق لا يوجد مشاكل نوع أو ترميز
XLSX openpyxl / pandas يدعم الأنواع، أوراق متعددة استخدام ذاكرة مرتفع
JSON json / orjson مهيكل، متوافق مع Pydantic حجم ملف كبير
Parquet pyarrow تخزين عمودي، نسبة ضغط عالية يتطلب مكتبات إضافية

6. تنزيل الملفات

(1) FileResponse وStreamingResponse

(1) ▶ مثال: FileResponse — تنزيل ملف

PYTHON
from fastapi import FastAPI
from fastapi.responses import FileResponse
from pathlib import Path

app = FastAPI()

@app.get("/download/prices/csv")
async def download_prices_csv(
    user=Depends(require_subscription("pro")),
):
    # توليد ملف CSV (أو استخدام ملف مُولَّد مسبقًا)
    file_path = Path("exports/prices.csv")
    return FileResponse(
        path=file_path,
        filename="price_data.csv",
        media_type="text/csv",
    )

الناتج:

TEXT
# تم تعريف الدالة بنجاح

(2) ▶ مثال: StreamingResponse — توليد CSV متدفق

PYTHON
from fastapi.responses import StreamingResponse
import csv
import io
from app.core.deps import get_db

@app.get("/api/v1/export/prices")
async def export_prices(
    category: str | None = None,
    db: AsyncSession = Depends(get_db),
    user=Depends(require_subscription("pro")),
):
    async def generate_csv():
        output = io.StringIO()
        writer = csv.writer(output)
        writer.writerow(["product_id", "name", "price", "currency", "recorded_at"])
        yield output.getvalue()
        output.seek(0)
        output.truncate(0)
        
        # تدفق الصفوف على دفعات
        offset = 0
        batch_size = 5000
        while True:
            rows = await fetch_price_batch(db, category, offset, batch_size)
            if not rows:
                break
            for row in rows:
                writer.writerow([
                    row.product_id, row.name,
                    row.price, row.currency, row.recorded_at,
                ])
                yield output.getvalue()
                output.seek(0)
                output.truncate(0)
            offset += batch_size
    
    return StreamingResponse(
        generate_csv(),
        media_type="text/csv",
        headers={"Content-Disposition": "attachment; filename=prices.csv"},
    )

الناتج:

TEXT
# تم تعريف الدالة بنجاح
نوع الاستجابة حالة الاستخدام استخدام الذاكرة
FileResponse ملفات موجودة منخفض (تدفق على مستوى نظام التشغيل)
StreamingResponse توليد ديناميكي منخفض جدًا (يُولَّد سطرًا بسطر)

❓ أسئلة شائعة

س متى تُحذف ملفات UploadFile المؤقتة؟
ج تُحذف تلقائيًا بعد انتهاء الطلب. إذا كنت بحاجة للاحتفاظ بها، انسخها إلى دليل دائم أثناء عملية الرفع.
س هل يوجد حد لحجم ملف الرفع؟
ج FastAPI لا يفرض حدودًا على مستوى الإطار، لكن Uvicorn يفرض حدًا افتراضيًا على حجم جسم الطلب. في بيئة الإنتاج، يتحكم في ذلك client_max_body_size في Nginx.
س ماذا أفعل إذا كانت بعض الصفوف غير صالحة أثناء استيراد CSV؟
ج يعتمد النهج على احتياجات عملك — الوضع الصارم (تراجع جميع الصفوف غير الصالحة) أو الوضع المتساهل (تخطَّ الصفوف غير الصالحة واستورد الصالحة). يستخدم PriceTracker الوضع المتساهل ويُعيد عدد الصفوف المستوردة وعدد الصفوف المتخطاة.
س كيف يجب التعامل مع ملفات Excel؟
ج استخدم read_excel() من pandas للتحليل، لكن ملفات XLSX لا يمكن قراءتها بطريقة متدفقة (يجب تحميلها بالكامل). للملفات بملايين الصفوف، يُوصى بتحويلها لتنسيق CSV قبل الرفع.
س هل يجب أن يكون المُولِّد لـ StreamingResponse من النوع async؟
ج لا، ليس شرطًا؛ المُولِّد المتزامن مقبول أيضًا. لكن المُولِّد async لا يحظر حلقة الأحداث، لذا يُوصى به.
س كيف يمكنني تقييد أنواع الملفات القابلة للرفع؟
ج تحقق من file.content_type (مثل text/csv) وامتداد الملف. لاحظ أن content_type يمكن تزويره؛ استخدم هذا فقط للتلميحات في الواجهة الأمامية — لا يزال التحليل والتحقق من جانب الخادم ضروريًا.

📖 ملخص


📝 تمارين

  1. مسألة أساسية (الصعوبة ⭐): نفّذ نقطة نهاية رفع CSV تقبل UploadFile، تقرأ المحتوى، وتُعيد عدد الصفوف وأسماء الأعمدة. تلميح: file: UploadFile = File(...) + csv.DictReader(io.StringIO(content))
  2. تمرين متقدم (الصعوبة ⭐⭐): نفّذ مسار رفع CSV → تحليل غير متزامن بـ Celery: احفظ الملف المرفوع في دليل مؤقت، شغّل مهمة Celery، أعد task_id، واستخدم نقطة نهاية أخرى لاستعلام حالة المهمة ونتائج الاستيراد. تلميح: shutil.copyfileobj + import_prices.delay(file_path)
  3. تحدٍ (الصعوبة: ⭐⭐⭐): نفّذ نقطة نهاية تصدير StreamingResponse: دفّع بيانات الأسعار من قاعدة البيانات، وُلِّد وأعد بيانات CSV سطرًا بسطر، حافظ على استخدام الذاكرة أقل من 1 MB، وادعم التصفية حسب فئة المنتج. تلميح: async def generate() + yield + StreamingResponse(generate(), media_type="text/csv")

---|

Web-Tutorial.com

فريق Web-Tutorial التقني

منصة دروس برمجية يديرها عدة مطورين. كل درس يتم كتابته ومراجعته بواسطة مطورين متخصصين في المجال. نعمل على ضمان دقة وموثوقية المحتوى — إذا لاحظت أي مشكلة، فيرجى إخبارنا.

100%