FASTQ クラウドバッチ — ノートブックから AWS への移行
このトピックを終えると
教科書で学んだ AWS デプロイメント と S3 ストレージ を組み合わせて、ローカルの Python スクリプトをコンテナでラップし、AWS Batch で大量実行するパイプラインを自分で作成できるようになります。例として FASTQ ファイルを使用していますが、このパターンはあらゆる種類の大量データ処理に適用できます。
この記事は 教育目的の一般的な例 です。実際のプロダクションシーケンシングパイプラインでは、Nextflow、Snakemake、WDL などのワークフロー言語を使用します。
「ノートパソコンのディスクがいっぱい」——ローカル処理の限界
シーケンシングコアから200GBのFASTQファイルを受け取ったとします。このファイルに対して、次の処理を行いたいとします。
- クオリティトリミング(品質の低いリードの除去)
- 参照ゲノムへのアライメント
- バリアントの呼び出し
ローカルアプローチ:
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)を使用してアクセスします。
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 の例:
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)にプッシュします。
docker build -t my-fastq-pipeline .aws ecr get-login-password | docker login --username AWS --password-stdin <account>.dkr.ecr.<region>.amazonaws.comdocker tag my-fastq-pipeline:latest <account>.dkr.ecr.<region>.amazonaws.com/my-fastq-pipeline:latestdocker push <account>.dkr.ecr.<region>.amazonaws.com/my-fastq-pipeline:latestコンポーネント3:AWS Batch ジョブ定義
Batch は、3つの階層の概念で構成されます。
コンピューティング環境: 実際のEC2インスタンスを管理する階層。最小vCPU、最大vCPU、インスタンスタイプを指定します。
ジョブキュー: ジョブが待機するキュー。優先度を指定します。
ジョブ定義: どのコンテナイメージを、どのリソースで、どのように実行するかを定義します。
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つのサンプルに対するジョブを送信します。
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 は次のようになっています。
import argparseimport subprocessfrom pathlib import Pathimport 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個のサンプルを自動送信します。
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:進捗状況ポーリング
送信したジョブの状態を定期的に確認します。
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:失敗時の再試行ロジック
スポットインスタンスは低コストだが、途中で終了することがあります。失敗したジョブを自動的に再試行します。
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"]ヒント:
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] = 完成したコードで提供するツールの概念。