1.はじめに
本ポートフォリオでは「Go言語による高速なAPI受付」と「Pythonによる重い画像解析(EasyOCR)」を、Google CloudのPub/Subを介して完全に切り離した、非同期のOCRバックエンド基盤を構築しました。
この基盤構築に関する詳細説明を行います。
💡 本記事は、3部作の「詳細説明編」です。
・【改善編】詳細説明編:「Go/Python/Terraform/Helm」の詳細解説。
・【改善編】総括編: Pub/Subを用Push型データパスの要点・効果を解説。
・【改善編】検証編: 負荷試験(k6)による大量リクエスト流入時の挙動を検証・考察。
| ◆構成 |
|---|
| 1.はじめに |
| 2.技術スタックとシステム構成 |
| 3.APIサーバー(Go) |
| 4.推論サーバ(python) |
| 5.Terraformによるインフラ自動化 |
| 6.Helmによるワークロード管理 |
| 7.デプロイ手順 |
| 8.おわりに |
2.技術スタックとシステム構成
2-1. 技術スタック
本システムでは、以下の技術スタックを採用しています。前回の負荷集中の課題を解決するため、「疎結合」「非同期バッファリング」「イベント駆動型オートスケール」 を軸とした構成へ刷新しました。
| カテゴリ | 技術・ツール | 用途 |
|---|---|---|
| Frontend API | Go (Standard) | HTTPリクエスト受付、GCSアップロード、Pub/Subメッセージ送信、構造化ログ出力 |
| Inference Engine | Python (EasyOCR) | Pub/Subメッセージ受信(Pull)、GCS画像ダウンロード、文字認識処理 |
| Messaging (Buffer) | Cloud Pub/Sub | Go-Python間の通信を仲介するメッセージキュー。スパイク負荷のバッファリング |
| Autoscaling | KEDA | Pub/Subの未処理メッセージ数(Backlog)を検知し、Podを高速にイベント駆動スケール |
| Infrastructure | GKE (Standard) | 完全管理型K8s環境。スポットインスタンスを活用した超低コスト運用 |
| IaC | Terraform | VPC、Cloud NAT、GKE、Pub/Sub、GCS、Log Sinkのコード管理(IaC) |
| Observability | Cloud Logging / BigQuery | 構造化ログの集約、およびSQLによる推論精度・性能分析 |
| Testing | k6 | 大量リクエストを用いた負荷試験。キューの滞留とKEDAのスケーリング挙動の検証 |
2-2. 主要ライブラリ
新アーキテクチャのパフォーマンスとコスト最適化を最大化するために採用した主要ライブラリ・ツールです。
| 対象 | ライブラリ | 採用理由 |
|---|---|---|
| Go | google.com | GCSへの画像アップロードおよびPub/Subへの超高速なメッセージ発行(Publish)のため |
| Go | gopsutil | 1リクエストの処理期間中における純粋なCPU使用率を正確に計測し、ログへ含めるため |
| Python | EasyOCR | PyTorchベースで精度が高く、日本語・英語に標準対応。今回は起動時にプレロードしてコールドスタートを排除。 |
| Python | psutil | OCR推論時のCPU負荷をノンブロッキングで計測し、BigQuery分析用の構造化ログに出力するため。 |
| K8s | KEDA (GCP Pub/Sub Scaler) | Pub/Subのキュー滞留数に応じてPod数を「0台⇄最大5台」までミリ秒単位で高感度スケールさせるため |
| GCP | Workload Identity | サービスアカウントの秘密鍵(JSON)の発行を完全に撤廃し、Pod単位で安全にGCP権限を借用するため |
2-3. ディレクトリ構成
インフラ管理(Terraform)コマンドで、デプロイできる以下の構成となっています。
このデータセットは、GitHubレポジトリ に登録しています。
.
├── terraform/ # GKE, VPC, Artifact Registry, Log Sinkの定義
├── chart/ # Kubernetesデプロイ用リソース (Helm Chart形式)
│ └── templates/ # Deployment, Service, HPA, RBACのマニフェスト
├── go-api/ # フロントエンドAPI (Go / Gin)
├── python-api/ # 推論エンジン (Python / EasyOCR / gRPC Server)
├── pb/ # gRPC定義ファイル (.proto) と自動生成コード
└── k6/ # 負荷試験スクリプト
└── test_images/ # 検証用画像および正解データ(mapping.json)
2-3. システム構成
1. 「データフロー」処理フェーズ(①〜⑪)
2. 「デプロイ・通信」フェーズ(①〜⑤)
3. 「認証・セキュリティ」フェーズ(①〜⑤)
3.APIサーバー(Go)
本システムの入り口となるGo APIサーバー(uploadHandler)は、クライアントからの画像アップロードを受け付け、クラウド基盤(GCS/Pub/Sub)へ迅速にタスクを委譲します。同期処理のボトルネックを徹底的に排除・可視化するため、処理の裏側でミリ秒単位のタイムスタンプ計測とピンポイントなCPU使用率のプロファイリングを同時に実行しています。全体の処理フローと計測タイミングは以下の通りです。
3-1. 構造化ログのデータ設計
BigQueryでの検索・集計効率を最大化するため、ログをテキスト(文字列)ではなくJSON形式で一貫して管理する構造体 CloudLogPayload を定義しています。
// Cloud Logging 兼 BigQuery 蓄積用の構造化ログ定義
type CloudLogPayload struct {
RequestID string `json:"requestid"`
Step string `json:"step"`
DurationStart int64 `json:"duration_start"` // GCS書き込み完了までのミリ秒
DurationEnd int64 `json:"duration_end"` // Pub/Sub送信完了までのミリ秒
CpuUsage float64 `json:"cpuusage"`
PodName string `json:"podname"`
Expected string `json:"expected"`
}
・トレーサビリティの確保:
RequestID や、処理位置を示す Step(この場合は "go")をログに含めることで、分散システム内でも1本のリクエストを確実に追跡できるようにしています。
・運用に役立つメトリクス:
パフォーマンス分析のため、GCS書き込み完了時点とPub/Sub送信完了時点の時間をミリ秒(int64)で保持します。
・インフラ負荷の可視化:
どのPod(PodName)で、どれだけのCPU(CpuUsage)を消費したかを同時に記録し、リソースの逼迫度を追跡可能にしています。
3-2. パフォーマンス計測と外的要因の排除
HTTPハンドラー(uploadHandler)の開始直後に、正確なプロファイリングを行うための初期化を行っています。
func uploadHandler(w http.ResponseWriter, r *http.Request) {
// 全体の処理時間を計測するためのスタート時間
startTime := time.Now()
// PodのCPUストップウォッチを「リセット(初期化)」
_, _ = cpu.Percent(0, false)
// 【外的要因の排除】現在の絶対時刻を「秒.ミリ秒」の文字列(例: "1717693201.234") として生成
goRequestTime := fmt.Sprintf("%f", float64(startTime.UnixNano())/1e9)
・正確なCPU使用率の計測:
github.com/shirou/gopsutil/v4/cpu を使用しています。ハンドラーの先頭で一度関数を呼び出して内部カウンターをリセット(基準点を作成)することで、リクエスト処理中だけの純粋なCPU負荷を測定できるようにしています。
・分散システム間の時刻同期:
startTime.UnixNano() をベースに、ナノ秒精度から秒・ミリ秒単位の浮動小数点文字列(goRequestTime)を生成しています。これをPub/Subの属性として後続のPython Workerへ引き継ぐことで、ネットワーク遅延に左右されない「システム全体の本当の処理遅延」を測定可能にしています。
3-3. 補助関数によるJSONログの標準出力
Google Cloud(GKEやCloud Runなど)の環境では、コンテナが標準出力(stdout)に書き出したJSON文字列を、Cloud Loggingが自動的に構造化ログ(jsonPayload)としてキャプチャしてくれます。
// 構造化JSONログを標準出力に吐き出す補助関数
func writeCloudLog(payload CloudLogPayload) {
jsonData, err := json.Marshal(payload)
if err != nil {
log.Printf("❌ ログのJSON化に失敗: %v", err)
return
}
fmt.Println(string(jsonData))
}
3-4.全コード
コード詳細は、こちらを参照して下さい。
./go/main.go
package main
import (
"context"
"encoding/json"
"fmt"
"io"
"log"
"net/http"
"os"
"path/filepath"
"time"
"github.com/shirou/gopsutil/v4/cpu"
"cloud.google.com/go/pubsub"
"cloud.google.com/go/storage"
)
// Cloud Logging 兼 BigQuery 蓄積用の構造化ログ定義
type CloudLogPayload struct {
RequestID string `json:"RequestID"`
Step string `json:"Step"`
DurationStart int64 `json:"duration_start"` // GCS書き込み完了までのミリ秒
DurationEnd int64 `json:"duration_end"` // Pub/Sub送信完了までのミリ秒
CPUUsage float64 `json:"CPUUsage"`
PodName string `json:"PodName"`
Expected string `json:"Expected"`
}
// 構造化JSONログを標準出力に吐き出す補助関数
func writeCloudLog(payload CloudLogPayload) {
jsonData, err := json.Marshal(payload)
if err != nil {
log.Printf("❌ ログのJSON化に失敗: %v", err)
return
}
fmt.Println(string(jsonData))
}
func main() {
http.HandleFunc("/upload", uploadHandler)
fmt.Println("[Go API] サーバーがポート 8080 で起動しました。")
if err := http.ListenAndServe(":8080", nil); err != nil {
fmt.Printf("サーバー起動エラー: %v\n", err)
}
}
func uploadHandler(w http.ResponseWriter, r *http.Request) {
// 全体の処理時間を計測するためのスタート時間
startTime := time.Now()
// PodのCPUストップウォッチを「リセット(初期化)」
_, _ = cpu.Percent(0, false)
// 【外的要因の排除】現在の絶対時刻を「秒.ミリ秒」の文字列(例: "1717693201.234")として生成
goRequestTime := fmt.Sprintf("%f", float64(startTime.UnixNano())/1e9)
if r.Method != http.MethodPost {
http.Error(w, "POSTメソッドのみ受け付けます", http.StatusMethodNotAllowed)
return
}
projectID := os.Getenv("GCP_PROJECT_ID")
podName := os.Getenv("POD_NAME")
if podName == "" {
podName = "local-go-pod"
}
topicID := "ocr-task-topic"
bucketName := fmt.Sprintf("ocr-images-bucket-%s", projectID)
if projectID == "" {
log.Println("❌ エラー: 環境変数 GCP_PROJECT_ID が設定されていません")
http.Error(w, "サーバー設定エラー (GCP_PROJECT_ID が不足しています)", http.StatusInternalServerError)
return
}
requestID := r.FormValue("request_id")
expectedText := r.FormValue("expected")
if requestID == "" {
requestID = fmt.Sprintf("gen-%d", time.Now().UnixNano())
}
file, header, err := r.FormFile("image")
if err != nil {
http.Error(w, "ファイルの取得に失敗しました: "+err.Error(), http.StatusBadRequest)
return
}
defer file.Close()
ext := filepath.Ext(header.Filename)
if ext == "" {
ext = ".png"
}
fileName := fmt.Sprintf("%s%s", requestID, ext)
ctx := context.Background()
// GCSへのアップロード
storageClient, err := storage.NewClient(ctx)
if err != nil {
log.Printf("❌ GCSクライアント作成失敗: %v\n", err)
http.Error(w, "GCSクライアント作成失敗: "+err.Error(), http.StatusInternalServerError)
return
}
defer storageClient.Close()
gcsObject := storageClient.Bucket(bucketName).Object(fileName)
gcsWriter := gcsObject.NewWriter(ctx)
if _, err := io.Copy(gcsWriter, file); err != nil {
log.Printf("❌ GCS書き込み失敗: %v\n", err)
http.Error(w, "GCS書き込み失敗", http.StatusInternalServerError)
return
}
gcsWriter.Close()
// GCS書き込み完了までの期間をミリ秒で計測
durationStart := time.Since(startTime).Milliseconds()
realGCSPath := fmt.Sprintf("gs://%s/%s", bucketName, fileName)
// 2. Pub/Subへのメッセージ送信
pubsubClient, err := pubsub.NewClient(ctx, projectID)
if err != nil {
log.Printf("❌ Pub/Subクライアント作成失敗: %v\n", err)
http.Error(w, "Pub/Subクライアント作成失敗: "+err.Error(), http.StatusInternalServerError)
return
}
defer pubsubClient.Close()
t := pubsubClient.Topic(topicID)
result := t.Publish(ctx, &pubsub.Message{
Data: []byte(realGCSPath),
Attributes: map[string]string{
"request_id": requestID,
"expected": expectedText,
"go_request_time": goRequestTime,
},
})
// 送信完了を待機
_, err = result.Get(ctx)
if err != nil {
log.Printf("❌ Pub/Sub送信失敗: %v\n", err)
http.Error(w, "Pub/Sub送信失敗", http.StatusInternalServerError)
return
}
// 送り出し(Pub/Sub送信完了)までの期間をミリ秒で計測
durationEnd := time.Since(startTime).Milliseconds()
// 純粋なリクエスト処理期間(数十〜数百ミリ秒間)の平均CPU使用率を取得
percent, _ := cpu.Percent(0, false)
var cpuVal float64
if len(percent) > 0 {
cpuVal = percent[0]
}
// クラウドロギング(構造化ログ)の出力
writeCloudLog(CloudLogPayload{
RequestID: requestID,
Step: "go",
DurationStart: durationStart,
DurationEnd: durationEnd,
CPUUsage: cpuVal,
PodName: podName,
Expected: expectedText,
})
w.WriteHeader(http.StatusOK)
w.Write([]byte("画像アップロード & パイプライン投入 成功!\n"))
}
4.推論サーバ(python)
Python Worker(python_worker.py)は、Go APIがパブリッシュしたPub/Subメッセージをトリガーに動作する、非同期の機械学習(OCR)推論エンジンです。単にメッセージを処理するだけでなく、コンテナ起動時のオーバーヘッド削減、キューの滞留時間計測、そしてKubernetes環境での安全なPod縮退(オートスケール)に耐えるライフサイクルを構築しています。全体のフローは以下の通りです。
4-1.処理速度を考慮したモデルプレロード
重い機械学習モデル(EasyOCR)をリクエストごとにロードすると、大幅な遅延とCPU/メモリの無駄遣いが発生します。そのため、コンテナ起動時に1度だけグローバルメモリに展開(プレロード)する設計にしました。
# 【コスト最適化】起動時に1度だけグローバルでEasyOCRモデルをメモリに展開(プレロード)
# AutopilotのCPU環境を想定し、明示的に gpu=False を指定して無駄な探索を省きます
print("[Python Worker] EasyOCRモデルを初期化中 (日本語/英語)...", flush=True)
READER = easyocr.Reader(['ja', 'en'], gpu=False)
・ウォームスタートの実現:
コンテナが起動した時点でモデルをロードしておくことで、Pub/Subメッセージ受信からOCR推論開始までのタイムラグを極限まで減らしています。
4-2.Kubernetes / KEDA 連携のための Graceful Shutdown
Kubernetes環境、特にKEDA(Kubernetes Event-driven Autoscaling)による自動スケールイン(ノードの縮小)が発生した際、処理中のメッセージを強制終了させずに安全に停止させる仕組み(Graceful Shutdown) を実装しました。
# Graceful Shutdown 用の終了フラグ
shutdown_requested = False
def receive_sigterm(signum, frame):
"""KEDAやKubernetesからの終了シグナルを安全に検知する"""
global shutdown_requested
print("[Python Worker] KEDAからの終了シグナル(SIGTERM)を検知しました。現在のメッセージ処理が完了次第、安全に停止します。")
shutdown_requested = True
def main():
signal.signal(signal.SIGTERM, receive_sigterm)
signal.signal(signal.SIGINT, receive_sigterm)
while not shutdown_requested:
has_message = pull_message()
if has_message:
continue
time.sleep(0.5) # ポーリング負荷による無駄なノードCPU消費(=課金)を抑える
・シグナル捕捉(SIGTERM / SIGINT):
KEDAやKubernetesがPodを削除する際、最初にコンテナに送られる SIGTERM シグナルを検知します。
・データロストの防止:
シグナルを受け取っても即座に sys.exit() せず、フラグ(shutdown_requested)を立てるだけに留めます。これにより、pull_message() 内で 現在実行中のOCR処理やGCSへの成果物保存、Pub/SubへのAck送信をすべて安全に完了させてから、 ループを抜けて正常終了(sys.exit(0))します。
4-3.全コード
コード詳細は、こちらを参照して下さい。
./python/worker.py
import os
import time
import signal
import sys
import json
import psutil
import urllib.request
from google.cloud import storage
from google.cloud import pubsub_v1
import easyocr
# グローバルでクライアントを1度だけ初期化
PROJECT_ID = os.environ.get("GCP_PROJECT_ID")
SUBSCRIBER = pubsub_v1.SubscriberClient()
STORAGE_CLIENT = storage.Client()
POD_NAME = os.environ.get("POD_NAME", "local-python-pod")
# 起動時に1度だけグローバルでEasyOCRモデルをメモリに展開(プレロード)
# AutopilotのCPU環境を想定し、明示的に gpu=False を指定して無駄な探索を省きます
print("[Python Worker] EasyOCRモデルを初期化中 (日本語/英語)...", flush=True)
READER = easyocr.Reader(['ja', 'en'], gpu=False)
print("[Python Worker] EasyOCRの準備が完了しました。メッセージを待機します。", flush=True)
# Graceful Shutdown 用の終了フラグと psutil インスタンスの初期化
shutdown_requested = False
PROCESS_MONITOR = psutil.Process(os.getpid())
PROCESS_MONITOR.cpu_percent(interval=None) # 初回呼び出しでカウンターを初期化
def receive_sigterm(signum, frame):
"""KEDAやKubernetesからの終了シグナルを安全に検知する"""
global shutdown_requested
print("[Python Worker] KEDAからの終了シグナル(SIGTERM)を検知しました。現在のメッセージ処理が完了次第、安全に停止します。")
shutdown_requested = True
def write_cloud_log(payload):
"""Cloud Logging / BigQueryに自動マッピングされる構造化JSONログを標準出力に吐き出す"""
print(json.dumps(payload), flush=True) # flush=True でリアルタイム転送を確実化
def get_cpu_usage():
"""現在のPythonプロセスの実際のCPU使用率(%)をノンブロッキングで計測する"""
try:
return PROCESS_MONITOR.cpu_percent(interval=None)
except Exception as e:
print(f"⚠ CPU使用率の計測に失敗: {e}", file=sys.stderr)
return 0.0
def download_image_from_gcs(gcs_path):
"""GCSから画像をローカル環境にダウンロードする"""
path_without_scheme = gcs_path.replace("gs://", "")
parts = path_without_scheme.split("/", 1)
bucket_name, blob_name = parts[0], parts[1]
local_filename = os.path.basename(blob_name)
bucket = STORAGE_CLIENT.bucket(bucket_name)
blob = bucket.blob(blob_name)
blob.download_to_filename(local_filename)
print(f"[Python Worker] ダウンロード完了: {os.path.abspath(local_filename)}")
return local_filename
def save_result_to_gcs(gcs_path, detected_text):
"""【修正版】EasyOCRが解析したテキストを結果保存用バケットに保存する"""
result_bucket_name = f"ocr-results-bucket-{PROJECT_ID}"
original_filename = os.path.basename(gcs_path)
base, _ = os.path.splitext(original_filename)
result_filename = f"{base}_result.txt"
bucket = STORAGE_CLIENT.bucket(result_bucket_name)
blob = bucket.blob(result_filename)
log_content = (
f"Timestamp: {time.strftime('%Y-%m-%d %H:%M:%S')}\n"
f"FileName: {gcs_path}\n"
f"Status: SUCCESS\n"
f"DetectedText:\n{detected_text}\n"
)
blob.upload_from_string(log_content, content_type="text/plain; charset=utf-8")
print(f"[Python Worker] 成果物をGCSに保存しました: gs://{result_bucket_name}/{result_filename}")
def pull_message():
subscription_id = "ocr-task-sub"
subscription_path = SUBSCRIBER.subscription_path(PROJECT_ID, subscription_id)
try:
response = SUBSCRIBER.pull(
request={"subscription": subscription_path, "max_messages": 1, "return_immediately": True}
)
except Exception as e:
print(f"Pub/Subからの受信エラー: {e}", file=sys.stderr)
return False
if not response.received_messages:
return False
received_message = response.received_messages[0]
ack_id = received_message.ack_id
message = received_message.message
# Go側から渡された時刻を復元
go_request_time = float(message.attributes.get("go_request_time", time.time()))
gcs_path = message.data.decode("utf-8")
request_id = message.attributes.get("request_id", f"gen-py-{int(time.time()*1000)}")
expected_text = message.attributes.get("expected", "")
print(f"\n[Python Worker] 受信成功! RequestID: {request_id}, GCSパス: {gcs_path}")
local_filename = None
try:
# 1. GCSダウンロードを行う「前」に受信完了時間を計測
duration_start = int((time.time() - go_request_time) * 1000)
# 2. ダウンロード処理
local_filename = download_image_from_gcs(gcs_path)
# CPUタイマーのリセット
get_cpu_usage()
# 3. メイン処理(本物のEasyOCRによる推論)
print(f"[Python Worker] 本物のEasyOCRモデルで解析中: {local_filename} ...")
ocr_results = READER.readtext(local_filename, detail=0)
detected_text = " ".join(ocr_results)
print(f"[Python Worker] OCR解析結果: {detected_text}")
# CPUタイマーの計測完了
cpu_val = get_cpu_usage()
# 4. 成果物保存
save_result_to_gcs(gcs_path, detected_text)
# 使い終わったローカルの画像ファイルは即座に削除してディスク圧迫を防ぐ
if local_filename and os.path.exists(local_filename):
os.remove(local_filename)
# 成功時のみAckを送信
SUBSCRIBER.acknowledge(request={"subscription": subscription_path, "ack_ids": [ack_id]})
print("[Python Worker] Ack送信完了。キューから削除されました。\n")
# 5. すべてが完了した瞬間を記録
duration_end = int((time.time() - go_request_time) * 1000)
# Go側の期待値と実際のOCR結果が部分一致しているか判定(1: 一致, 0: 不一致)
is_match = 1 if expected_text in detected_text else 0
# クラウドロギング(構造化ログ)の出力
write_cloud_log({
"RequestID": request_id,
"Step": "python",
"duration_start": duration_start,
"duration_end": duration_end,
"CPUUsage": cpu_val,
"PodName": POD_NAME,
"Expected": expected_text,
"Detected": detected_text,
"IsMatch": is_match
})
return True
except Exception as e:
print(f"[重大エラー] 処理全体、またはGCS操作でエラーが発生: {e}", file=sys.stderr)
# エラー発生時も一時ファイルが残っていれば安全に消去する
if local_filename and os.path.exists(local_filename):
os.remove(local_filename)
# 処理が失敗したためメッセージを即座にPub/Subに返却(Nack)
SUBSCRIBER.modify_ack_deadline(
request={"subscription": subscription_path, "ack_ids": [ack_id], "ack_deadlines": [0]}
)
print("[Python Worker] エラーのためメッセージをキューに即座に返却しました。\n", file=sys.stderr)
return True # 処理サイクル自体は終わったためTrueを返し、メインループの判定へ繋ぐ
def main():
"""プログラム全体のメインエントリーポイント"""
signal.signal(signal.SIGTERM, receive_sigterm)
signal.signal(signal.SIGINT, receive_sigterm)
print("[Python Worker] 起動しました。ループを開始します。")
# shutdown_requested が True になるまで安全にループを実行
while not shutdown_requested:
has_message = pull_message()
if has_message:
continue
# ポーリング負荷による無駄なノードCPU消費(=課金)を抑えるため待機
time.sleep(0.5)
print("[Python Worker] 安全にシャットダウンを完了しました。プロセスを終了します。")
sys.exit(0)
if __name__ == "__main__":
if not PROJECT_ID:
print("❌ エラー: 環境変数 GCP_PROJECT_ID が設定されていません。", file=sys.stderr)
sys.exit(1)
main()
5.Terraformによるインフラ自動化
システムの土台となるGCPリソースはすべてTerraformで構成(IaC:Infrastructure as Code)しています。インフラをコード化するにあたり、単にリソースを並べるだけでなく、「Workload Identityによる安全な認証基盤」と「ログ転送の自動化」を導入しました。
5-1.Workload Identityによる安全な認証基盤
本システムでは、鍵の紛失や漏洩のリスクを回避するため、GKEの推奨機能である Workload Identity を採用し、KubernetesのServiceAccountとGCPのIAMサービスアカウントを安全にバインドしています。
# GKE内部設定 (Namespace & ServiceAccountの作成)
resource "kubernetes_namespace" "keda" {
metadata {
name = "keda"
}
depends_on = [google_container_cluster.primary]
}
# アプリが身分を証明するための会員証(ServiceAccount)を作成
resource "kubernetes_service_account" "app_k8s_sa" {
metadata {
name = "ocr-app-sa"
namespace = "default"
annotations = {
# この注釈(アノテーション)を入れることで、GCPのSAとガチッとリンクします
"iam.gke.io/gcp-service-account" = google_service_account.app_sa.email
}
}
depends_on = [google_container_cluster.primary]
}
・認証情報の隠蔽(脱・秘密鍵)
クラスタ内部に生のGCP秘密鍵(JSON)が一切存在しない状態を作れるため、漏洩リスクを根本から排除しています。
・最小権限の原則(ポッド単位の制御)
作成した ocr-app-sa をPodに割り当てることで、該当コンテナ(Go/Python)に必要なPub/SubやGCSへのアクセス権限だけを、インフラ側から安全に委譲しています。
5-2.Cloud LoggingからBigQueryへのログ転送
GoやPythonのコンテナが標準出力(stdout)に吐き出した構造化JSONログは、Cloud Loggingを介して自動的にBigQueryへとリアルタイム転送(シンク)されるようにパイプラインを構築しました。
# ログを蓄積するための BigQuery データセットを作成
resource "google_bigquery_dataset" "logging_dataset" {
dataset_id = var.dataset_id
friendly_name = "GKE OCR Pipeline Logs"
description = "Go と Python から出力されるカスタムログの蓄積用データセット"
location = var.region
delete_contents_on_destroy = true # 試験環境用:terraform destroy時にテーブルも削除
}
# ログルーターシンクの定義(GKEの特定のログのみを抽出してBigQueryに送る)
resource "google_logging_project_sink" "bq_log_sink" {
name = "ocr-pipeline-log-sink"
description = "GoとPythonコンテナの標準出力をBigQueryへ転送"
destination = "://googleapis.com{var.project_id}/datasets/${var.dataset_id}"
# フィルタ:GKE上の指定したコンテナ、またはログ内に特定のキー(RequestID)が含まれるものを対象にする
filter = <<EOT
resource.type="k8s_container"
jsonPayload.RequestID:*
EOT
unique_writer_identity = true
}
# ログルーター(シンク)が BigQuery にデータを書き込めるように権限(IAM)を付与
resource "google_project_iam_member" "sink_bq_writer" {
project = var.project_id
role = "roles/bigquery.dataEditor"
member = google_logging_project_sink.bq_log_sink.writer_identity
}
5-3. 全コード
コード詳細は、こちらを参照して下さい。
./terraform/main.tf
# ===============================
# 0. ネットワーク定義
# ===============================
# 専用のVPCネットワークを作成
resource "google_compute_network" "ocr_vpc" {
name = "ocr-pipeline-vpc"
auto_create_subnetworks = false # 自動で全リージョンにサブネットを作らない
}
# GKEを配置するための専用サブネットを作成
resource "google_compute_subnetwork" "ocr_subnet" {
name = "ocr-pipeline-subnet"
ip_cidr_range = "10.0.0.0/24" # 10.0.0.1 〜 10.0.0.254
region = var.region
network = google_compute_network.ocr_vpc.id
private_ip_google_access = true # Google API(GCSやPub/Sub)へのプライベート接続を許可
secondary_ip_range {
range_name = "pod-ranges"
ip_cidr_range = "192.168.16.0/20"
}
secondary_ip_range {
range_name = "services-ranges"
ip_cidr_range = "192.168.32.0/20"
}
}
# GKEノードがインターネットからコンテナを引っ張ったり、
# 外部と通信できるようにするためのNATゲートウェイ(ルーター)
resource "google_compute_router" "ocr_router" {
name = "ocr-pipeline-router"
region = var.region
network = google_compute_network.ocr_vpc.id
}
resource "google_compute_router_nat" "ocr_nat" {
name = "ocr-pipeline-nat"
router = google_compute_router.ocr_router.name
region = google_compute_router.ocr_router.region
nat_ip_allocate_option = "AUTO_ONLY"
source_subnetwork_ip_ranges_to_nat = "ALL_SUBNETWORKS_ALL_IP_RANGES"
}
# ===============================
# 1. プロバイダー & 共通データ定義
# ===============================
provider "google" {
project = var.project_id
region = var.region
}
data "google_client_config" "default" {}
provider "kubernetes" {
host = "https://${google_container_cluster.primary.endpoint}"
token = data.google_client_config.default.access_token
# cluster_ca_certificate = base64decode(google_container_cluster.primary.master_auth[0].cluster_ca_certificate)
cluster_ca_certificate = base64decode(
google_container_cluster.primary.master_auth[0].cluster_ca_certificate
)
}
# GKEノードプールが使用するデフォルトのコンピュートサービスアカウント
data "google_compute_default_service_account" "default" {}
# ===============================
# 2. クラウドストレージ (GCS バケット)
# ===============================
# 画像保存用バケット
resource "google_storage_bucket" "ocr_images" {
name = "ocr-images-bucket-${var.project_id}"
location = "ASIA-NORTHEAST1"
force_destroy = true
public_access_prevention = "enforced"
}
# 処理結果(テキスト)保存用バケット
resource "google_storage_bucket" "ocr_results" {
name = "ocr-results-bucket-${var.project_id}"
location = "ASIA-NORTHEAST1"
force_destroy = true
public_access_prevention = "enforced"
}
# ===============================
# 3. メッセージング (Pub/Sub)
# ===============================
resource "google_pubsub_topic" "ocr_task" {
name = "ocr-task-topic"
}
resource "google_pubsub_subscription" "ocr_task_sub" {
name = "ocr-task-sub"
topic = google_pubsub_topic.ocr_task.name
ack_deadline_seconds = 60
}
# ===============================
# 4. コンテナレジストリ (Artifact Registry)
# ===============================
resource "google_artifact_registry_repository" "ocr_repo" {
location = var.region
repository_id = "ocr-pipeline-repo"
format = "DOCKER"
}
# ===============================
# 5. GKE クラスター(Autopilotモード)
# ===============================
resource "google_container_cluster" "primary" {
name = "ocr-pipeline-cluster"
location = var.region # Autopilotは高可用性を担保するため「リージョン」指定になります
deletion_protection = false
# Autopilotモードを有効化(ノードプール定義は丸ごと削除可能になります)
enable_autopilot = true
# カスタムVPCとサブネットを指定
network = google_compute_network.ocr_vpc.id
subnetwork = google_compute_subnetwork.ocr_subnet.id
# 非公開(プライベート)クラスター構成
private_cluster_config {
enable_private_nodes = true
enable_private_endpoint = false # 外部(手元のPC)からkubectl操作を行えるようにする
}
}
# ===============================
# 5.5 GKE カスタムノードプール(Autopilot化に伴い削除)
# ===============================
# ===============================
# 6. IAM (サービスアカウントとGCP権限の定義)
# ===============================
# --- 6-1. KEDA専用サービスアカウント(変更なし) ---
resource "google_service_account" "keda_sa" {
account_id = "keda-metrics-scaler"
display_name = "KEDA Metrics Scaler Service Account"
}
resource "google_project_iam_member" "keda_pubsub_viewer" {
project = var.project_id
role = "roles/pubsub.viewer"
member = "serviceAccount:${google_service_account.keda_sa.email}"
}
resource "google_project_iam_member" "keda_pubsub_subscriber" {
project = var.project_id
role = "roles/pubsub.subscriber"
member = "serviceAccount:${google_service_account.keda_sa.email}"
}
resource "google_project_iam_member" "keda_monitoring_viewer" {
project = var.project_id
role = "roles/monitoring.viewer"
member = "serviceAccount:${google_service_account.keda_sa.email}"
}
# --- 6-2. アプリケーション専用サービスアカウント(変更なし) ---
resource "google_service_account" "app_sa" {
account_id = "ocr-app-worker"
display_name = "OCR Application App Worker Service Account"
}
resource "google_project_iam_member" "app_pubsub_publisher" {
project = var.project_id
role = "roles/pubsub.publisher"
member = "serviceAccount:${google_service_account.app_sa.email}"
}
resource "google_project_iam_member" "app_storage_admin" {
project = var.project_id
role = "roles/storage.objectAdmin"
member = "serviceAccount:${google_service_account.app_sa.email}"
}
resource "google_project_iam_member" "app_pubsub_subscriber" {
project = var.project_id
role = "roles/pubsub.subscriber"
member = "serviceAccount:${google_service_account.app_sa.email}"
}
# ===============================
# Workload Identity User
# ===============================
# ------------------------------------------------------------------------------
# 🌟 Workload Identity User (GKEの世界とGCPの世界の紐付け)
# ------------------------------------------------------------------------------
# ① KEDA用の紐付け1:作業者(kedaネームスペースの「keda-operator」という名前のK8sアカウントと結合)
resource "google_service_account_iam_member" "keda_wi" {
service_account_id = google_service_account.keda_sa.name
role = "roles/iam.workloadIdentityUser"
member = "serviceAccount:${var.project_id}.svc.id.goog[keda/keda-operator]"
# 🌟【追記】GKEクラスターとノードプールが完全に出来上がってから作らせる
depends_on = [
google_container_cluster.primary
]
}
# KEDA用の紐付け2:Metrics Server(集計係)用の Workload Identity 紐付け
resource "google_service_account_iam_member" "keda_metrics_server_wi" {
service_account_id = google_service_account.keda_sa.name
role = "roles/iam.workloadIdentityUser"
# 末尾のターゲットを [keda/keda-metrics-apiserver] に変更
member = "serviceAccount:${var.project_id}.svc.id.goog[keda/keda-metrics-apiserver]"
depends_on = [
google_container_cluster.primary
]
}
# ② アプリ用の紐付け(defaultネームスペースの「ocr-app-sa」という名前のK8sアカウントと結合)
resource "google_service_account_iam_member" "app_wi" {
service_account_id = google_service_account.app_sa.name
role = "roles/iam.workloadIdentityUser"
member = "serviceAccount:${var.project_id}.svc.id.goog[default/ocr-app-sa]"
# 🌟【追記】GKEクラスターとノードプールが完全に出来上がってから作らせる
depends_on = [
google_container_cluster.primary
]
}
# ===============================
# 7. Kubernetes 内部設定 (Namespace & 🌟ServiceAccountの作成)
# ===============================
resource "kubernetes_namespace" "keda" {
metadata {
name = "keda"
}
depends_on = [google_container_cluster.primary]
}
# defaultネームスペース側に、アプリが身分を証明するための会員証(ServiceAccount)を作成
resource "kubernetes_service_account" "app_k8s_sa" {
metadata {
name = "ocr-app-sa"
namespace = "default"
annotations = {
# この注釈(アノテーション)を入れることで、GCPのSAとガチッとリンクします
"iam.gke.io/gcp-service-account" = google_service_account.app_sa.email
}
}
depends_on = [google_container_cluster.primary]
}
# ===============================
# 8. クラウドロギングから BigQuery へのログ転送(シンク設定)
# ===============================
# ログを蓄積するための BigQuery データセットを作成
resource "google_bigquery_dataset" "logging_dataset" {
dataset_id = var.dataset_id
friendly_name = "GKE OCR Pipeline Logs"
description = "Go と Python から出力されるカスタムログの蓄積用データセット"
location = var.region
delete_contents_on_destroy = true # 試験環境用:terraform destroy時にテーブルも削除
}
# ログルーターシンクの定義(GKEの特定のログのみを抽出してBigQueryに送る)
resource "google_logging_project_sink" "bq_log_sink" {
name = "ocr-pipeline-log-sink"
description = "GoとPythonコンテナの標準出力をBigQueryへ転送"
destination = "bigquery.googleapis.com/projects/${var.project_id}/datasets/${var.dataset_id}"
# フィルタ:GKE上の指定したコンテナ、またはログ内に特定のキー(RequestID)が含まれるものを対象にする
# ※GoやPython側で `jsonPayload.RequestID` が出力されるため、これをフックします
filter = <<EOT
resource.type="k8s_container"
jsonPayload.RequestID:*
EOT
# 一意のライター(サービスアカウント)を自動生成させて権限付与に利用する
unique_writer_identity = true
}
# ログルーター(シンク)が BigQuery にデータを書き込めるように権限(IAM)を付与
resource "google_project_iam_member" "sink_bq_writer" {
project = var.project_id
role = "roles/bigquery.dataEditor"
member = google_logging_project_sink.bq_log_sink.writer_identity
}
6.Helmによるワークロード管理
本システムでは、KubernetesのパッケージマネージャーであるHelm(ヘルム)を採用し、環境依存のパラメータをすべて values.yaml に切り出して一元管理しています。
6-1. values.yaml
本システムにおける values.yaml は、「環境依存のパラメータ」および「運用時に頻繁に変更するパラメータ」のみを厳選して定義する最小限の構成にしています。
./helm-chart/values.yaml
# GCPの基本設定
gcp:
projectId: ""
region: "asia-northeast1"
# Go API (Publisher) の設定
goApi:
replicaCount: 1
# Python Worker (Subscriber) の設定
pythonWorker:
replicaCount: 1
6-2. go-api.yaml
いくつかのポイントをピックアップして説明します。
1. Goアプリへの「動的な環境変数注入」とイメージタグの自動生成
printf 関数を利用して、Artifact Registryのコンテナイメージのフルパスを values.yaml の値から動的に組み立てています。
image: {{ printf "%s-docker.pkg.dev/%s/ocr-pipeline-repo/go-api:v1" .Values.gcp.region .Values.gcp.projectId }}
これを行うことで、GCPのプロジェクト名やリージョンが変わっても、マニフェスト自体を一切修正することなく追従させることができます。
2. Kubernetes Downward APIによるPod名の取得
- name: POD_NAME
valueFrom:
fieldRef:
fieldPath: metadata.name
Goアプリ側が自分自身の「Pod名(動的に割り振られる一意の識別子)」を知るために、Kubernetesの Downward API を利用しています。
これを環境変数 POD_NAME としてGoアプリに渡すことで、BigQueryに「どのPodがリクエストを処理したか」のログ(トレーサビリティ)を正確に記録できるようになり、負荷分散の状況やエラー発生時の調査が劇的にスムーズになります。
3. 外部負荷テスト(k6)を見据えたロードバランサー化
type: LoadBalancer
APIサーバーを外部(インターネット)に公開し、負荷テストツール(k6)などからダイレクトに大量のリクエストを流し込めるよう、Serviceの型を ClusterIP から LoadBalancer に変更しています。これにより、クラウド側(GCP)のロードバランサーが自動でプロビジョニングされ、パブリックIP経由での高負荷検証が容易になります。
./helm-chart/templates/go-api.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: go-api
namespace: default
spec:
replicas: {{ .Values.goApi.replicaCount }}
selector:
matchLabels:
app: go-api
template:
metadata:
labels:
app: go-api
spec:
serviceAccountName: ocr-app-sa
containers:
- name: go-api
image: {{ printf "%s-docker.pkg.dev/%s/ocr-pipeline-repo/go-api:v1" .Values.gcp.region .Values.gcp.projectId }}
imagePullPolicy: Always
ports:
- containerPort: 8080
env:
- name: GCP_PROJECT_ID
value: {{ .Values.gcp.projectId }}
# Goアプリ側で実際のPod名を取得してBigQueryに送るための設定
- name: POD_NAME
valueFrom:
fieldRef:
fieldPath: metadata.name
---
apiVersion: v1
kind: Service
metadata:
name: go-api-service
namespace: default
spec:
# k6(外部)からのトラフィックを直接受けるためにロードバランサーへ変更
type: LoadBalancer
ports:
- port: 8080
targetPort: 8080
protocol: TCP
selector:
app: go-api
6-3. keda-scaler.yaml
keda-scaler.yamlは、GCP Workload Identityを用いた安全な認証と、pollingIntervalを2秒、cooldownPeriodを10秒に設定した「超高感度」なスケーリングを特徴としています。さらにstabilizationWindowSeconds: 0とsubscriptionSizeMetricsMode: "pubsub"を組み合わせ、デモや負荷テスト用に即時のスケールイン・アウトを実現する構成です。
./helm-chart/templates/keda-scaler.yaml
apiVersion: keda.sh/v1alpha1
kind: TriggerAuthentication
metadata:
name: keda-gcp-auth
namespace: default
spec:
# Secretではなく、Pod自身の持つWorkload Identity(gcp)を使うように指定します
podIdentity:
provider: gcp
---
apiVersion: keda.sh/v1alpha1
kind: ScaledObject
metadata:
name: python-worker-scaler
namespace: default
spec:
scaleTargetRef:
apiVersion: apps/v1
kind: Deployment
name: python-worker
# キューが空なら1台、初動を早くしたい場合は1に変更
minReplicaCount: 1
maxReplicaCount: 5
# 処理が終わってからPodを0台に「減らす」までの待ち時間を10秒に縮める
# (デフォルトの30秒だと、次の検証まで30秒待たないと0台に戻りません)
cooldownPeriod: 10
# Pub/Subの残件を見に行く間隔を最速の「1秒」にする
# (メッセージが届いたことを1秒以内に超高速検知してPodを起動させます)
pollingInterval: 2
# HPAの縮小バッファを無効化(0秒)にする
advanced:
horizontalPodAutoscalerConfig:
behavior:
scaleDown:
stabilizationWindowSeconds: 0
triggers:
- type: gcp-pubsub
metadata:
value: "2"
subscriptionName: "projects/{{ .Values.gcp.projectId }}/subscriptions/ocr-task-sub"
# Monitoring経由ではなく、直接の生件数を見る
subscriptionSizeMetricsMode: "pubsub"
# 10秒間の平均を見て超高感度に検知
subscriptionSizeSamplePeriod: "10"
authenticationRef:
name: keda-gcp-auth
6-4. python-worker.yaml
いくつかのポイントをピックアップし説明します。
1. GKE Autopilot「スポットインスタンス」によるコスト削減
nodeSelector:
://google.com: "true"
本システムでは、GKEスポットインスタンスを採用し通常のインフラコストを削減しています。
terminationGracePeriodSeconds: 25
スポットインスタンスに伴う設定として、GCP側の都合で突然強制終了(回収)される 25秒前 に停止シグナル(SIGTERM)が送信される設定とし、Python側でしデータ損失を防ぐ仕組みとなっています。
3. Autopilotの制約をクリアする絶妙なリソース設計(OOM防止)
resources:
requests:
cpu: "500m" # 0.5コア
memory: "2048Mi" # 2.0GB
limits:
cpu: "1000m" # 1コア
memory: "2048Mi" # 2.0GB(requestsと同値)
-
メモリ不足(OOM)の防止: 実用レベルの推論速度と安定性を確保するため、メモリを
2048Mi (2GB)まで引き上げています。 -
Autopilot特有のルール: GKE Autopilotの仕様上、メモリの
limitsはrequestsと「完全に同じ値」にする必要があります。 -
CPUのバースト対応:
limitsのみ1000m (1コア)に引き上げ、柔軟にバースト対応できるようにしています。
./helm-chart/templates/python-worker.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
name: python-worker
namespace: default
spec:
replicas: {{ .Values.pythonWorker.replicaCount }}
selector:
matchLabels:
app: python-worker
template:
metadata:
labels:
app: python-worker
spec:
spec:
# Autopilotにスポット(格安枠)で動かすよう命令
nodeSelector:
#://google.com: "true"
cloud.google.com/gke-spot: "true"
# スポット回収シグナルが来てから退避するまでの猶予(worker.pyの設計と同期)
terminationGracePeriodSeconds: 25
serviceAccountName: ocr-app-sa
containers:
- name: python-worker
image: {{ printf "%s-docker.pkg.dev/%s/ocr-pipeline-repo/python-worker:v1" .Values.gcp.region .Values.gcp.projectId }}
imagePullPolicy: Always
resources:
# Autopilotの最小・最適ルールを適用
resources:
requests:
cpu: "500m" # 0.5コアあればEasyOCRのCPU推論は十分実用レベルです
memory: "2048Mi" # 本物のモデルを展開するために2.0GBへ引き上げ(OOM防止)
limits:
cpu: "1000m" # 突発的な重い画像用。Autopilotはlimitsに応じた追加課金はありません
memory: "2048Mi" # Autopilotの仕様上、メモリのlimitsはrequestsと同値にする必要が
env:
- name: GCP_PROJECT_ID
value: {{ .Values.gcp.projectId }}
- name: POD_NAME
valueFrom:
fieldRef:
fieldPath: metadata.name
- name: PYTHONUNBUFFERED
value: "1"
7.デプロイ手順
本システムのコードは GitHub に公開しています。以下の手順で、インフラの構築からアプリケーションのデプロイまでを再現可能です。
7-1. インフラの構築 (Terraform)
まず、GCP上にGKEクラスタやArtifact Registryを構築します。
セキュリティ保護のため、プロジェクトID等を含む terraform.tfvars は GitHub に含めていません。以下の手順に従って、ご自身の環境に合わせたファイルを作成してください
1. terraform.tfvars の作成
terraformディレクトリ直下にterraform.tfvarsという名前でファイルを作成し、以下の内容を記述します。
project_id = "あなたのGCPプロジェクトID"
region = "asia-northeast1"
dataset_id = "log_dataset" # 任意の名前に変更可能
2. インフラのデプロイ
準備ができたら、以下のコマンドでリソースを作成します。
cd terraform
# 初期化
terraform init
# 実行(内容を確認して 'yes' を入力)
terraform apply
7-2. イメージのビルド・リネーム・プッシュ
GoとPythonの各イメージをビルドし、Artifact Registryへプッシュします。
以下のコマンドの <PROJECT_ID> 部分は、ご自身のGCPプロジェクトIDに置き換えて実行してください。(ここでは Windows環境での実行説明となっています)
# goディレクトリに移動する
cd .\go\
# go-api のイメージをbuild、プッシュ(アップロード)するスクリプト
# このファイル内に変数として「プロジェクトID」がありますのでご自身の物を設定して下さい
.\cmd_docker.ps1
# python ディレクトリに移動する
cd ..\python\
# python-api のイメージをbuild、プッシュ(アップロード)するスクリプト
# このファイル内に変数として「プロジェクトID」がありますのでご自身の物を設定して下さい
.\cmd_docker.ps1
7-3. アプリケーションのデプロイ (Helm)
Helmを使用して、K8sリソースを一括デプロイします。
# 接続先クラスタを切り替え、KEDAのインストール
# このファイル内に変数として「プロジェクトID」がありますのでご自身の物を設定して下さい
.\cmd_helm_1.ps1
# Helmチャートをインストール
# このファイル内に変数として「プロジェクトID」がありますのでご自身の物を設定して下さい
.\cmd_helm_2.ps1
# Podの起動確認(Running になるまで待ちます)
kubectl get pods
# Serviceの確認(外部IPが割り当てられているか確認)
kubectl get svc
# (参考)アプリケーションを削除する場合
# helm uninstall ocr-pipeline
7-4. k6 テスト前の準備
テストを実行する前に、接続先のIPアドレスの設定と、テスト用データの生成が必要です。
1.接続先IPアドレスの設定
GKE上のServiceに割り当てられた外部IPアドレスを、k6のスクリプトに反映させます。
以下のコマンドで go-api-service の EXTERNAL-IP を確認します。
kubectl get svc
./k6/script.js の 20行目にあるIPアドレスを、確認したIPアドレスに書き換えます。
// 修正前: const url = 'http://35.221.115.130:8080/predict';
const url = 'http://<あなたのEXTERNAL-IP>:8080/predict';
2. テスト用画像の生成
負荷試験に使用する画像(100枚)はリポジトリに含まれていないため、スクリプトで生成します。
ディレクトリ./k6に移動してから以下の作業を実施して下さい。
画像生成に必要なライブラリをインストールします。
pip install Pillow
生成スクリプトを実行します
python ./gen_image.py
実行後、./test_imagesディレクトリに100枚の画像と、正解ラベルとなるmapping.jsonが作成されれば準備完了です。
8.おわりに
本記事では、GoとPython、そしてGKEを組み合わせた 「スケーラブルなOCR基盤」 の
詳細な設定やコマンドについて解説しました。
💡 本記事は、3部作の「詳細説明編」です。
・【改善編】詳細説明編:「Go/Python/Terraform/Helm」の詳細解説。
・【改善編】総括編: Pub/Subを用Push型データパスの要点・効果を解説。
・【改善編】検証編: 負荷試験(k6)による大量リクエスト流入時の挙動を検証・考察。


