رفع وتنزيل الملفات — الملفات الكبيرة والاستيراد الدفعي
رفع الملفات كاستلام طرد — الحزم الصغيرة تُوقَّع فورًا (في الذاكرة)، بينما الشحنات الكبيرة تتطلب تفريغًا على دفعات (كتابة متدفقة)؛ التنزيل كشحن طرد — طرود مفردة (FileResponse)، والشحنات الدفعية تسير على سير متحرك (StreamingResponse).
1. ما ستتعلمه
- معاملات
UploadFileوFile(): الرفع القائم على التدفق مقابل الرفع في الذاكرة - معالجة الملفات الكبيرة: قراءة كتلة بكتلة، كتابة متدفقة بـ
shutil.copyfileobj - تحليل CSV/Excel:
pandas+ التحقق المشترك بنموذج Pydantic - استجابة الملفات: نقاط نهاية تنزيل
FileResponse/StreamingResponse - سيناريو Alice: يرفع Bob ملف CSV يحتوي على ملايين صفوف بيانات الأسعار → تحليل غير متزامن بواسطة Celery في الخلفية → إشعار عند اكتمال الاستيراد
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) ▶ مثال: رفع ملف بسيط
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,
}
الناتج:
# تم تعريف الدالة بنجاح
| الطريقة | استخدام الذاكرة | حجم الملف المناسب | API |
|---|---|---|---|
bytes |
تحميل الكل في الذاكرة | < 2MB | file: bytes = File() |
UploadFile |
ملفات مؤقتة (تدفق) | غير محدود | file: UploadFile = File() |
(2) ▶ مثال: رفع ملفات متعددة
@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}
الناتج:
# تم تعريف الدالة بنجاح
4. المعالجة المتدفقة للملفات الكبيرة
(1) القراءة كتلة بكتلة والكتابة القائمة على التدفق
(1) ▶ مثال: رفع ملفات كبيرة متدفقة
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),
}
الناتج:
# تم تعريف الدالة بنجاح
(2) ▶ مثال: تخصيص حجم الكتلة
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}
الناتج:
# تم تعريف الدالة بنجاح
5. الاستيراد الدفعي CSV/Excel
(1) مسار رفع الملفات الكبيرة + المعالجة غير المتزامنة
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
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.",
}
الناتج:
# تم تعريف الدالة بنجاح
(2) ▶ مثال: مهمة Celery تحلل ملف CSV سطرًا بسطر
# 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),
}
الناتج:
# تم تعريف الدالة بنجاح
(2) مصفوفة معالجة تنسيقات الملفات
| التنسيق | مكتبة التحليل | المميزات | العيوب |
|---|---|---|---|
| CSV | CSV / pandas | خفيف، قائم على التدفق | لا يوجد مشاكل نوع أو ترميز |
| XLSX | openpyxl / pandas | يدعم الأنواع، أوراق متعددة | استخدام ذاكرة مرتفع |
| JSON | json / orjson | مهيكل، متوافق مع Pydantic | حجم ملف كبير |
| Parquet | pyarrow | تخزين عمودي، نسبة ضغط عالية | يتطلب مكتبات إضافية |
6. تنزيل الملفات
(1) FileResponse وStreamingResponse
(1) ▶ مثال: FileResponse — تنزيل ملف
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",
)
الناتج:
# تم تعريف الدالة بنجاح
(2) ▶ مثال: StreamingResponse — توليد CSV متدفق
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"},
)
الناتج:
# تم تعريف الدالة بنجاح
| نوع الاستجابة | حالة الاستخدام | استخدام الذاكرة |
|---|---|---|
FileResponse |
ملفات موجودة | منخفض (تدفق على مستوى نظام التشغيل) |
StreamingResponse |
توليد ديناميكي | منخفض جدًا (يُولَّد سطرًا بسطر) |
❓ أسئلة شائعة
client_max_body_size في Nginx.read_excel() من pandas للتحليل، لكن ملفات XLSX لا يمكن قراءتها بطريقة متدفقة (يجب تحميلها بالكامل). للملفات بملايين الصفوف، يُوصى بتحويلها لتنسيق CSV قبل الرفع.StreamingResponse من النوع async؟async لا يحظر حلقة الأحداث، لذا يُوصى به.file.content_type (مثل text/csv) وامتداد الملف. لاحظ أن content_type يمكن تزويره؛ استخدم هذا فقط للتلميحات في الواجهة الأمامية — لا يزال التحليل والتحقق من جانب الخادم ضروريًا.📖 ملخص
UploadFileيُدفق البيانات إلى ملف مؤقت، لا يستخدم الذاكرة، ومناسب لملفات بأي حجم- استخدم
shutil.copyfileobj()للملفات الكبيرة أوread()للكتابة المتدفقة بكتل إلى القرص - رفع CSV + تحليل غير متزامن بـ Celery: تُعيد API
task_idفورًا؛ العملية الخلفية تتحقق من كل صف وتنفذ إدراجًا دفعيًا - Pydantic يتحقق من بيانات CSV سطرًا بسطر، يتخطَّى ويسجّل الأسطر غير الصالحة، ويستخدم وضعًا متساهلًا لضمان الاستيراد الجزئي
FileResponseلتنزيل ملفات موجودة؛StreamingResponseللتوليد الديناميكي وإرجاع البيانات في الوقت الفعلي
📝 تمارين
- مسألة أساسية (الصعوبة ⭐): نفّذ نقطة نهاية رفع CSV تقبل
UploadFile، تقرأ المحتوى، وتُعيد عدد الصفوف وأسماء الأعمدة. تلميح:file: UploadFile = File(...)+csv.DictReader(io.StringIO(content)) - تمرين متقدم (الصعوبة ⭐⭐): نفّذ مسار رفع CSV → تحليل غير متزامن بـ Celery: احفظ الملف المرفوع في دليل مؤقت، شغّل مهمة Celery، أعد
task_id، واستخدم نقطة نهاية أخرى لاستعلام حالة المهمة ونتائج الاستيراد. تلميح:shutil.copyfileobj+import_prices.delay(file_path) - تحدٍ (الصعوبة: ⭐⭐⭐): نفّذ نقطة نهاية تصدير StreamingResponse: دفّع بيانات الأسعار من قاعدة البيانات، وُلِّد وأعد بيانات CSV سطرًا بسطر، حافظ على استخدام الذاكرة أقل من 1 MB، وادعم التصفية حسب فئة المنتج. تلميح:
async def generate()+yield+StreamingResponse(generate(), media_type="text/csv")
---|



