一覧へ

FASTQクラウドバッチ処理:ノートPCからAWSへの移行

ローカル環境で200GBのFASTQファイルを処理できない場合、Docker、S3、AWS Batchを組み合わせて、Pythonでスケーラブルなシーケンシングパイプラインを構築します。

上級
|
120
|
検証済み (2026-07)
FASTQファイルの処理AWS BatchS3クラウドパイプラインシーケンスデータBoto3Docker
進捗0/19 (0%)

FASTQ クラウドバッチ — ノートブックから AWS への移行

このトピックを終えると

教科書で学んだ AWS デプロイメントS3 ストレージ を組み合わせて、ローカルの Python スクリプトをコンテナでラップし、AWS Batch で大量実行するパイプラインを自分で作成できるようになります。例として FASTQ ファイルを使用していますが、このパターンはあらゆる種類の大量データ処理に適用できます。

この記事は 教育目的の一般的な例 です。実際のプロダクションシーケンシングパイプラインでは、Nextflow、Snakemake、WDL などのワークフロー言語を使用します。


「ノートパソコンのディスクがいっぱい」——ローカル処理の限界

シーケンシングコアから200GBのFASTQファイルを受け取ったとします。このファイルに対して、次の処理を行いたいとします。

  1. クオリティトリミング(品質の低いリードの除去)
  2. 参照ゲノムへのアライメント
  3. バリアントの呼び出し

ローカルアプローチ:

python
import subprocess
subprocess.run(["fastp", "-i", "sample.fastq.gz", "-o", "trimmed.fastq.gz"])
subprocess.run(["bwa", "mem", "reference.fa", "trimmed.fastq.gz", "-o", "aligned.sam"])
subprocess.run(["samtools", "sort", "aligned.sam", "-o", "sorted.bam"])
subprocess.run(["bcftools", "call", "sorted.bam", "-o", "variants.vcf"])

このアプローチには、3つの現実的な問題があります。

問題1: ディスク。元の200GB + 複数の中間ファイル = 1TB以上。ノートパソコンのディスク容量では足りません。

問題2: メモリと時間。BWAアライメントには、32GBのメモリと数時間かかります。ノートパソコンで他の作業を同時に行うことはできません。

問題3: 拡張性。今回はサンプルが1つですが、次のプロジェクトでは100個になります。ローカルで逐次的に処理すると、数週間かかります。

現実的なアプローチは、このパイプラインをコンテナにまとめて、AWS Batchで実行することです。元のファイルはS3に置き、Batchジョブが大きなEC2インスタンスでコンテナを起動し、ファイルを処理した後、結果を再びS3に保存します。100個のサンプルは、100個の並列ジョブとして実行されます。

ブラックボックスからコンポーネントへ:クラウドパイプラインの仕組み

主要なコンポーネントは4つあります。

コンポーネント1:S3 - 無限大のストレージ

S3 は、事実上無限大の容量を持つオブジェクトストレージです。1つのファイルの最大サイズは5TBで、バケット内のファイル数の制限はありません。ライフサイクルポリシーを使用して、古いファイルを自動的に低コストのティア(Glacier)に移行できます。

Python SDK(boto3)を使用してアクセスします。

python
import boto3
s3 = boto3.client("s3")
s3.upload_file("sample.fastq.gz", "my-bucket", "raw/sample.fastq.gz")
s3.download_file("my-bucket", "results/variants.vcf", "variants.vcf")
for obj in s3.list_objects_v2(Bucket="my-bucket", Prefix="raw/")["Contents"]:
print(obj["Key"], obj["Size"])

注意: S3の料金は、ストレージ、転送、およびリクエストによって決まります。同じリージョンのEC2からアクセスする場合、転送費用は発生しません。ファイルを処理するリージョンで実行することが、コスト削減の第一歩です。

コンポーネント2:Dockerコンテナ

パイプラインスクリプトに必要なものすべて(Python + fastp + bwa + samtools + bcftools + 自分のコード)を1つのイメージにまとめます。

Dockerfile の例:

dockerfile
FROM ubuntu:22.04

RUN apt-get update && apt-get install -y \
    python3.12 python3-pip \
    fastp bwa samtools bcftools \
    && rm -rf /var/lib/apt/lists/*

COPY requirements.txt /app/
RUN pip3 install -r /app/requirements.txt

COPY pipeline.py /app/
WORKDIR /app

ENTRYPOINT ["python3", "pipeline.py"]

ビルドしてECR(Elastic Container Registry)にプッシュします。

bash
docker build -t my-fastq-pipeline .
aws ecr get-login-password | docker login --username AWS --password-stdin <account>.dkr.ecr.<region>.amazonaws.com
docker tag my-fastq-pipeline:latest <account>.dkr.ecr.<region>.amazonaws.com/my-fastq-pipeline:latest
docker push <account>.dkr.ecr.<region>.amazonaws.com/my-fastq-pipeline:latest

コンポーネント3:AWS Batch ジョブ定義

Batch は、3つの階層の概念で構成されます。

コンピューティング環境: 実際のEC2インスタンスを管理する階層。最小vCPU、最大vCPU、インスタンスタイプを指定します。

ジョブキュー: ジョブが待機するキュー。優先度を指定します。

ジョブ定義: どのコンテナイメージを、どのリソースで、どのように実行するかを定義します。

python
import boto3
batch = boto3.client("batch")
batch.register_job_definition(
jobDefinitionName="fastq-pipeline",
type="container",
containerProperties={
"image": "<account>.dkr.ecr.<region>.amazonaws.com/my-fastq-pipeline:latest",
"vcpus": 8,
"memory": 32000,
"command": [
"--input", "Ref::input_s3",
"--output", "Ref::output_s3"
],
"jobRoleArn": "arn:aws:iam::<account>:role/BatchJobRole"
}
)

Ref::input_s3 は、ジョブを送信する際に設定されるパラメーターです。

コンポーネント4:ジョブ送信

1つのサンプルに対するジョブを送信します。

python
def submit_sample_job(sample_id: str, input_s3: str, output_s3: str) -> str:
response = batch.submit_job(
jobName=f"fastq-{sample_id}",
jobQueue="my-fastq-queue",
jobDefinition="fastq-pipeline",
parameters={
"input_s3": input_s3,
"output_s3": output_s3
}
)
return response["jobId"]
job_id = submit_sample_job(
sample_id="sample01",
input_s3="s3://my-bucket/raw/sample01.fastq.gz",
output_s3="s3://my-bucket/results/sample01/"
)
print(f"Submitted job: {job_id}")

4つの要素を組み合わせて — コンテナ内のロジック

コンテナ内で実行される pipeline.py は次のようになっています。

python
import argparse
import subprocess
from pathlib import Path
import boto3
def parse_s3_uri(uri: str) -> tuple[str, str]:
parts = uri.replace("s3://", "").split("/", 1)
return parts[0], parts[1] if len(parts) > 1 else ""
def download_from_s3(s3_uri: str, local_path: Path) -> None:
bucket, key = parse_s3_uri(s3_uri)
s3 = boto3.client("s3")
s3.download_file(bucket, key, str(local_path))
def upload_to_s3(local_path: Path, s3_uri: str) -> None:
bucket, key = parse_s3_uri(s3_uri)
s3 = boto3.client("s3")
s3.upload_file(str(local_path), bucket, key)
def run_pipeline(input_s3: str, output_s3: str, work_dir: Path) -> None:
work_dir.mkdir(exist_ok=True)
raw = work_dir / "sample.fastq.gz"
trimmed = work_dir / "trimmed.fastq.gz"
aligned = work_dir / "aligned.bam"
sorted_bam = work_dir / "sorted.bam"
variants = work_dir / "variants.vcf"
print(f"[1/5] Downloading {input_s3}")
download_from_s3(input_s3, raw)
print("[2/5] Quality trimming with fastp")
subprocess.run(["fastp", "-i", str(raw), "-o", str(trimmed)], check=True)
print("[3/5] Alignment with BWA")
with open(aligned, "w") as f:
subprocess.run(
["bwa", "mem", "/reference/hg38.fa", str(trimmed)],
stdout=f, check=True
)
print("[4/5] Sort BAM")
subprocess.run(
["samtools", "sort", str(aligned), "-o", str(sorted_bam)],
check=True
)
print("[5/5] Variant calling")
subprocess.run(
["bcftools", "call", "-o", str(variants), str(sorted_bam)],
check=True
)
print(f"Uploading {variants} to {output_s3}variants.vcf")
upload_to_s3(variants, f"{output_s3}variants.vcf")
if __name__ == "__main__":
parser = argparse.ArgumentParser()
parser.add_argument("--input", required=True)
parser.add_argument("--output", required=True)
args = parser.parse_args()
run_pipeline(args.input, args.output, Path("/tmp/work"))

このスクリプトは、あなたのノートブック(小さなファイルの場合)と、Batch コンテナ(大きなファイルの場合)の両方で実行されます。違いは、Python コードではなく実行環境です。


フェーディング — 埋めるべき3つの空白

空白1:大量ジョブ送信オーケストレーター

単一のサンプルではなく、100個のサンプルを自動送信します。

python
def submit_batch(sample_manifest: str, output_bucket: str) -> list[str]:
"""
sample_manifest: カラム [sample_id, input_s3] を持つCSV
各行に対してジョブを送信し、ジョブIDのリストを返します。
"""
import csv
job_ids = []
with open(sample_manifest) as f:
reader = csv.DictReader(f)
for row in reader:
# TODO: submit_sample_job を呼び出す
# output_s3 を output_bucket + sample_id で連結する
# 返された job_id を job_ids に追加する
pass
return job_ids

ヒント: output_s3 = f"s3://{output_bucket}/results/{row['sample_id']}/"

空白2:進捗状況ポーリング

送信したジョブの状態を定期的に確認します。

python
def wait_for_jobs(job_ids: list[str], poll_interval: int = 60) -> dict:
"""
各ジョブの最終状態を返します。
状態: SUBMITTED, PENDING, RUNNABLE, STARTING, RUNNING, SUCCEEDED, FAILED
"""
import time
final_states = {}
pending = set(job_ids)
while pending:
# TODO: batch.describe_jobs(jobs=list(pending)) を呼び出す
# 各ジョブの状態を確認し、完了したものを final_states に記録し、pending から削除する
# 残りの pending があれば、poll_interval 秒待機する
pass
return final_states

ヒント: response["jobs"] の各項目の status フィールドを確認する。SUCCEEDED, FAILED のみを最終状態として扱う。

空白3:失敗時の再試行ロジック

スポットインスタンスは低コストだが、途中で終了することがあります。失敗したジョブを自動的に再試行します。

python
def submit_with_retry(sample_id: str, input_s3: str, output_s3: str, max_attempts: int = 3) -> str:
response = batch.submit_job(
jobName=f"fastq-{sample_id}",
jobQueue="my-fastq-queue",
jobDefinition="fastq-pipeline",
parameters={"input_s3": input_s3, "output_s3": output_s3},
# TODO: retryStrategy を追加する
# attempts=max_attempts, evaluateOnExit 条件を指定する
)
return response["jobId"]

ヒント:

python
retryStrategy={
"attempts": max_attempts,
"evaluateOnExit": [
{"onStatusReason": "Host EC2*", "action": "RETRY"},
{"onExitCode": "0", "action": "EXIT"}
]
}

考察:このパイプラインは、実運用におけるシーケンスパイプラインとどのように異なるか

ワークフロー言語: 実運用では、Nextflow、Snakemake、WDLなどのワークフロー言語を使用します。皆さんのPythonスクリプトは各ステップを順番に実行しますが、ワークフロー言語はDAG(有向非巡回グラフ)としてステップ間の依存関係を表現し、並列実行とキャッシングを自動化します。

参照データ管理: 皆さんのコンテナには、hg38参照がハードコードされていますが、実運用ではEFS(Elastic File System)や別の大きなボリュームに参照を格納し、複数のジョブで共有します。

コスト最適化: 実運用では、スポットインスタンスを積極的に活用してコストを70%削減します。ただし、スポットインスタンスの終了に対応するロジックが必須です。また、各ステップのCPU/メモリ要件を分離し、小さなステップは小さなインスタンスで実行します。

セキュリティ: 実運用では、IAMロールの最小権限の原則、S3バケットの暗号化、VPCエンドポイントによるプライベートネットワーク、監査ログ(CloudTrail)などが標準です。

可視性: 各ジョブのstdout/stderr、リソース使用率、実行時間をCloudWatchで収集します。失敗時には、即座に通知(SNS)を行います。

拡張プロジェクト

1. Nextflow への移植: 上記の Python パイプラインを Nextflow DSL2 に書き換える。各ステップを process として分離し、並列化を確認する。

2. 結果表示ダッシュボード: 完了したジョブの結果 VCF を自動的に要約し、ウェブダッシュボードに表示する。Streamlit を使用すると高速化できる。

3. コスト追跡: 各ジョブで使用されたリソースとコストを計算し、サンプルごとのコストレポートを作成する。Cost Explorer API を活用する。

4. ローカル開発モード: LocalStack を使用して、ローカルで S3/Batch をシミュレートする開発環境を構築する。実際の AWS コストをかけずにパイプラインを開発する。

このパートの構成要素

  • [F] AWSデプロイ: バッチコンピューティング環境、ジョブキュー、ジョブ定義の設計。
  • [F] S3: 大量のデータの保存とアクセス。リージョン、コスト、転送の理解。
  • [W] venv・file-io: Pythonの仮想環境とローカルファイルの処理(コンテナ内で使用)。
  • [W] Docker: コンテナイメージのビルドとECRへのプッシュ(完成したスクリプトを提供)。

[F] = 自分で実装する / [W] = 完成したコードで提供するツールの概念。

💬 質問・コメント

0件のコメント

ログインせずに投稿できます。ゲスト投稿は投稿者自身で編集・削除できません。

0/2000

読み込み中...