Qwen3.6-35B-A3B を Amazon SageMaker で動かしてストリーミング応答を得てみた

Qwen3.6-35B-A3B を Amazon SageMaker で動かしてストリーミング応答を得てみた

Qwen3.6 の Q4_K_M GGUF を ml.g5.2xlarge の 24 GB GPU へロードし、回答を逐次受信できました。ウォーム TTFT は 2,058 ms でした。
2026.07.26

はじめに

大規模な 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 の回答を逐次表示した様子です。

Next.js 検証用アプリから 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 を使って生成結果を逐次受信したい方

参考資料

検証構成と技術選定

要件は次の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をビルドします。

Dockerfile
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 として転送します。入力検証とエラー処理を含む全文は付録に掲載しています。

adapter.py
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 では、検証に使用したインスタンスとタイムアウトを固定しました。

main.tf
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 イベントの途中で分割されても復元します。

invoke-stream.ts
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 の呼び出しも中止します。

invoke-stream.ts
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 を直接呼び出して取得しました。

Next.js 検証用アプリから Qwen3.6 の回答が逐次表示される様子

計測値

この入力を使い、同じエンドポイントに対して正常完了とクライアント中断を確認しました。生成設定は max_tokens = 192temperature = 0.2top_p = 0.9seed = 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 の全文
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 の全文
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 の全文
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 の全文
versions.tf
terraform {
  required_version = "= 1.15.8"

  required_providers {
    aws = {
      source  = "hashicorp/aws"
      version = "= 6.49.0"
    }
  }

  backend "s3" {}
}
variables.tf
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."
  }
}
main.tf
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
}
terraform.tfvars
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 を有効にします。

backend.s3.tfbackend
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 の全文
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

AI白書2026 配布中

クラスメソッドが独自に行なったAI診断調査をもとに、企業のAI活用の現在地を調査レポートとしてまとめました。企業規模別の活用度傾向に加え、規模を超えてAI活用を進める企業に共通する取り組みまで、自社の現在地を捉えるためのヒントにぜひ。

AI白書2026

無料でダウンロードする

この記事をシェアする

関連記事