404 Not Found

404 Not Found


nginx

ファイルアップロードとダウンロード — 大容量ファイルと一括インポート

ファイルのアップロードは荷物の受け取りのようなものです。小包は即座にサインして受け取ります (メモリ), 大量の荷物は分割して荷下ろしします (ストリーム書き込み)。ダウンロードは荷物の発送のようなものです。単品はFileResponse, 大量発送はコンベアベルトに乗せます (StreamingResponse)。

1. 学ぶ内容


2. Aliceのリアルストーリー

(1) ペインポイント:100万行のCSVインポートでメモリ不足エラー

Bobは毎月百万行の価格データを含むCSVファイルをPriceTrackerにアップロードします。以前の実装ではファイル全体をメモリに読み込んでから解析しており, 1GBのCSVファイルがPythonプロセスをOOM (Out of Memory)エラーでクラッシュさせていました。さらに悪いことに, CSVに不正データ (負の価格, 無効な通貨)が含まれており, インポートが途中で失敗した後, データベースが不整合な状態に残りました。

(2) ストリームベースアップロードと非同期処理のソリューション

FastAPIのUploadFileはデフォルトでファイルを一時ファイルとして保存します (メモリを消費しません)。Celeryと組み合わせて百万行のデータを非同期処理し, Pydanticで行単位の検証を行って不正データをスキップし, トランザクショナルなバッチ挿入で整合性を確保します。

(3) 収益

1GBのCSVインポートがOOMエラーによるクラッシュから安定動作に変わり, メモリ使用量が2GBから50MBに削減されました。Celery Workerが100万行を約10分で処理し, Bobはタスクステータスエンドポイントでリアルタイムに進捗を確認できます。


3. UploadFileの基礎

(1) UploadFile vs. 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 Frontend
    participant API as FastAPI
    participant Temp as Temp File
    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: Progress: 75%
    Celery-->>Celery: タスク完了
    Bob->>API: GET /tasks/{task_id}
    API-->>Bob: Status: SUCCESS, imported: 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 file with price data"),
    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) ▶サンプル:CSVファイルを行単位で解析するCeleryタスク

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 既存ファイル 低 (OSレベルのストリーミング)
StreamingResponse 動的生成 非常に低い (行単位で生成)

❓ よくある質問

Q UploadFileの一時ファイルはいつ削除されますか?
A リクエスト終了後に自動的に削除されます。保持する必要がある場合は, アップロード中に永続ディレクトリにコピーしてください。
Q アップロードのファイルサイズに制限はありますか?
A FastAPI自体には制限がありませんが, Uvicornがリクエストボディサイズにデフォルト制限を設けています。本番環境ではNginxのclient_max_body_sizeで制御します。
Q CSVインポートで一部の行が不正な場合はどうすればよいですか?
A ビジネス要件によります。厳格モード (不正行があれば全てロールバック)か寛容モード (不正行をスキップして有効行のみインポート)かです。PriceTrackerは寛容モードを使用し, インポート件数とスキップ件数を返します。
Q Excelファイルはどのように処理すべきですか?
A pandasのread_excel()で解析しますが, XLSXファイルはストリーミング読み込みができません (全体を読み込む必要があります)。百万行規模のファイルはCSV形式に変換してからアップロードすることを推奨します。
Q StreamingResponseのジェネレータはasyncでなければなりませんか?
A いいえ, 同期ジェネレータでも構いません。ただし, asyncジェネレータはイベントループをブロックしないため推奨されます。
Q アップロードできるファイルの種類を制限するにはどうすればよいですか?
A file.content_type (例:text/csv)とファイル拡張子をチェックします。content_typeは偽装可能なため, フロントエンドのヒントとしてのみ使用し, サーバー側の解析と検証は必須です。

📖 まとめ


📝 練習問題

  1. 基本問題 (難易度 ⭐):UploadFileを受け取り, 内容を読み込んで行数と列名を返すCSVアップロードエンドポイントを実装します。ヒント:file: UploadFile = File(...) + csv.DictReader(io.StringIO(content))
  2. 応用問題 (難易度 ⭐⭐):CSVアップロード → Celery非同期解析ワークフローを実装します。アップロードされたファイルを一時ディレクトリに保存し, Celeryタスクをトリガーし, task_idを返し, 別のエンドポイントでタスクステータスとインポート結果を照会します。ヒント:shutil.copyfileobj + import_prices.delay(file_path)
  3. チャレンジ (難易度 ⭐⭐⭐):StreamingResponseエクスポートエンドポイントを実装します。データベースから価格データをストリーミングし, CSVデータを行単位で生成・返却し, メモリ使用量を1MB以下に抑え, プロダクトカテゴリによるフィルタリングをサポートします。ヒント:async def generate() + yield + StreamingResponse(generate(), media_type="text/csv")

---|

Web-Tutorial.com

Web-Tutorial 技術チーム

複数の開発者によって共同維持されているプログラミングチュートリアルプラットフォーム。各チュートリアルは専門分野の開発者が執筆・レビューしています。正確で信頼性の高いコンテンツを目指しています — 問題を見つけた場合はお知らせください。

100%