Qwen3.6-35B-A3B を Amazon SageMaker で動かしてストリーミング応答を得てみた
はじめに
大規模な LLM を試してみたくても、手元の GPU ではメモリが足りないことがあります。かといって、検証のためだけに GPU サーバーを用意し、OS や CUDA を管理するのも大変です。
そこで、Qwen3.6-35B-A3B の 21 GB 級 Q4_K_M GGUF と固定した llama.cpp をカスタムコンテナへ組み込み、SageMaker Real-time Inference の ml.g5.2xlargeで実行しました。モデルはコンテナへ含めず、Amazon S3 から SageMaker に配置させています。
次の GIF は、Next.js 検証用アプリから SageMaker のエンドポイントを呼び出し、Qwen3.6 の回答を逐次表示した様子です。

SageMaker Real-time Inference とは
SageMaker Real-time Inference は、推論コンテナを常時稼働するエンドポイントとしてホストする機能です。カスタムコンテナを利用でき、Response Streaming にも対応します。
Qwen3.6 とは
Qwen3.6 は、Qwen チームが公開しているオープンウェイトモデルです。本記事で使用した Qwen3.6-35B-A3B は、総パラメータ数が 35B、推論時の有効パラメータ数が 3B の MoE モデルです。公式のモデルウェイトではなく、LM Studio Community が配布する Q4_K_M GGUF を検証しました。
検証環境
- リージョン:
us-west-2 - インスタンス:
ml.g5.2xlarge - GPU: NVIDIA A10G、24 GB
- モデル:
lmstudio-community/Qwen3.6-35B-A3B-GGUF - GGUF:
Qwen3.6-35B-A3B-Q4_K_M.gguf - GGUF サイズ: 21,166,757,728 bytes
- コンテキスト長: 16,384
- 並列スロット: 1
- llama.cpp:
86a9c79f866799eb0e7e89c03578ccfbcc5d808e - 検証時のイメージダイジェスト:
sha256:bbc7ad869bc41d90644e170e3d077aa5ff93a1cf93333cfe9b3ad9789519bedc - CUDA: 12.4.1
- Terraform: 1.15.8
- AWS Provider: 6.49.0
ml.g5.2xlargeは、24 GB の GPU メモリ、8 vCPU、32 GiB のホストメモリを備えています。インスタンス仕様は Amazon EC2 G5 Instances で確認できます。
実行前に、Service Quotas の ml.g5.2xlarge for endpoint usageを1以上にします。検証時のクォータコードは L-9614C779でした。
対象読者
- 手元の GPU に収まらない LLM を AWS 上で試したい方
- GGUF と llama.cpp を SageMaker のカスタムコンテナで動かしたい方
- SageMaker Response Streaming を使って生成結果を逐次受信したい方
参考資料
- Custom Inference Code with Hosting Services
- Deploying uncompressed models
- InvokeEndpointWithResponseStream
- Invoke models for real-time inference
- SageMaker Pricing
- Hugging Face Hub CLI
- Qwen3.6-35B-A3B
- Qwen3.6-35B-A3B-GGUF
検証構成と技術選定
要件は次の3つです。
- 任意の GGUF と固定した llama.cpp を使えること
- 生成結果を逐次受信できること
- 検証後に GPU を停止できること
構成は次のようにしました。
Amazon Bedrock へ GGUF と llama.cpp をそのまま持ち込むことはできません。 GPU 付き Amazon EC2 では OS と推論プロセスの管理が必要です。Processing Job と Asynchronous Inference は非同期処理となるため、Response Streaming に対応する Real-time Inference を選びました。
モデルとコンテナの準備
GGUF の固定
第三者配布の GGUF では、変換元の正確なリビジョンと量子化コマンドが公開されていません。そのため、配布リポジトリのリビジョン、ファイルサイズ、実ファイルの SHA-256 を固定しました。
| 項目 | 値 |
|---|---|
| 配布元 | lmstudio-community/Qwen3.6-35B-A3B-GGUF |
| リビジョン | c7ed48a94a6a082167b768bab350745695824e0a |
| ファイル | Qwen3.6-35B-A3B-Q4_K_M.gguf |
| 量子化 | Q4_K_M |
| サイズ | 21,166,757,728 bytes |
| SHA-256 | 4ac6a06bce551257267f49ad2226f8671a22519ccc1a4dde9d5b433d1f2a410d |
| ライセンスメタデータ | Apache-2.0 |
stat --format='%s' ./Qwen3.6-35B-A3B-Q4_K_M.gguf
sha256sum ./Qwen3.6-35B-A3B-Q4_K_M.gguf
SHA-256 の一致を確認した後、バージョニングと Block Public Access を有効にした非公開 S3 バケットへ配置しました。SageMaker の ModelDataSourceでは、単一の非圧縮オブジェクトを指定しています。AWS の仕様では、S3DataType = "S3Object"かつCompressionType = "None"とすると、ファイルは /opt/ml/modelへ配置されます。
llama.cpp コンテナ
コンテナには、コミットハッシュを固定して CUDA 対応にした llama.cpp と、SageMaker 用アダプターを含めました。モデル変更のたびに巨大なイメージを再構築しないよう、GGUF は含めません。
llama.cpp のコミットハッシュを固定し、CUDA を有効にして llama-serverをビルドします。
ARG LLAMA_CPP_COMMIT=86a9c79f866799eb0e7e89c03578ccfbcc5d808e
WORKDIR /src
RUN git init llama.cpp \
&& cd llama.cpp \
&& git remote add origin https://github.com/ggml-org/llama.cpp.git \
&& git fetch --depth 1 origin "${LLAMA_CPP_COMMIT}" \
&& git checkout --detach FETCH_HEAD \
&& test "$(git rev-parse HEAD)" = "${LLAMA_CPP_COMMIT}"
RUN cmake -S /src/llama.cpp -B /src/llama.cpp/build -G Ninja \
-DCMAKE_BUILD_TYPE=Release \
-DCMAKE_CUDA_ARCHITECTURES="86;89" \
-DGGML_CUDA=ON \
-DLLAMA_BUILD_SERVER=ON \
&& cmake --build /src/llama.cpp/build --target llama-server --parallel
アダプターの責務は次のとおりです。
- SageMaker が要求するポート
8080で/pingと/invocationsを提供する - GGUF のサイズと SHA-256 を確認してから llama.cpp を起動する
- llama.cpp の OpenAI 互換 SSE を HTTP chunked response として転送する
- プロンプトと生成本文を通常ログへ出さない
/pingは、llama.cpp の /healthが成功した後に 200を返します。/invocationsでは、llama.cpp の SSE を chunked response として転送します。入力検証とエラー処理を含む全文は付録に掲載しています。
def do_GET(self) -> None:
request_id = self.headers.get("X-Request-Id") or str(uuid.uuid4())
if self.path != "/ping":
self.send_json(404, {"error": "not_found"}, request_id)
return
ready = upstream_ready()
self.send_json(
200 if ready else 503,
{"status": "ok" if ready else "unavailable"},
request_id,
)
def do_POST(self) -> None:
# 入力検証と upstream_request の組み立ては省略
...
with request.urlopen(upstream_request, timeout=UPSTREAM_TIMEOUT_SECONDS) as response:
if response.headers.get_content_type() != "text/event-stream":
self.send_json(502, {"error": "model_server_invalid_content_type"}, request_id)
return
self.send_response(response.status)
self.send_header("Content-Type", "text/event-stream")
self.send_header("Transfer-Encoding", "chunked")
self.end_headers()
while chunk := response.read1(64 * 1024):
self.wfile.write(f"{len(chunk):X}\r\n".encode("ascii"))
self.wfile.write(chunk)
self.wfile.write(b"\r\n")
self.wfile.flush()
self.wfile.write(b"0\r\n\r\n")
Terraform によるエンドポイントの作成
Terraform では、ロググループ、実行ロール、SageMaker Model、EndpointConfig、エンドポイントを作成しました。S3 VersionId と ECR イメージダイジェストを入力値へ固定し、terraform plan時に現在のオブジェクトと照合しています。
EndpointConfig は ml.g5.2xlargeを1台起動します。モデルダウンロードとヘルスチェックのタイムアウトは、それぞれ 1,800 秒です。
モデルは非圧縮の S3 オブジェクトとして配置し、コンテナへファイル名、サイズ、SHA-256 を渡します。EndpointConfig では、検証に使用したインスタンスとタイムアウトを固定しました。
resource "aws_sagemaker_model" "llama" {
name = "${var.name_prefix}-model"
execution_role_arn = aws_iam_role.sagemaker.arn
enable_network_isolation = true
primary_container {
image = var.container_image_uri
mode = "SingleModel"
model_data_source {
s3_data_source {
compression_type = "None"
s3_data_type = "S3Object"
s3_uri = "s3://${var.artifact_bucket_name}/${var.model_object_key}"
}
}
environment = {
MODEL_FILE = "Qwen3.6-35B-A3B-Q4_K_M.gguf"
MODEL_ARTIFACT_SIZE_BYTES = tostring(var.model_size_bytes)
MODEL_ARTIFACT_SHA256 = var.model_sha256
CONTEXT_SIZE = "16384"
PARALLEL = "1"
GPU_LAYERS = "999"
}
}
}
resource "aws_sagemaker_endpoint_configuration" "realtime" {
name = "${var.name_prefix}-config"
production_variants {
variant_name = "AllTraffic"
model_name = aws_sagemaker_model.llama.name
initial_instance_count = 1
instance_type = "ml.g5.2xlarge"
model_data_download_timeout_in_seconds = 1800
container_startup_health_check_timeout_in_seconds = 1800
}
}
resource "aws_sagemaker_endpoint" "realtime" {
name = var.name_prefix
endpoint_config_name = aws_sagemaker_endpoint_configuration.realtime.name
}
initial_instance_count = 1のため、エンドポイントの作成と同時に GPU の準備が始まります。Service Quotas API の適用値が1以上でも、キャパシティ不足によって InsufficientInstanceCapacityとなる場合があります。
ストリーミング推論
コンテナは llama.cpp の SSE を HTTP chunked response で返します。呼び出し元では、AWS SDK for JavaScript の InvokeEndpointWithResponseStreamCommandを使いました。
AWS SDK が返す PayloadPartは、コンテナ側の chunk と同じ境界になるとは限りません。各バイト列を streaming TextDecoderへ順に渡し、UTF-8 文字や SSE イベントの途中で分割されても復元します。
async function* readSse(
chunks: AsyncIterable<Uint8Array>,
): AsyncGenerator<string> {
const decoder = new TextDecoder("utf-8");
let buffer = "";
for await (const chunk of chunks) {
buffer += decoder.decode(chunk, { stream: true });
while (true) {
const match = buffer.match(/\r?\n\r?\n/);
if (!match || match.index === undefined) break;
const block = buffer.slice(0, match.index);
buffer = buffer.slice(match.index + match[0].length);
const data = block
.split(/\r?\n/)
.map((line) => line.match(/^data:\s?(.*)$/)?.[1])
.filter((line): line is string => line !== undefined)
.join("\n");
if (data) yield data;
}
}
buffer += decoder.decode();
if (buffer.trim()) {
const data = buffer
.split(/\r?\n/)
.map((line) => line.match(/^data:\s?(.*)$/)?.[1])
.filter((line): line is string => line !== undefined)
.join("\n");
if (data) yield data;
}
}
InvokeEndpointWithResponseStreamCommandの応答から PayloadPart.Bytesだけを取り出し、上記の関数へ渡します。読み取りを中止した場合は、AbortControllerで AWS SDK の呼び出しも中止します。
const response = await client.send(
new InvokeEndpointWithResponseStreamCommand({
EndpointName: endpointName,
ContentType: "application/json",
Accept: "text/event-stream",
Body: new TextEncoder().encode(JSON.stringify(requestBody)),
}),
{ abortSignal: abortController.signal },
);
if (!response.Body) throw new Error("response body is missing");
const payloads = (async function* () {
for await (const event of response.Body!) {
if (!isPayloadPart(event) || !event.PayloadPart.Bytes) {
throw new Error("SageMaker response stream failed");
}
yield event.PayloadPart.Bytes;
}
})();
for await (const data of readSse(payloads)) {
if (data === "[DONE]") break;
const event = JSON.parse(data);
const delta = event?.choices?.[0]?.delta?.content;
if (typeof delta === "string") process.stdout.write(delta);
}
リクエスト本文の組み立てを含む完全なコードは付録に掲載しています。
実行結果
入力と出力
公開可能な合成データを使い、次のメッセージを送信しました。
system:
あなたは簡潔で安全な日本語アシスタントです。
user:
雨の日に家で楽しめる創作活動を三つ、各一文で提案してください。
生成結果の原文は次のとおりです。表記を含めて変更していません。
1. 好きな曲に合わせて、その日の気分や雨の音をイメージした短歌や詩を書いてみましょう。
2. 家に転がっている古紙や布を使って、自分だけのオリジナルカードや小物を作ってみましょう。
3. 窓から見える雨景色をスケッチしたり、写真に収めたりして、アート作品に仕上げてみましょう。
Next.js アプリからのストリーミング確認
同じ構成を Next.js 検証用アプリから呼び出し、SageMaker のイベントを画面表示用の差分へ変換しました。
次の GIF は、エンドポイントが利用可能になった後にプロンプトを送信し、回答が逐次表示されて完了するまでの画面です。次節の計測値は、AWS SDK を直接呼び出して取得しました。

計測値
この入力を使い、同じエンドポイントに対して正常完了とクライアント中断を確認しました。生成設定は max_tokens = 192、temperature = 0.2、top_p = 0.9、seed = 42、thinking 無効です。
| 項目 | 実測値 |
|---|---|
エンドポイント作成から InService |
302.771 秒 |
| アダプター起動 | エンドポイント作成から約 245.7 秒 |
| コンテナ内の SHA-256 確認 | 20.019 秒 |
| llama.cpp のモデルロード | 約 3.95 秒 |
| ウォーム TTFT | 2,058 ms |
| 正常完了まで | 2,601 ms |
PayloadPart |
92 |
| 非空の表示用デルタ | 76 |
| 生成トークン | 77 |
| llama.cpp が記録した生成速度 | 123.71 tokens/second |
| クライアント中断 | 最初のデルタを受信後、578 msで終了 |
terraform destroy |
17.280 秒 |
エンドポイント作成から InServiceまでの 302.771 秒には、インスタンス準備、S3 からのモデル配置、コンテナ起動、SHA-256 確認、モデルロードが含まれます。llama.cpp のモデルロード自体は約 3.95 秒でした。
us-west-2における ml.g5.2xlargeの確認時点の単価は 1.515 USD/hour でした。terraform apply開始からterraform destroy完了までの約 910 秒を全量課金として計算すると、約 0.383 USD です。
計測値は、短い単一プロンプトを一回実行した結果です。
削除結果
terraform destroyは 17.280 秒で完了しました。SageMaker の一覧と Terraform state を照合し、エンドポイント、EndpointConfig、SageMaker Model が残っていないことを確認しました。S3 のモデルと ECR イメージは再利用のため残しています。
考察
エンドポイント作成から InServiceまでは 302.771 秒、コンテナ内のモデルロードは約 3.95 秒、ウォーム TTFT は 2,058 ms でした。停止状態から利用する場合の待ち時間は、モデルの推論よりもインスタンス準備やモデル配置の影響を強く受けると考えられます。Web アプリへ組み込む場合は、エンドポイントの起動とプロンプト送信を別の操作として扱うと、利用者が待ち時間を把握しやすくなるでしょう。
モデルを S3、llama.cpp を ECR へ分離したため、コンテナイメージを作り直さずに GGUF を差し替えられます。同じ llama.cpp のコミットと推論条件を維持できることから、複数の量子化モデルを同じ実行環境で比較する用途にも利用できると考えられます。
Next.js 検証用アプリでも回答を逐次表示できました。SageMaker Real-time Inference はモデルのロード確認だけでなく、独自のモデルを使う Web アプリの試作環境としても利用できそうです。
まとめ
Qwen3.6 の 21 GB 級 Q4_K_M GGUF を ml.g5.2xlargeへロードし、SageMaker Response Streaming で回答を逐次受信できました。ウォーム TTFT は 2,058 ms、エンドポイント作成から InServiceまでは 302.771 秒、作成から削除までの概算料金は約 0.383 USD でした。
GPU ホストの OS を管理せずに任意の GGUF を試せる一方、カスタムコンテナの保守とエンドポイントの作成・削除は必要です。この運用を許容できる短時間の検証では、SageMaker Real-time Inference を利用できます。
付録
ここからは、再現に必要な固定値とコードを掲載します。AWS アカウント ID、実際の ARN、バケット名、ECR リポジトリ名、S3 VersionId はプレースホルダーへ置き換えています。
モデルの取得と配置
Hugging Face CLI の取得対象を revision とファイル名で固定します。ダウンロード後にサイズと SHA-256 が一致しなければ、S3 へアップロードしないでください。
python -m pip install --upgrade huggingface_hub
hf download \
lmstudio-community/Qwen3.6-35B-A3B-GGUF \
Qwen3.6-35B-A3B-Q4_K_M.gguf \
--revision c7ed48a94a6a082167b768bab350745695824e0a \
--local-dir ./model
sha256sum ./model/Qwen3.6-35B-A3B-Q4_K_M.gguf
aws s3 cp \
./model/Qwen3.6-35B-A3B-Q4_K_M.gguf \
s3://<PRIVATE_MODEL_BUCKET>/models/Qwen3.6-35B-A3B-Q4_K_M.gguf \
--region us-west-2
aws s3api head-object \
--bucket <PRIVATE_MODEL_BUCKET> \
--key models/Qwen3.6-35B-A3B-Q4_K_M.gguf \
--region us-west-2 \
--query '{VersionId:VersionId,ContentLength:ContentLength}'
コンテナ
Dockerfile の全文
# syntax=docker/dockerfile:1.7
ARG CUDA_DEVEL_IMAGE=nvidia/cuda:12.4.1-devel-ubuntu22.04@sha256:5645fec64549cc35930eee9d85aafd2b0006c0c3f22632be5a1d85e2604e9749
ARG CUDA_RUNTIME_IMAGE=nvidia/cuda:12.4.1-runtime-ubuntu22.04@sha256:cff3a0d82d2c2b47bab252d67fa9b34a20ef4c50781d98501b5c7367ea9afd10
FROM ${CUDA_DEVEL_IMAGE} AS builder
ARG LLAMA_CPP_COMMIT=86a9c79f866799eb0e7e89c03578ccfbcc5d808e
ARG CUDA_ARCHITECTURES=86;89
RUN apt-get update \
&& DEBIAN_FRONTEND=noninteractive apt-get install -y --no-install-recommends \
build-essential ca-certificates cmake git ninja-build \
&& rm -rf /var/lib/apt/lists/*
WORKDIR /src
RUN git init llama.cpp \
&& cd llama.cpp \
&& git remote add origin https://github.com/ggml-org/llama.cpp.git \
&& git fetch --depth 1 origin "${LLAMA_CPP_COMMIT}" \
&& git checkout --detach FETCH_HEAD \
&& test "$(git rev-parse HEAD)" = "${LLAMA_CPP_COMMIT}"
RUN cmake -S /src/llama.cpp -B /src/llama.cpp/build -G Ninja \
-DCMAKE_BUILD_TYPE=Release \
-DCMAKE_CUDA_ARCHITECTURES="${CUDA_ARCHITECTURES}" \
-DGGML_CUDA=ON \
-DGGML_NATIVE=OFF \
-DBUILD_SHARED_LIBS=OFF \
-DLLAMA_CURL=OFF \
-DLLAMA_BUILD_TESTS=OFF \
-DLLAMA_BUILD_EXAMPLES=ON \
-DLLAMA_BUILD_SERVER=ON \
&& cmake --build /src/llama.cpp/build --target llama-server --parallel
RUN printf '%s\n' "${LLAMA_CPP_COMMIT}" > /src/llama.cpp/build/LLAMA_CPP_COMMIT
FROM builder AS metadata-inspector
RUN apt-get update \
&& DEBIAN_FRONTEND=noninteractive apt-get install -y --no-install-recommends \
python3 python3-numpy python3-yaml \
&& rm -rf /var/lib/apt/lists/*
FROM ${CUDA_RUNTIME_IMAGE} AS runtime
ARG LLAMA_CPP_COMMIT=86a9c79f866799eb0e7e89c03578ccfbcc5d808e
ARG CUDA_ARCHITECTURES=86;89
RUN apt-get update \
&& DEBIAN_FRONTEND=noninteractive apt-get install -y --no-install-recommends \
ca-certificates libgomp1 python3 util-linux \
&& rm -rf /var/lib/apt/lists/* \
&& groupadd --system --gid 10001 model-server \
&& useradd --system --uid 10001 --gid 10001 --create-home --home-dir /home/model-server model-server
COPY --from=builder /src/llama.cpp/build/bin/llama-server /opt/llama/bin/llama-server
COPY --from=builder /src/llama.cpp/build/LLAMA_CPP_COMMIT /opt/llama/LLAMA_CPP_COMMIT
COPY adapter.py /opt/program/adapter.py
COPY --chmod=0755 processing-entrypoint.sh /opt/program/processing-entrypoint.sh
ENV PYTHONUNBUFFERED=1 \
PORT=8080 \
LLAMA_SERVER_URL=http://127.0.0.1:8081 \
START_LLAMA_SERVER=1 \
MODEL_DIR=/opt/ml/model \
CONTEXT_SIZE=16384 \
PARALLEL=1 \
GPU_LAYERS=999 \
LLAMA_CPP_COMMIT=${LLAMA_CPP_COMMIT} \
CUDA_ARCHITECTURES=${CUDA_ARCHITECTURES}
LABEL org.opencontainers.image.source="https://github.com/ggml-org/llama.cpp" \
org.opencontainers.image.revision="${LLAMA_CPP_COMMIT}" \
aws-llm.model-in-image="false"
USER 0
EXPOSE 8080
ENTRYPOINT ["/opt/program/processing-entrypoint.sh"]
CMD ["python3", "/opt/program/adapter.py", "serve"]
processing-entrypoint.sh の全文
#!/bin/sh
set -eu
runtime_uid=10001
runtime_gid=10001
processing_output=/opt/ml/processing/output
if [ "$(id -u)" -ne 0 ]; then
echo "processing entrypoint must start as root" >&2
exit 70
fi
if [ "${PREPARE_PROCESSING_OUTPUT:-0}" = "1" ]; then
if [ "${SPIKE_OUTPUT_DIR:-$processing_output}" != "$processing_output" ]; then
echo "refusing to prepare an unexpected output path" >&2
exit 70
fi
mkdir -p "$processing_output"
chown "$runtime_uid:$runtime_gid" "$processing_output"
chmod 0750 "$processing_output"
fi
# SageMaker Hosting replaces the image command with the single argument
# `serve`. Convert only that exact invocation; Processing commands remain
# unchanged.
if [ "$#" -eq 1 ] && [ "$1" = "serve" ]; then
set -- python3 /opt/program/adapter.py serve
fi
exec setpriv \
--reuid="$runtime_uid" \
--regid="$runtime_gid" \
--init-groups \
-- "$@"
アダプターは、次の完全版を adapter.pyとして保存します。
adapter.py の全文
"""Minimal SageMaker /ping and /invocations adapter for a local llama-server.
Request and generated text are deliberately excluded from application logs.
"""
from __future__ import annotations
import hashlib
import http.client
import json
import logging
import os
import re
import signal
import subprocess
import threading
import time
import uuid
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from pathlib import Path
from urllib import error, request
logging.basicConfig(
level=os.environ.get("LOG_LEVEL", "INFO"),
format="%(asctime)s %(levelname)s %(message)s",
)
LOG = logging.getLogger("sagemaker-llama-adapter")
PORT = int(os.environ.get("PORT", "8080"))
UPSTREAM = os.environ.get("LLAMA_SERVER_URL", "http://127.0.0.1:8081").rstrip("/")
START_LLAMA_SERVER = os.environ.get("START_LLAMA_SERVER", "1") == "1"
UPSTREAM_TIMEOUT_SECONDS = int(os.environ.get("UPSTREAM_TIMEOUT_SECONDS", "1800"))
MAX_REQUEST_BYTES = int(os.environ.get("MAX_REQUEST_BYTES", "1048576"))
REQUEST_LOCK = threading.Lock()
CHILD: subprocess.Popen[bytes] | None = None
SHUTTING_DOWN = threading.Event()
def mount_point_for(path: Path, mountinfo: str | None = None) -> Path:
"""Resolve the longest Linux mountinfo entry containing path."""
resolved = path.resolve(strict=True)
if mountinfo is None:
try:
mountinfo = Path("/proc/self/mountinfo").read_text(encoding="utf-8")
except OSError:
mountinfo = ""
candidates: list[Path] = []
for line in mountinfo.splitlines():
fields = line.split()
if len(fields) < 5:
continue
decoded = re.sub(r"\\([0-7]{3})", lambda match: chr(int(match.group(1), 8)), fields[4])
candidate = Path(decoded)
try:
resolved.relative_to(candidate)
except ValueError:
continue
candidates.append(candidate)
if candidates:
return max(candidates, key=lambda candidate: len(str(candidate)))
candidate = resolved.parent
while candidate.parent != candidate and not os.path.ismount(candidate):
candidate = candidate.parent
return candidate
def storage_telemetry(model_path: Path) -> dict[str, int | str]:
"""Return non-content storage evidence for the SageMaker model mount."""
resolved = model_path.resolve(strict=True)
mount_point = mount_point_for(resolved)
stats = os.statvfs(resolved)
block_size = stats.f_frsize or stats.f_bsize
return {
"model_path": str(resolved),
"mount_point": str(mount_point),
"model_size_bytes": resolved.stat().st_size,
"filesystem_total_bytes": stats.f_blocks * block_size,
"filesystem_free_bytes": stats.f_bfree * block_size,
"filesystem_available_bytes": stats.f_bavail * block_size,
}
def log_storage_telemetry(model_path: Path) -> None:
telemetry = storage_telemetry(model_path)
LOG.info(
"event=model_storage model_path=%s mount_point=%s model_size_bytes=%s "
"filesystem_total_bytes=%s filesystem_free_bytes=%s filesystem_available_bytes=%s",
telemetry["model_path"],
telemetry["mount_point"],
telemetry["model_size_bytes"],
telemetry["filesystem_total_bytes"],
telemetry["filesystem_free_bytes"],
telemetry["filesystem_available_bytes"],
)
def configured_model_path() -> Path:
model_file = os.environ.get("MODEL_FILE", "").strip()
if not model_file:
raise ValueError("MODEL_FILE must be configured")
if Path(model_file).name != model_file:
raise ValueError("MODEL_FILE must contain only a filename")
return Path(os.environ.get("MODEL_DIR", "/opt/ml/model")) / model_file
def llama_command() -> list[str]:
model_path = configured_model_path()
if not model_path.is_file():
raise FileNotFoundError(f"model file not found at configured MODEL_DIR/MODEL_FILE")
log_storage_telemetry(model_path)
verify_model(model_path)
return [
"/opt/llama/bin/llama-server",
"--model",
str(model_path),
"--host",
"127.0.0.1",
"--port",
"8081",
"--ctx-size",
os.environ.get("CONTEXT_SIZE", "16384"),
"--parallel",
os.environ.get("PARALLEL", "1"),
"--n-gpu-layers",
os.environ.get("GPU_LAYERS", "999"),
"--jinja",
"--metrics",
"--no-webui",
]
def verify_model(model_path: Path) -> None:
expected_size = int(os.environ.get("MODEL_ARTIFACT_SIZE_BYTES", "0"))
expected_sha256 = os.environ.get("MODEL_ARTIFACT_SHA256", "").lower()
if expected_size <= 0 or len(expected_sha256) != 64:
raise ValueError("model size and SHA-256 must be configured")
actual_size = model_path.stat().st_size
if actual_size != expected_size:
raise ValueError("model size verification failed")
started = time.monotonic()
digest = hashlib.sha256()
with model_path.open("rb") as handle:
while block := handle.read(64 * 1024 * 1024):
if SHUTTING_DOWN.is_set():
raise InterruptedError("shutdown requested during model verification")
digest.update(block)
if digest.hexdigest() != expected_sha256:
raise ValueError("model SHA-256 verification failed")
LOG.info("event=model_artifact_verified bytes=%s elapsed_ms=%s", actual_size, round((time.monotonic() - started) * 1000))
def start_child() -> None:
global CHILD
if not START_LLAMA_SERVER:
LOG.info("event=llama_child_external")
return
try:
CHILD = subprocess.Popen(llama_command())
LOG.info("event=llama_child_started pid=%s", CHILD.pid)
except Exception as exc: # startup failure must remain visible without prompt data
LOG.error("event=llama_child_start_failed type=%s", type(exc).__name__)
def upstream_ready(timeout: float = 2.0) -> bool:
if START_LLAMA_SERVER and (CHILD is None or CHILD.poll() is not None):
return False
try:
with request.urlopen(f"{UPSTREAM}/health", timeout=timeout) as response:
return response.status == 200
except (error.URLError, TimeoutError, OSError):
return False
def validate_payload(value: object) -> dict:
if not isinstance(value, dict):
raise ValueError("request body must be a JSON object")
messages = value.get("messages")
if not isinstance(messages, list) or not messages:
raise ValueError("messages must be a non-empty array")
for item in messages:
if not isinstance(item, dict):
raise ValueError("each message must be an object")
if item.get("role") not in {"system", "user", "assistant"}:
raise ValueError("message role is invalid")
if not isinstance(item.get("content"), str) or not item["content"]:
raise ValueError("message content must be a non-empty string")
result = {
"messages": messages,
"stream": True,
"max_tokens": int(value.get("max_tokens", 256)),
"temperature": float(value.get("temperature", 0.7)),
"top_p": float(value.get("top_p", 0.9)),
}
if "seed" in value:
result["seed"] = int(value["seed"])
if "enable_thinking" in value:
result["chat_template_kwargs"] = {
"enable_thinking": bool(value["enable_thinking"])
}
return result
class Handler(BaseHTTPRequestHandler):
server_version = "sagemaker-llama-adapter/1"
protocol_version = "HTTP/1.1"
def log_message(self, _format: str, *_args: object) -> None:
return
def send_json(self, status: int, payload: dict, request_id: str) -> None:
body = json.dumps(payload, separators=(",", ":")).encode("utf-8")
self.send_response(status)
self.send_header("Content-Type", "application/json")
self.send_header("Content-Length", str(len(body)))
self.send_header("X-Request-Id", request_id)
self.end_headers()
self.wfile.write(body)
def do_GET(self) -> None:
request_id = self.headers.get("X-Request-Id") or str(uuid.uuid4())
if self.path != "/ping":
self.send_json(404, {"error": "not_found"}, request_id)
return
ready = upstream_ready()
self.send_json(200 if ready else 503, {"status": "ok" if ready else "unavailable"}, request_id)
def do_POST(self) -> None:
started = time.monotonic()
request_id = self.headers.get("X-Request-Id") or str(uuid.uuid4())
status = 500
response_started = False
try:
if self.path != "/invocations":
status = 404
self.send_json(status, {"error": "not_found"}, request_id)
return
if self.headers.get_content_type() != "application/json":
status = 415
self.send_json(status, {"error": "content_type_must_be_json"}, request_id)
return
try:
length = int(self.headers.get("Content-Length", "0"))
except ValueError:
length = -1
if length <= 0 or length > MAX_REQUEST_BYTES:
status = 413
self.send_json(status, {"error": "invalid_request_size"}, request_id)
return
try:
payload = validate_payload(json.loads(self.rfile.read(length)))
except (json.JSONDecodeError, UnicodeDecodeError, ValueError, TypeError, OverflowError):
status = 400
self.send_json(status, {"error": "invalid_request"}, request_id)
return
if not upstream_ready():
status = 503
self.send_json(status, {"error": "model_server_unavailable"}, request_id)
return
if not REQUEST_LOCK.acquire(blocking=False):
status = 429
self.send_json(status, {"error": "request_in_progress"}, request_id)
return
try:
upstream_request = request.Request(
f"{UPSTREAM}/v1/chat/completions",
data=json.dumps(payload).encode("utf-8"),
headers={"Content-Type": "application/json", "X-Request-Id": request_id},
method="POST",
)
with request.urlopen(upstream_request, timeout=UPSTREAM_TIMEOUT_SECONDS) as response:
status = response.status
content_type = response.headers.get_content_type()
if content_type != "text/event-stream":
status = 502
self.send_json(status, {"error": "model_server_invalid_content_type"}, request_id)
return
self.send_response(status)
self.send_header("Content-Type", "text/event-stream")
self.send_header("Cache-Control", "no-cache")
self.send_header("Transfer-Encoding", "chunked")
self.send_header("X-Request-Id", request_id)
self.end_headers()
response_started = True
while chunk := response.read1(64 * 1024):
self.wfile.write(f"{len(chunk):X}\r\n".encode("ascii"))
self.wfile.write(chunk)
self.wfile.write(b"\r\n")
self.wfile.flush()
self.wfile.write(b"0\r\n\r\n")
self.wfile.flush()
except error.HTTPError as exc:
status = 502
self.send_json(status, {"error": "model_server_error", "upstream_status": exc.code}, request_id)
except (BrokenPipeError, ConnectionResetError, http.client.IncompleteRead, http.client.RemoteDisconnected) as exc:
LOG.warning("event=stream_interrupted request_id=%s type=%s", request_id, type(exc).__name__)
self.close_connection = True
except (error.URLError, TimeoutError, OSError) as exc:
if response_started:
LOG.warning("event=stream_interrupted request_id=%s type=%s", request_id, type(exc).__name__)
self.close_connection = True
else:
status = 504
self.send_json(status, {"error": "model_server_timeout_or_unavailable"}, request_id)
finally:
REQUEST_LOCK.release()
finally:
elapsed_ms = round((time.monotonic() - started) * 1000)
LOG.info("event=invocation request_id=%s status=%s elapsed_ms=%s", request_id, status, elapsed_ms)
def shutdown(_signum: int, _frame: object) -> None:
if SHUTTING_DOWN.is_set():
return
SHUTTING_DOWN.set()
LOG.info("event=shutdown_requested")
if CHILD is not None and CHILD.poll() is None:
CHILD.terminate()
try:
CHILD.wait(timeout=20)
except subprocess.TimeoutExpired:
CHILD.kill()
def main() -> None:
signal.signal(signal.SIGTERM, shutdown)
signal.signal(signal.SIGINT, shutdown)
server = ThreadingHTTPServer(("0.0.0.0", PORT), Handler)
server.daemon_threads = True
server.timeout = 1
LOG.info("event=adapter_started port=%s", PORT)
startup_thread = threading.Thread(target=start_child, name="llama-startup", daemon=True)
startup_thread.start()
child_exit_reported = False
try:
while not SHUTTING_DOWN.is_set():
server.handle_request()
if START_LLAMA_SERVER and CHILD is not None and CHILD.poll() is not None and not child_exit_reported:
LOG.error("event=llama_child_stopped exit_code=%s", CHILD.returncode)
child_exit_reported = True
finally:
server.server_close()
shutdown(signal.SIGTERM, None)
LOG.info("event=adapter_stopped")
if __name__ == "__main__":
main()
コンテナをビルドして ECR へ push した後、タグではなく digest を取得します。
aws ecr create-repository \
--region us-west-2 \
--repository-name <ECR_REPOSITORY_NAME>
aws ecr get-login-password --region us-west-2 \
| docker login --username AWS --password-stdin <ACCOUNT_ID>.dkr.ecr.us-west-2.amazonaws.com
docker build --platform linux/amd64 -t <ECR_REPOSITORY_URI>:qwen36 .
docker push <ECR_REPOSITORY_URI>:qwen36
aws ecr describe-images \
--region us-west-2 \
--repository-name <ECR_REPOSITORY_NAME> \
--image-ids imageTag=qwen36 \
--query 'imageDetails[0].imageDigest'
Terraform
実行ロールには、対象のモデルオブジェクト、ECR リポジトリ、CloudWatch Logs だけを許可します。Terraform を実行する主体には、SageMaker リソースの操作と、この実行ロールに限定した iam:PassRoleが別途必要です。
Terraform の全文
terraform {
required_version = "= 1.15.8"
required_providers {
aws = {
source = "hashicorp/aws"
version = "= 6.49.0"
}
}
backend "s3" {}
}
variable "aws_region" {
type = string
default = "us-west-2"
}
variable "aws_profile" {
type = string
}
variable "name_prefix" {
type = string
default = "qwen36-gguf-poc"
}
variable "artifact_bucket_name" {
type = string
}
variable "artifact_bucket_arn" {
type = string
}
variable "model_object_key" {
type = string
}
variable "model_s3_version_id" {
type = string
}
variable "model_sha256" {
type = string
default = "4ac6a06bce551257267f49ad2226f8671a22519ccc1a4dde9d5b433d1f2a410d"
}
variable "model_size_bytes" {
type = number
default = 21166757728
}
variable "ecr_repository_name" {
type = string
}
variable "ecr_repository_arn" {
type = string
}
variable "container_image_uri" {
type = string
validation {
condition = can(regex("@sha256:[0-9a-f]{64}$", var.container_image_uri))
error_message = "container_image_uri must be pinned by digest."
}
}
provider "aws" {
region = var.aws_region
profile = var.aws_profile
}
data "aws_s3_object" "model" {
bucket = var.artifact_bucket_name
key = var.model_object_key
version_id = var.model_s3_version_id
}
data "aws_s3_object" "current_model" {
bucket = var.artifact_bucket_name
key = var.model_object_key
}
data "aws_ecr_image" "container" {
repository_name = var.ecr_repository_name
image_digest = split("@", var.container_image_uri)[1]
}
check "immutable_inputs" {
assert {
condition = (
data.aws_s3_object.model.version_id == var.model_s3_version_id &&
data.aws_s3_object.model.content_length == var.model_size_bytes &&
data.aws_s3_object.current_model.version_id == var.model_s3_version_id &&
data.aws_s3_object.current_model.content_length == var.model_size_bytes &&
data.aws_ecr_image.container.image_digest == split("@", var.container_image_uri)[1]
)
error_message = "The current model object or image differs from the reviewed input."
}
}
resource "aws_cloudwatch_log_group" "endpoint" {
name = "/aws/sagemaker/Endpoints/${var.name_prefix}"
retention_in_days = 7
}
data "aws_iam_policy_document" "assume" {
statement {
actions = ["sts:AssumeRole"]
principals {
type = "Service"
identifiers = ["sagemaker.amazonaws.com"]
}
}
}
resource "aws_iam_role" "sagemaker" {
name = "${var.name_prefix}-execution"
assume_role_policy = data.aws_iam_policy_document.assume.json
}
data "aws_iam_policy_document" "execution" {
statement {
actions = ["s3:GetObject"]
resources = ["${var.artifact_bucket_arn}/${var.model_object_key}"]
}
statement {
actions = ["s3:GetBucketLocation"]
resources = [var.artifact_bucket_arn]
}
statement {
actions = ["ecr:GetAuthorizationToken"]
resources = ["*"]
}
statement {
actions = [
"ecr:BatchCheckLayerAvailability",
"ecr:BatchGetImage",
"ecr:GetDownloadUrlForLayer",
]
resources = [var.ecr_repository_arn]
}
statement {
actions = [
"logs:CreateLogStream",
"logs:DescribeLogStreams",
"logs:PutLogEvents",
]
resources = [
aws_cloudwatch_log_group.endpoint.arn,
"${aws_cloudwatch_log_group.endpoint.arn}:*",
]
}
}
resource "aws_iam_role_policy" "execution" {
name = "${var.name_prefix}-execution"
role = aws_iam_role.sagemaker.id
policy = data.aws_iam_policy_document.execution.json
}
resource "aws_sagemaker_model" "llama" {
name = "${var.name_prefix}-model"
execution_role_arn = aws_iam_role.sagemaker.arn
enable_network_isolation = true
primary_container {
image = var.container_image_uri
mode = "SingleModel"
model_data_source {
s3_data_source {
compression_type = "None"
s3_data_type = "S3Object"
s3_uri = "s3://${var.artifact_bucket_name}/${var.model_object_key}"
}
}
environment = {
MODEL_FILE = "Qwen3.6-35B-A3B-Q4_K_M.gguf"
MODEL_ARTIFACT_SIZE_BYTES = tostring(var.model_size_bytes)
MODEL_ARTIFACT_SHA256 = var.model_sha256
CONTEXT_SIZE = "16384"
PARALLEL = "1"
GPU_LAYERS = "999"
UPSTREAM_TIMEOUT_SECONDS = "480"
}
}
depends_on = [aws_iam_role_policy.execution]
}
resource "aws_sagemaker_endpoint_configuration" "realtime" {
name = "${var.name_prefix}-config"
production_variants {
variant_name = "AllTraffic"
model_name = aws_sagemaker_model.llama.name
initial_instance_count = 1
instance_type = "ml.g5.2xlarge"
model_data_download_timeout_in_seconds = 1800
container_startup_health_check_timeout_in_seconds = 1800
}
}
resource "aws_sagemaker_endpoint" "realtime" {
name = var.name_prefix
endpoint_config_name = aws_sagemaker_endpoint_configuration.realtime.name
depends_on = [aws_cloudwatch_log_group.endpoint]
}
output "endpoint_name" {
value = aws_sagemaker_endpoint.realtime.name
}
aws_profile = "<AWS_PROFILE>"
artifact_bucket_name = "<PRIVATE_MODEL_BUCKET>"
artifact_bucket_arn = "arn:aws:s3:::<PRIVATE_MODEL_BUCKET>"
model_object_key = "models/Qwen3.6-35B-A3B-Q4_K_M.gguf"
model_s3_version_id = "<S3_VERSION_ID>"
ecr_repository_name = "<ECR_REPOSITORY_NAME>"
ecr_repository_arn = "arn:aws:ecr:us-west-2:<ACCOUNT_ID>:repository/<ECR_REPOSITORY_NAME>"
container_image_uri = "<ACCOUNT_ID>.dkr.ecr.us-west-2.amazonaws.com/<ECR_REPOSITORY_NAME>@sha256:<IMAGE_DIGEST>"
リモート state を使う場合の backend 設定例です。state 用バケットは事前に作成し、バージョニング、暗号化、Block Public Access を有効にします。
bucket = "<TERRAFORM_STATE_BUCKET>"
key = "qwen36-gguf-poc/terraform.tfstate"
region = "us-west-2"
encrypt = true
use_lockfile = true
terraform init -backend-config=backend.s3.tfbackend
terraform fmt -check
terraform validate
terraform plan -out=realtime.tfplan
terraform apply realtime.tfplan
ストリーミング呼び出し
AWS SDK for JavaScript v3 の @aws-sdk/client-sagemaker-runtimeを使用します。
invoke-stream.ts の全文
import {
InvokeEndpointWithResponseStreamCommand,
SageMakerRuntimeClient,
type ResponseStream,
} from "@aws-sdk/client-sagemaker-runtime";
function isPayloadPart(
event: ResponseStream,
): event is ResponseStream.PayloadPartMember {
return "PayloadPart" in event;
}
async function* readSse(
chunks: AsyncIterable<Uint8Array>,
): AsyncGenerator<string> {
const decoder = new TextDecoder("utf-8");
let buffer = "";
for await (const chunk of chunks) {
buffer += decoder.decode(chunk, { stream: true });
while (true) {
const match = buffer.match(/\r?\n\r?\n/);
if (!match || match.index === undefined) break;
const block = buffer.slice(0, match.index);
buffer = buffer.slice(match.index + match[0].length);
const data = block
.split(/\r?\n/)
.map((line) => line.match(/^data:\s?(.*)$/)?.[1])
.filter((line): line is string => line !== undefined)
.join("\n");
if (data) yield data;
}
}
buffer += decoder.decode();
if (buffer.trim()) {
const data = buffer
.split(/\r?\n/)
.map((line) => line.match(/^data:\s?(.*)$/)?.[1])
.filter((line): line is string => line !== undefined)
.join("\n");
if (data) yield data;
}
}
const endpointName = process.env.SAGEMAKER_ENDPOINT_NAME;
if (!endpointName) throw new Error("SAGEMAKER_ENDPOINT_NAME is required");
const client = new SageMakerRuntimeClient({ region: "us-west-2" });
const abortController = new AbortController();
try {
const response = await client.send(
new InvokeEndpointWithResponseStreamCommand({
EndpointName: endpointName,
ContentType: "application/json",
Accept: "text/event-stream",
Body: new TextEncoder().encode(JSON.stringify({
messages: [
{
role: "system",
content: "あなたは簡潔で安全な日本語アシスタントです。",
},
{
role: "user",
content: "雨の日に家で楽しめる創作活動を三つ、各一文で提案してください。",
},
],
stream: true,
max_tokens: 192,
temperature: 0.2,
top_p: 0.9,
seed: 42,
enable_thinking: false,
})),
}),
{ abortSignal: abortController.signal },
);
if (!response.Body) throw new Error("response body is missing");
const payloads = (async function* () {
for await (const event of response.Body!) {
if (!isPayloadPart(event) || !event.PayloadPart.Bytes) {
throw new Error("SageMaker response stream failed");
}
yield event.PayloadPart.Bytes;
}
})();
for await (const data of readSse(payloads)) {
if (data === "[DONE]") break;
const event = JSON.parse(data);
const delta = event?.choices?.[0]?.delta?.content;
if (typeof delta === "string") process.stdout.write(delta);
}
} finally {
client.destroy();
}
削除
検証後は、保存した plan ではなく現在の Terraform 設定と state を使って削除対象を確認します。
terraform plan -destroy
terraform destroy
aws sagemaker list-endpoints \
--region us-west-2 \
--name-contains qwen36-gguf-poc
aws sagemaker list-endpoint-configs \
--region us-west-2 \
--name-contains qwen36-gguf-poc
aws sagemaker list-models \
--region us-west-2 \
--name-contains qwen36-gguf-poc








