ファイルアップロードとダウンロード — 大容量ファイルと一括インポート
ファイルのアップロードは荷物の受け取りのようなものです。小包は即座にサインして受け取ります (メモリ), 大量の荷物は分割して荷下ろしします (ストリーム書き込み)。ダウンロードは荷物の発送のようなものです。単品はFileResponse, 大量発送はコンベアベルトに乗せます (StreamingResponse)。
1. 学ぶ内容
UploadFileとFile()パラメータ:ストリームベースアップロード vs. インメモリアップロード- 大容量ファイル処理:ブロック単位の読み込み,
shutil.copyfileobjストリーム書き込み - CSV/Excel解析:
pandas+ Pydanticモデルの合同検証 - ファイルレスポンス:
FileResponse/StreamingResponseダウンロードエンドポイント - Aliceシナリオ:Bobが百万行の価格データを含むCSVファイルをアップロード → Celeryがバックグラウンドで非同期に解析 → インポート完了時に通知
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) 大容量ファイルアップロード + 非同期処理ワークフロー
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は偽装可能なため, フロントエンドのヒントとしてのみ使用し, サーバー側の解析と検証は必須です。📖 まとめ
UploadFileはデータを一時ファイルにストリームし, メモリを使用せず, 任意のサイズのファイルに適しています- 大容量ファイルには
shutil.copyfileobj()またはread()でチャンクストリーミングしてディスクに書き込みます - CSVアップロード + Celery非同期解析:APIは即座に
task_idを返し, バックグラウンドプロセスが各行を検証してバッチ挿入します - PydanticはCSVデータを行単位で検証し, 不正行をスキップ・ログ記録し, 寛容モードで部分インポートを確保します
FileResponseは既存ファイルのダウンロードに,StreamingResponseは動的生成・リアルタイム返却に使用します
📝 練習問題
- 基本問題 (難易度 ⭐):
UploadFileを受け取り, 内容を読み込んで行数と列名を返すCSVアップロードエンドポイントを実装します。ヒント:file: UploadFile = File(...)+csv.DictReader(io.StringIO(content)) - 応用問題 (難易度 ⭐⭐):CSVアップロード → Celery非同期解析ワークフローを実装します。アップロードされたファイルを一時ディレクトリに保存し, Celeryタスクをトリガーし,
task_idを返し, 別のエンドポイントでタスクステータスとインポート結果を照会します。ヒント:shutil.copyfileobj+import_prices.delay(file_path) - チャレンジ (難易度 ⭐⭐⭐):StreamingResponseエクスポートエンドポイントを実装します。データベースから価格データをストリーミングし, CSVデータを行単位で生成・返却し, メモリ使用量を1MB以下に抑え, プロダクトカテゴリによるフィルタリングをサポートします。ヒント:
async def generate()+yield+StreamingResponse(generate(), media_type="text/csv")
---|



