シリーズ構成
- Go の基本文法
- CLI ツールを作る
- HTTP サーバーと REST API
- データベース連携
- 実務的な周辺技術(本記事)
- ポートフォリオを作る
- 実践知識を深める
この回の目標は 「API を本番相当の構成で動かせるようにし、Go の並行処理を自信を持って書けるようになる」 ことです。
前半は運用に必要な道具(Docker・設定・ログ・認証・lint)、後半は Go の最大の武器である 並行処理 を体系的に扱います。並行処理はこの記事の中で最も長い章ですが、バックエンドエンジニアとして避けて通れない部分です。
目次
Part 1:運用のための道具
- Docker とマルチステージビルド
- Docker Compose で開発環境を作る
- 設定管理(環境変数)
- 構造化ログ(log/slog)
- 認証:パスワードハッシュと JWT
- Makefile・lint・CI
Part 2:並行処理
7. goroutine の基礎と落とし穴
8. sync パッケージ:Mutex・RWMutex・WaitGroup・Once
9. channel の設計
10. select
11. context:キャンセルとタイムアウト
12. 並行処理パターン集
13. errgroup
14. データ競合の検出とデバッグ
15. 演習問題
Part 1:運用のための道具
1. Docker とマルチステージビルド
Go は 静的リンクされた単一バイナリ を生成できるため、極めて小さいコンテナイメージが作れます。
Dockerfile
# ---------- ビルドステージ ----------
FROM golang:1.23-alpine AS builder
WORKDIR /app
# 依存だけ先にダウンロードしてレイヤーキャッシュを効かせる
COPY go.mod go.sum ./
RUN go mod download
COPY . .
# CGO_ENABLED=0 で完全静的リンク。-ldflags でサイズ削減
RUN CGO_ENABLED=0 GOOS=linux go build \
-ldflags="-s -w -X main.version=$(git describe --tags --always 2>/dev/null || echo dev)" \
-o /server ./cmd/server
# ---------- 実行ステージ ----------
FROM gcr.io/distroless/static-debian12:nonroot
COPY --from=builder /server /server
USER nonroot:nonroot
EXPOSE 8080
ENTRYPOINT ["/server"]
docker build -t todo-api .
docker images todo-api # 数 MB〜十数 MB
docker run -p 8080:8080 -e DATABASE_URL=... todo-api
なぜこう書くか
| 行 | 理由 |
|---|---|
COPY go.mod go.sum を先に | ソースが変わっても依存が同じならキャッシュが効く |
CGO_ENABLED=0 | C ライブラリに依存しない静的バイナリ。distroless/static で動く |
-s -w | シンボルテーブルとデバッグ情報を削除。サイズ 30% 減 |
distroless | シェルもパッケージマネージャもない。攻撃面が最小 |
nonroot | root で動かさない |
デバッグでシェルが必要なら alpine を使います。
FROM alpine:3.20
RUN apk add --no-cache ca-certificates tzdata
COPY --from=builder /server /server
ca-certificates は HTTPS 通信に、tzdata はタイムゾーン処理に必要です。distroless にはどちらも含まれています。
.dockerignore
.git
*.md
tmp/
*_test.go
.env
2. Docker Compose で開発環境を作る
compose.yaml
services:
db:
image: postgres:16-alpine
environment:
POSTGRES_USER: app
POSTGRES_PASSWORD: secret
POSTGRES_DB: todo
ports:
- "5432:5432"
volumes:
- pgdata:/var/lib/postgresql/data
healthcheck:
test: ["CMD-SHELL", "pg_isready -U app -d todo"]
interval: 3s
timeout: 3s
retries: 10
migrate:
image: migrate/migrate:v4.17.0
volumes:
- ./db/migrations:/migrations
command: ["-path", "/migrations", "-database", "postgres://app:secret@db:5432/todo?sslmode=disable", "up"]
depends_on:
db:
condition: service_healthy
api:
build: .
environment:
PORT: "8080"
DATABASE_URL: postgres://app:secret@db:5432/todo?sslmode=disable
JWT_SECRET: dev-secret-change-me
LOG_LEVEL: debug
ports:
- "8080:8080"
depends_on:
migrate:
condition: service_completed_successfully
volumes:
pgdata:
docker compose up --build # 起動
docker compose logs -f api # ログ
docker compose down -v # 停止してボリュームも削除
ホットリロード(開発時)
Go はコンパイルが速いですが、保存で自動再起動すると更に快適です。air が定番です。
go install github.com/air-verse/air@latest
air init
air
Compose で使う場合は api サービスのイメージを golang:1.23 にし、ソースをマウントして air を起動します。
3. 設定管理(環境変数)
12 Factor App の原則に従い、設定は環境変数から読みます。設定ファイルをイメージに焼き込まないことで、同じイメージを dev / staging / prod で使えます。
// internal/config/config.go
package config
import (
"errors"
"fmt"
"os"
"strconv"
"time"
)
type Config struct {
Port int
DatabaseURL string
JWTSecret string
JWTExpiry time.Duration
LogLevel string
AllowedOrigins []string
ShutdownTimeout time.Duration
}
func Load() (Config, error) {
var errs []error
cfg := Config{
Port: envInt("PORT", 8080, &errs),
DatabaseURL: envRequired("DATABASE_URL", &errs),
JWTSecret: envRequired("JWT_SECRET", &errs),
JWTExpiry: envDuration("JWT_EXPIRY", 24*time.Hour, &errs),
LogLevel: envString("LOG_LEVEL", "info"),
AllowedOrigins: envList("ALLOWED_ORIGINS", []string{"http://localhost:3000"}),
ShutdownTimeout: envDuration("SHUTDOWN_TIMEOUT", 10*time.Second, &errs),
}
if len(cfg.JWTSecret) < 32 && cfg.JWTSecret != "" {
errs = append(errs, errors.New("JWT_SECRET must be at least 32 characters"))
}
if len(errs) > 0 {
return Config{}, errors.Join(errs...)
}
return cfg, nil
}
func envString(key, def string) string {
if v := os.Getenv(key); v != "" {
return v
}
return def
}
func envRequired(key string, errs *[]error) string {
v := os.Getenv(key)
if v == "" {
*errs = append(*errs, fmt.Errorf("%s is required", key))
}
return v
}
func envInt(key string, def int, errs *[]error) int {
v := os.Getenv(key)
if v == "" {
return def
}
n, err := strconv.Atoi(v)
if err != nil {
*errs = append(*errs, fmt.Errorf("%s: %w", key, err))
return def
}
return n
}
func envDuration(key string, def time.Duration, errs *[]error) time.Duration {
v := os.Getenv(key)
if v == "" {
return def
}
d, err := time.ParseDuration(v) // "30s", "5m", "1h"
if err != nil {
*errs = append(*errs, fmt.Errorf("%s: %w", key, err))
return def
}
return d
}
func envList(key string, def []string) []string {
v := os.Getenv(key)
if v == "" {
return def
}
return strings.Split(v, ",")
}
エラーを集めて一度に返す ことで、起動失敗時に「足りない設定を全部」把握できます。
ローカル開発では .env
go get github.com/joho/godotenv
import "github.com/joho/godotenv"
func main() {
_ = godotenv.Load() // .env があれば読む。なければ無視
cfg, err := config.Load()
// ...
}
.env は .gitignore に入れ、.env.example をコミットします。
# .env.example
PORT=8080
DATABASE_URL=postgres://app:secret@localhost:5432/todo?sslmode=disable
JWT_SECRET=change-me-to-a-random-32-plus-character-string
LOG_LEVEL=debug
ライブラリを使う場合
caarlos0/env や kelseyhightower/envconfig を使えば構造体タグで宣言できます。
type Config struct {
Port int `env:"PORT" envDefault:"8080"`
DatabaseURL string `env:"DATABASE_URL,required"`
}
cfg := Config{}
err := env.Parse(&cfg)
4. 構造化ログ(log/slog)
Go 1.21 で標準入りした log/slog を使います。JSON 形式 で出せば、Datadog・CloudWatch・Loki などのログ基盤でフィルタ・集計できます。
セットアップ
// internal/logging/logging.go
package logging
import (
"log/slog"
"os"
"strings"
)
func Setup(level string) *slog.Logger {
var lvl slog.Level
switch strings.ToLower(level) {
case "debug":
lvl = slog.LevelDebug
case "warn":
lvl = slog.LevelWarn
case "error":
lvl = slog.LevelError
default:
lvl = slog.LevelInfo
}
handler := slog.NewJSONHandler(os.Stdout, &slog.HandlerOptions{
Level: lvl,
AddSource: lvl == slog.LevelDebug, // debug 時はファイル名と行番号
})
logger := slog.New(handler)
slog.SetDefault(logger) // slog.Info() などで使えるように
return logger
}
使い方
slog.Info("server started", "port", 8080, "version", version)
slog.Warn("slow query", "duration_ms", 1523, "query", "SELECT ...")
slog.Error("failed to connect", "err", err)
// 属性をまとめる
slog.Info("user created",
slog.Int64("user_id", u.ID),
slog.String("email", u.Email),
slog.Group("request", "method", "POST", "path", "/users"),
)
{"time":"2025-01-15T10:30:00.123Z","level":"INFO","msg":"user created","user_id":42,"email":"a@example.com","request":{"method":"POST","path":"/users"}}
リクエストごとのロガー
リクエスト ID を すべてのログに自動で載せる には、context にロガーを入れます。
type loggerKey struct{}
func WithLogger(ctx context.Context, l *slog.Logger) context.Context {
return context.WithValue(ctx, loggerKey{}, l)
}
func FromContext(ctx context.Context) *slog.Logger {
if l, ok := ctx.Value(loggerKey{}).(*slog.Logger); ok {
return l
}
return slog.Default()
}
// ミドルウェア
func RequestLogger(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
reqID := middleware.GetReqID(r.Context())
l := slog.Default().With("request_id", reqID, "method", r.Method, "path", r.URL.Path)
ctx := WithLogger(r.Context(), l)
start := time.Now()
rec := &statusRecorder{ResponseWriter: w}
next.ServeHTTP(rec, r.WithContext(ctx))
l.Info("request completed", "status", rec.status, "duration_ms", time.Since(start).Milliseconds())
})
}
// どの層からでも
func (s *TodoService) Create(ctx context.Context, ...) {
log := logging.FromContext(ctx)
log.Debug("creating todo", "title", title)
}
ログの指針
- 何を出すか:リクエスト開始/完了、エラー、外部呼び出し、状態変化(ユーザー作成など)
- 何を出さないか:パスワード、トークン、クレジットカード、個人情報の生データ
- レベル:
Debugは開発用、Infoは通常の出来事、Warnは注意(リトライ成功など)、Errorは対応が必要 - 1 行 1 イベント:複数行のメッセージにしない
fmt.Printlnでログを出さない
機密情報をマスクする
type User struct {
Email string
PasswordHash string
}
// LogValue を実装するとログ出力時に自動で呼ばれる
func (u User) LogValue() slog.Value {
return slog.GroupValue(
slog.String("email", u.Email),
slog.String("password_hash", "[REDACTED]"),
)
}
5. 認証:パスワードハッシュと JWT
パスワードハッシュ
bcrypt または argon2id を使います。SHA-256 は高速すぎて総当たりに弱いので不可。
go get golang.org/x/crypto/bcrypt
// internal/auth/password.go
package auth
import "golang.org/x/crypto/bcrypt"
const bcryptCost = 12 // 10〜14。数値が 1 増えると 2 倍遅い
func HashPassword(password string) (string, error) {
hash, err := bcrypt.GenerateFromPassword([]byte(password), bcryptCost)
if err != nil {
return "", err
}
return string(hash), nil
}
func VerifyPassword(hash, password string) bool {
return bcrypt.CompareHashAndPassword([]byte(hash), []byte(password)) == nil
}
bcrypt は 72 バイトまでしか見ません。長いパスワードを許容するなら先に SHA-256 でハッシュしてから bcrypt に渡すか、argon2id を使います。
JWT(JSON Web Token)
go get github.com/golang-jwt/jwt/v5
// internal/auth/token.go
package auth
import (
"errors"
"fmt"
"time"
"github.com/golang-jwt/jwt/v5"
)
var ErrInvalidToken = errors.New("invalid token")
type Claims struct {
UserID int64 `json:"uid"`
jwt.RegisteredClaims
}
type TokenService struct {
secret []byte
expiry time.Duration
issuer string
}
func NewTokenService(secret string, expiry time.Duration) *TokenService {
return &TokenService{secret: []byte(secret), expiry: expiry, issuer: "todo-api"}
}
func (s *TokenService) Issue(userID int64) (string, error) {
now := time.Now()
claims := Claims{
UserID: userID,
RegisteredClaims: jwt.RegisteredClaims{
Issuer: s.issuer,
Subject: fmt.Sprint(userID),
IssuedAt: jwt.NewNumericDate(now),
ExpiresAt: jwt.NewNumericDate(now.Add(s.expiry)),
},
}
return jwt.NewWithClaims(jwt.SigningMethodHS256, claims).SignedString(s.secret)
}
func (s *TokenService) Verify(tokenString string) (*Claims, error) {
claims := &Claims{}
_, err := jwt.ParseWithClaims(tokenString, claims,
func(t *jwt.Token) (any, error) {
// アルゴリズム混同攻撃を防ぐ
if _, ok := t.Method.(*jwt.SigningMethodHMAC); !ok {
return nil, fmt.Errorf("unexpected signing method: %v", t.Header["alg"])
}
return s.secret, nil
},
jwt.WithIssuer(s.issuer),
jwt.WithExpirationRequired(),
)
if err != nil {
return nil, fmt.Errorf("%w: %v", ErrInvalidToken, err)
}
return claims, nil
}
認証ミドルウェア
// internal/auth/middleware.go
package auth
type userIDKey struct{}
func UserIDFrom(ctx context.Context) (int64, bool) {
id, ok := ctx.Value(userIDKey{}).(int64)
return id, ok
}
func Middleware(tokens *TokenService) func(http.Handler) http.Handler {
return func(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
authz := r.Header.Get("Authorization")
const prefix = "Bearer "
if !strings.HasPrefix(authz, prefix) {
w.Header().Set("WWW-Authenticate", `Bearer realm="api"`)
writeError(w, http.StatusUnauthorized, "missing bearer token")
return
}
claims, err := tokens.Verify(strings.TrimPrefix(authz, prefix))
if err != nil {
writeError(w, http.StatusUnauthorized, "invalid token")
return
}
ctx := context.WithValue(r.Context(), userIDKey{}, claims.UserID)
next.ServeHTTP(w, r.WithContext(ctx))
})
}
}
サインアップとログインのハンドラ
// internal/handler/auth.go
type AuthHandler struct {
users UserRepo
tokens *auth.TokenService
}
type credentials struct {
Email string `json:"email"`
Password string `json:"password"`
}
func (h *AuthHandler) Signup(w http.ResponseWriter, r *http.Request) {
var in credentials
if err := decodeJSON(r, &in); err != nil {
writeError(w, http.StatusBadRequest, err.Error())
return
}
if !isValidEmail(in.Email) || len(in.Password) < 8 {
writeError(w, http.StatusUnprocessableEntity, "invalid email or password too short")
return
}
hash, err := auth.HashPassword(in.Password)
if err != nil {
writeError(w, http.StatusInternalServerError, "internal error")
return
}
u, err := h.users.Create(r.Context(), strings.ToLower(in.Email), hash)
if errors.Is(err, repository.ErrEmailTaken) {
writeError(w, http.StatusConflict, "email already registered")
return
}
if err != nil {
logging.FromContext(r.Context()).Error("create user", "err", err)
writeError(w, http.StatusInternalServerError, "internal error")
return
}
token, _ := h.tokens.Issue(u.ID)
writeJSON(w, http.StatusCreated, map[string]any{"token": token, "user_id": u.ID})
}
func (h *AuthHandler) Login(w http.ResponseWriter, r *http.Request) {
var in credentials
if err := decodeJSON(r, &in); err != nil {
writeError(w, http.StatusBadRequest, err.Error())
return
}
u, err := h.users.GetByEmail(r.Context(), strings.ToLower(in.Email))
// ユーザーが存在しない場合も同じメッセージ(列挙攻撃対策)
if err != nil || !auth.VerifyPassword(u.PasswordHash, in.Password) {
writeError(w, http.StatusUnauthorized, "invalid email or password")
return
}
token, err := h.tokens.Issue(u.ID)
if err != nil {
writeError(w, http.StatusInternalServerError, "internal error")
return
}
writeJSON(w, http.StatusOK, map[string]string{"token": token})
}
curl -X POST localhost:8080/signup -d '{"email":"a@example.com","password":"password123"}'
# {"token":"eyJ...","user_id":1}
curl localhost:8080/api/v1/todos -H "Authorization: Bearer eyJ..."
JWT の注意点
- 秘密鍵は 32 バイト以上のランダム値。
openssl rand -base64 32で生成 - 有効期限は短く(15 分〜24 時間)。長期セッションはリフレッシュトークンで
- JWT は 失効できない。ログアウトを厳密にしたいならサーバー側セッション(Redis など)を検討
- ペイロードは Base64 で誰でも読める。機密情報を入れない
- 本番では RS256 / ES256(公開鍵方式)も選択肢。複数サービスで検証する場合に有利
6. Makefile・lint・CI
Makefile
.PHONY: run build test lint fmt migrate-up migrate-down sqlc docker
BIN := bin/server
DATABASE_URL ?= postgres://app:secret@localhost:5432/todo?sslmode=disable
run:
go run ./cmd/server
build:
CGO_ENABLED=0 go build -ldflags="-s -w" -o $(BIN) ./cmd/server
test:
go test -race -cover -count=1 ./...
test-integration:
TEST_DATABASE_URL=$(DATABASE_URL) go test -race -tags=integration ./...
lint:
golangci-lint run ./...
fmt:
gofmt -s -w .
goimports -w .
migrate-up:
migrate -database "$(DATABASE_URL)" -path db/migrations up
migrate-down:
migrate -database "$(DATABASE_URL)" -path db/migrations down 1
migrate-new:
@read -p "name: " name; migrate create -ext sql -dir db/migrations -seq $$name
sqlc:
sqlc generate
docker:
docker compose up --build
golangci-lint
brew install golangci-lint
.golangci.yml
run:
timeout: 5m
linters:
enable:
- errcheck # エラーの無視を検出
- govet # 標準の静的解析
- staticcheck # 幅広いバグ検出
- unused
- gosimple
- ineffassign # 無効な代入
- gocritic
- revive # スタイル
- misspell
- bodyclose # resp.Body.Close() 忘れ
- noctx # context なしの HTTP リクエスト
- sqlclosecheck # rows.Close() 忘れ
- errorlint # errors.Is/As の使い忘れ
- gosec # セキュリティ
linters-settings:
revive:
rules:
- name: exported
disabled: true # 全公開関数にコメントを強制しない
golangci-lint run ./...
GitHub Actions
.github/workflows/ci.yml
name: CI
on:
push:
branches: [main]
pull_request:
jobs:
test:
runs-on: ubuntu-latest
services:
postgres:
image: postgres:16-alpine
env:
POSTGRES_USER: app
POSTGRES_PASSWORD: secret
POSTGRES_DB: todo_test
ports: ["5432:5432"]
options: >-
--health-cmd "pg_isready -U app"
--health-interval 5s
--health-timeout 5s
--health-retries 10
steps:
- uses: actions/checkout@v4
- uses: actions/setup-go@v5
with:
go-version-file: go.mod
cache: true
- name: Lint
uses: golangci/golangci-lint-action@v6
with:
version: latest
- name: Check sqlc is up to date
run: |
go install github.com/sqlc-dev/sqlc/cmd/sqlc@latest
sqlc generate
git diff --exit-code
- name: Migrate
run: |
go install -tags 'postgres' github.com/golang-migrate/migrate/v4/cmd/migrate@latest
migrate -database "postgres://app:secret@localhost:5432/todo_test?sslmode=disable" -path db/migrations up
- name: Test
run: go test -race -cover -count=1 ./...
env:
TEST_DATABASE_URL: postgres://app:secret@localhost:5432/todo_test?sslmode=disable
- name: Build
run: CGO_ENABLED=0 go build -o /dev/null ./cmd/server
Part 2:並行処理
7. goroutine の基礎と落とし穴
goroutine とは
Go ランタイムが管理する軽量スレッドです。初期スタックは 2KB 程度で、数万〜数十万個を同時に動かせます。OS スレッドとの対応は Go スケジューラが自動で行います(M:N モデル)。
go doSomething() // 関数呼び出しの前に go をつけるだけ
go func() {
// 無名関数
}()
落とし穴 1:main が先に終わる
func main() {
go fmt.Println("hello")
// main が終わると全 goroutine が強制終了。hello は出ないかもしれない
}
goroutine の完了を待つ仕組み(WaitGroup、channel)が必ず必要です。
落とし穴 2:goroutine リーク
終了条件のない goroutine は永遠にメモリを占有します。
// ✗ 受信側がいないと送信でブロックし、goroutine が永遠に残る
func leak() {
ch := make(chan int)
go func() {
ch <- 1 // 誰も受け取らない
}()
return
}
すべての goroutine に「いつ終わるか」を設計する のが原則です。多くの場合、context のキャンセルか channel のクローズが終了条件になります。
落とし穴 3:エラーが消える
go func() {
if err := process(); err != nil {
// ここで return しても誰も知らない
}
}()
goroutine のエラーは channel か errgroup で親に返します(後述)。
落とし穴 4:panic でプロセス全体が落ちる
goroutine 内の panic は recover されないとプロセスを終了させます。HTTP ハンドラは net/http が recover してくれますが、自分で起動した goroutine は自分で守ります。
go func() {
defer func() {
if r := recover(); r != nil {
slog.Error("goroutine panic", "err", r, "stack", string(debug.Stack()))
}
}()
process()
}()
8. sync パッケージ
Mutex:共有状態を守る
複数の goroutine が同じ変数を読み書きすると データ競合 になります。結果が不定になり、最悪クラッシュします。
type Counter struct {
mu sync.Mutex
n map[string]int
}
func (c *Counter) Inc(key string) {
c.mu.Lock()
defer c.mu.Unlock()
c.n[key]++
}
func (c *Counter) Get(key string) int {
c.mu.Lock()
defer c.mu.Unlock()
return c.n[key]
}
ルール
- Mutex は 守る対象の直上に置く(構造体の中で、フィールドの直前)
Lockしたらdefer Unlock- Lock 中に 外部呼び出し(DB、HTTP、他の Lock)をしない(デッドロックの原因)
- Mutex を含む構造体は コピーしない(ポインタで渡す)。
go vetが警告する
RWMutex:読み取りが多い場合
type Cache struct {
mu sync.RWMutex
data map[string]string
}
func (c *Cache) Get(k string) (string, bool) {
c.mu.RLock() // 複数の読み取りは同時に進める
defer c.mu.RUnlock()
v, ok := c.data[k]
return v, ok
}
func (c *Cache) Set(k, v string) {
c.mu.Lock() // 書き込みは排他
defer c.mu.Unlock()
c.data[k] = v
}
読み取りが書き込みの 10 倍以上多い場合に効果があります。そうでなければ普通の Mutex の方が単純で速いこともあります。
WaitGroup:完了を待つ
func fetchAll(urls []string) []Result {
var wg sync.WaitGroup
results := make([]Result, len(urls)) // インデックスで書き込むので競合しない
for i, u := range urls {
wg.Add(1)
go func() {
defer wg.Done()
results[i] = fetch(u)
}()
}
wg.Wait()
return results
}
Addは goroutine を起動する前 に呼ぶ(中で呼ぶと Wait が先に通り抜ける可能性がある)- Go 1.22 以降はループ変数
i,uをそのまま閉じ込めて安全
Once:一度だけ初期化
var (
instance *Config
once sync.Once
)
func GetConfig() *Config {
once.Do(func() {
instance = loadConfig()
})
return instance
}
Go 1.21 からは sync.OnceValue でより簡潔に書けます。
var GetConfig = sync.OnceValue(func() *Config {
return loadConfig()
})
atomic:単純な数値ならロック不要
var requests atomic.Int64
requests.Add(1)
n := requests.Load()
カウンタやフラグのような 単一の値 には sync/atomic の方が軽量です。複数の値をまとめて更新するなら Mutex。
sync.Map
読み取りが圧倒的に多く、キーが安定しているキャッシュ向け。通常は RWMutex + map で十分です。
9. channel の設計
channel は goroutine 間で値を安全に受け渡す 型付きのパイプです。「メモリを共有して通信するのではなく、通信によってメモリを共有する」が Go の哲学です。
基本操作
ch := make(chan int) // unbuffered
ch := make(chan int, 10) // buffered(容量 10)
ch <- 42 // 送信
v := <-ch // 受信
v, ok := <-ch // ok=false ならクローズ済みかつ空
close(ch) // クローズ
unbuffered vs buffered
| unbuffered | buffered | |
|---|---|---|
| 送信 | 受信側が受け取るまで ブロック | バッファに空きがあれば即時 |
| 用途 | 同期(ハンドオフ)、完了通知 | 生産者と消費者の速度差を吸収 |
| 注意 | 受信側がいないとデッドロック | 容量の根拠を説明できること |
バッファサイズは「なんとなく 100」にしないでください。 多くの場合、unbuffered か 1 か「ワーカー数と同じ」が正解です。
方向を型で制限する
関数の引数では 送信専用・受信専用 にすると、誤用をコンパイル時に防げます。
func producer(out chan<- int) { // 送信のみ
for i := 0; i < 5; i++ {
out <- i
}
close(out)
}
func consumer(in <-chan int) { // 受信のみ
for v := range in { // クローズされるまでループ
fmt.Println(v)
}
}
func main() {
ch := make(chan int)
go producer(ch)
consumer(ch)
}
close のルール
- 送信側だけがクローズする。受信側がクローズすると送信側が panic する
- クローズ済み channel への送信は panic
- クローズ済み channel からの受信は ゼロ値と ok=false を返す(ブロックしない)
- 二重クローズは panic
- クローズしなくてもリークはしない(GC される)。
rangeで受ける側にループ終了を伝えるためにクローズする
nil channel
nil channel への送受信は 永遠にブロック します。select で特定の case を無効化するテクニックに使います。
完了通知には struct{}
done := make(chan struct{})
go func() {
work()
close(done) // 値を送らずクローズで通知。複数の受信者に一斉に伝わる
}()
<-done
10. select
複数の channel 操作を 同時に待つ 構文です。
select {
case v := <-ch1:
fmt.Println("from ch1:", v)
case ch2 <- 42:
fmt.Println("sent to ch2")
case <-time.After(1 * time.Second):
fmt.Println("timeout")
case <-ctx.Done():
return ctx.Err()
}
- 複数の case が同時に準備できたら ランダムに 1 つ 選ばれる
defaultがあると、どの case も準備できていないとき即座にdefaultが実行される(ノンブロッキング)
ノンブロッキング送信
select {
case ch <- v:
// 送れた
default:
// バッファが満杯。捨てる or ログ
}
ループ内の select
func worker(ctx context.Context, jobs <-chan Job) {
for {
select {
case <-ctx.Done():
return
case j, ok := <-jobs:
if !ok {
return // クローズされた
}
process(j)
}
}
}
time.After のリーク
ループ内で time.After を使うと、タイマーが GC されるまでメモリを保持します。ループでは time.NewTimer を使い、Stop / Reset で再利用します(Go 1.23 以降は改善されましたが、癖として覚えておきます)。
timer := time.NewTimer(d)
defer timer.Stop()
for {
timer.Reset(d)
select {
case <-timer.C:
case v := <-ch:
}
}
11. context:キャンセルとタイムアウト
context.Context は「この処理はもう不要」「この時刻までに終わらせて」を 関数呼び出しの連鎖全体に伝える 仕組みです。
生成
ctx := context.Background() // ルート。main やテストで使う
ctx := context.TODO() // 後で適切な context に置き換える印
ctx, cancel := context.WithCancel(parent)
ctx, cancel := context.WithTimeout(parent, 5*time.Second)
ctx, cancel := context.WithDeadline(parent, time.Now().Add(5*time.Second))
defer cancel() // 必ず呼ぶ(リソース解放)
// Go 1.21+:キャンセル理由を付ける
ctx, cancel := context.WithCancelCause(parent)
cancel(errors.New("user logged out"))
context.Cause(ctx) // → その error
確認
select {
case <-ctx.Done():
return ctx.Err() // context.Canceled または context.DeadlineExceeded
default:
// 続行
}
// または単に
if err := ctx.Err(); err != nil {
return err
}
ツリー構造
子 context をキャンセルしても親は影響を受けませんが、親をキャンセルすると子孫すべてがキャンセル されます。
Background
└── HTTP リクエストの ctx(クライアント切断でキャンセル)
├── DB クエリ用 ctx(3 秒タイムアウト)
└── 外部 API 用 ctx(2 秒タイムアウト)
実践:長い処理をキャンセル可能にする
func processItems(ctx context.Context, items []Item) error {
for i, item := range items {
// 定期的にキャンセルを確認
if err := ctx.Err(); err != nil {
return fmt.Errorf("cancelled at item %d: %w", i, err)
}
if err := process(ctx, item); err != nil {
return err
}
}
return nil
}
実践:外部呼び出しに個別タイムアウト
func (s *Service) FetchProfile(ctx context.Context, id int64) (*Profile, error) {
// リクエスト全体が 30 秒でも、この API は 2 秒で諦める
ctx, cancel := context.WithTimeout(ctx, 2*time.Second)
defer cancel()
req, _ := http.NewRequestWithContext(ctx, "GET", url, nil)
resp, err := s.client.Do(req)
if errors.Is(err, context.DeadlineExceeded) {
return nil, ErrUpstreamTimeout
}
// ...
}
context のルール
- 第一引数に置き、名前は
ctx - 構造体のフィールドに 保存しない(リクエストごとに渡す)
nilを渡さない。迷ったらcontext.TODO()WithValueはリクエストスコープのメタデータ(リクエスト ID、認証ユーザー)のみ。関数の引数になるべき値を入れないcancelは必ず呼ぶ。defer cancel()が定石
12. 並行処理パターン集
パターン 1:ワーカープール
同時実行数を制限しつつ大量のジョブを処理します。
func WorkerPool[T, R any](ctx context.Context, jobs []T, workers int, fn func(context.Context, T) (R, error)) ([]R, error) {
type indexed struct {
i int
res R
err error
}
jobCh := make(chan int)
resCh := make(chan indexed)
var wg sync.WaitGroup
// ワーカー起動
for w := 0; w < workers; w++ {
wg.Add(1)
go func() {
defer wg.Done()
for i := range jobCh {
res, err := fn(ctx, jobs[i])
resCh <- indexed{i, res, err}
}
}()
}
// ジョブ投入(別 goroutine で。resCh の受信と並行させるため)
go func() {
defer close(jobCh)
for i := range jobs {
select {
case jobCh <- i:
case <-ctx.Done():
return
}
}
}()
// 全ワーカー終了後に resCh をクローズ
go func() {
wg.Wait()
close(resCh)
}()
results := make([]R, len(jobs))
var firstErr error
for r := range resCh {
if r.err != nil && firstErr == nil {
firstErr = r.err
}
results[r.i] = r.res
}
return results, firstErr
}
パターン 2:セマフォで同時実行数を制限
WaitGroup と buffered channel の組み合わせ。ワーカープールより簡潔です。
func fetchAllLimited(ctx context.Context, urls []string, limit int) []Result {
sem := make(chan struct{}, limit)
results := make([]Result, len(urls))
var wg sync.WaitGroup
for i, u := range urls {
wg.Add(1)
go func() {
defer wg.Done()
sem <- struct{}{} // 空きを待つ
defer func() { <-sem }() // 解放
results[i] = fetch(ctx, u)
}()
}
wg.Wait()
return results
}
golang.org/x/sync/semaphore に重み付きセマフォもあります。
パターン 3:パイプライン
処理を段階に分け、各段階を goroutine で繋ぎます。
func generate(ctx context.Context, nums ...int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for _, n := range nums {
select {
case out <- n:
case <-ctx.Done():
return
}
}
}()
return out
}
func square(ctx context.Context, in <-chan int) <-chan int {
out := make(chan int)
go func() {
defer close(out)
for n := range in {
select {
case out <- n * n:
case <-ctx.Done():
return
}
}
}()
return out
}
func main() {
ctx := context.Background()
for v := range square(ctx, generate(ctx, 1, 2, 3)) {
fmt.Println(v) // 1 4 9
}
}
パターン 4:Fan-out / Fan-in
1 つの入力を複数のワーカーで処理し(fan-out)、結果を 1 つに集約する(fan-in)。
func fanIn[T any](ctx context.Context, chans ...<-chan T) <-chan T {
out := make(chan T)
var wg sync.WaitGroup
for _, c := range chans {
wg.Add(1)
go func(c <-chan T) {
defer wg.Done()
for v := range c {
select {
case out <- v:
case <-ctx.Done():
return
}
}
}(c)
}
go func() {
wg.Wait()
close(out)
}()
return out
}
// 使い方
in := generate(ctx, nums...)
w1 := square(ctx, in)
w2 := square(ctx, in) // 同じ in から複数が読む = fan-out
w3 := square(ctx, in)
for v := range fanIn(ctx, w1, w2, w3) { // fan-in
fmt.Println(v)
}
パターン 5:タイムアウト付きの結果待ち
func withTimeout[T any](ctx context.Context, d time.Duration, fn func(context.Context) (T, error)) (T, error) {
ctx, cancel := context.WithTimeout(ctx, d)
defer cancel()
type result struct {
v T
err error
}
ch := make(chan result, 1) // buffered 1:goroutine がブロックせず終われる
go func() {
v, err := fn(ctx)
ch <- result{v, err}
}()
select {
case r := <-ch:
return r.v, r.err
case <-ctx.Done():
var zero T
return zero, ctx.Err()
}
}
ch を buffered 1 にするのが重要です。unbuffered だとタイムアウト後に goroutine が送信でブロックしてリークします。
パターン 6:定期実行
func runPeriodically(ctx context.Context, interval time.Duration, fn func(context.Context) error) {
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
if err := fn(ctx); err != nil {
slog.Error("periodic job failed", "err", err)
}
}
}
}
パターン 7:バックグラウンドワーカー付きサーバー
HTTP サーバーとバックグラウンド処理を同じプロセスで動かし、まとめて Graceful Shutdown します。
func main() {
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
var wg sync.WaitGroup
// バックグラウンドジョブ
wg.Add(1)
go func() {
defer wg.Done()
runPeriodically(ctx, time.Minute, cleanupExpiredTokens)
}()
// HTTP サーバー
srv := &http.Server{Addr: ":8080", Handler: router}
wg.Add(1)
go func() {
defer wg.Done()
if err := srv.ListenAndServe(); !errors.Is(err, http.ErrServerClosed) {
slog.Error("server", "err", err)
stop()
}
}()
<-ctx.Done()
slog.Info("shutting down")
shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
srv.Shutdown(shutdownCtx)
wg.Wait()
slog.Info("stopped")
}
13. errgroup
golang.org/x/sync/errgroup は WaitGroup + エラー収集 + context 連携をまとめたものです。実務で最も使う並行処理ライブラリ です。
go get golang.org/x/sync/errgroup
基本
func fetchDashboard(ctx context.Context, userID int64) (*Dashboard, error) {
g, ctx := errgroup.WithContext(ctx)
var (
user *User
todos []Todo
stats *Stats
)
g.Go(func() error {
var err error
user, err = userRepo.Get(ctx, userID)
return err
})
g.Go(func() error {
var err error
todos, err = todoRepo.List(ctx, userID)
return err
})
g.Go(func() error {
var err error
stats, err = statsRepo.Get(ctx, userID)
return err
})
if err := g.Wait(); err != nil {
return nil, err // 最初に発生したエラー。他の goroutine の ctx はキャンセル済み
}
return &Dashboard{User: user, Todos: todos, Stats: stats}, nil
}
3 つの DB クエリが 並行に 走り、1 つが失敗すると残りは ctx 経由でキャンセルされます。逐次実行なら合計 300ms かかる処理が、並行なら最遅の 100ms で終わります。
各 goroutine が 別々の変数に書く ので Mutex は不要です。同じスライスに append するなら Mutex が必要です。
同時実行数の制限
g, ctx := errgroup.WithContext(ctx)
g.SetLimit(10) // 同時に 10 個まで。超えると g.Go がブロック
for _, url := range urls {
g.Go(func() error {
return fetch(ctx, url)
})
}
return g.Wait()
TryGo
if !g.TryGo(fn) {
// 上限に達していて起動できなかった
}
errgroup vs 手書き
| errgroup | WaitGroup + channel | |
|---|---|---|
| エラー伝播 | 自動 | 自分で channel を設計 |
| 早期キャンセル | 自動(ctx) | 自分で実装 |
| 結果の収集 | 変数に直接書く | channel で集める |
| 適用場面 | 「全部成功か失敗か」 | 部分的な成功を許容、ストリーム処理 |
「複数の独立した処理を並行で走らせて全部の完了を待つ」なら まず errgroup を検討してください。
14. データ競合の検出とデバッグ
-race フラグ
Go には データ競合検出器 が組み込まれています。並行処理を書いたら必ず有効にしてテストします。
go test -race ./...
go run -race .
go build -race . # 本番では使わない(5〜10 倍遅い)
WARNING: DATA RACE
Write at 0x00c000012345 by goroutine 7:
main.increment()
/app/main.go:15 +0x44
Previous read at 0x00c000012345 by goroutine 6:
main.increment()
/app/main.go:15 +0x30
行番号が両方示されるので、原因箇所が特定できます。CI では 必ず -race を付ける ことを習慣にします。
競合を再現するテスト
func TestCounter_Concurrent(t *testing.T) {
c := NewCounter()
var wg sync.WaitGroup
for i := 0; i < 100; i++ {
wg.Add(1)
go func() {
defer wg.Done()
c.Inc("key")
}()
}
wg.Wait()
if got := c.Get("key"); got != 100 {
t.Errorf("got %d, want 100", got)
}
}
このテストは Mutex がなくても たまに通ります(それが競合の恐ろしさ)。-race を付けると確実に検出できます。
デッドロックの検出
全 goroutine がブロックすると Go ランタイムが自動で検出して終了します。
fatal error: all goroutines are asleep - deadlock!
ただし一部の goroutine だけが固まっている場合は検出されません。その場合は SIGQUIT(Ctrl+\)で全 goroutine のスタックをダンプするか、net/http/pprof の /debug/pprof/goroutine?debug=2 を見ます。
goroutine リークの検出
go.uber.org/goleak をテストに組み込むと、テスト終了時に残っている goroutine を報告します。
func TestMain(m *testing.M) {
goleak.VerifyTestMain(m)
}
pprof で goroutine 数を監視
import _ "net/http/pprof"
go http.ListenAndServe("localhost:6060", nil)
curl localhost:6060/debug/pprof/goroutine?debug=1 | head
goroutine 数が単調増加していたらリークしています。
15. 演習問題
Part 1:運用
- Docker イメージの最小化:
alpineベースとdistrolessベースの 2 つの Dockerfile を書き、イメージサイズを比較する。docker historyでレイヤーごとのサイズを確認する - 設定のテスト:
config.Load()に対し、必須項目が欠けたとき すべての不足項目 がエラーメッセージに含まれることをテストする(t.Setenvを使う) - ログのテスト:
slog.NewJSONHandlerの出力先をbytes.Bufferにし、リクエストログにrequest_idとstatusが含まれることを検証する - リフレッシュトークン:アクセストークン(15 分)とリフレッシュトークン(7 日、DB に保存)の 2 トークン方式を実装する。
POST /refreshでアクセストークンを再発行し、POST /logoutでリフレッシュトークンを失効させる - CI の整備:GitHub Actions で lint → sqlc 差分チェック → マイグレーション →
-raceテスト → ビルドが通るパイプラインを作り、バッジを README に貼る
Part 2:並行処理
- 安全なカウンタ:Mutex 版と
atomic版のCounterを実装し、100 goroutine × 1000 回インクリメントするテストを-race付きで通す。ベンチマークで性能を比較する - 並列ダウンローダ:URL のリストを受け取り、同時 5 接続で取得して結果をファイルに保存する CLI。
errgroup.SetLimitを使い、1 つ失敗したら残りをキャンセルする。進捗を標準エラーに表示する - タイムアウト付き集約 API:
GET /dashboardが 3 つのリポジトリを errgroup で並行に呼び、全体 2 秒でタイムアウトする。1 つが遅いとき 504 が返ることをhttptestでテストする(遅いリポジトリはモックでtime.Sleep) - ワーカープール with graceful stop:ジョブキュー(channel)とワーカー N 個を持つ
Pool型を作る。Stop(ctx)を呼ぶと新規受付を止め、処理中のジョブの完了をctxの期限まで待つ。リークしないことをgoleakで確認する - レートリミッタ:
golang.org/x/time/rateを使い、IP ごとに「1 秒 10 リクエスト、バースト 20」を許可するミドルウェアを作る。IP ごとのリミッタはsync.MapかMutex + mapで管理し、10 分アクセスがない IP のエントリを定期的に削除する - パイプライン処理:CSV を読む → 各行をパースする → バリデーションする → DB に書く、の 4 段パイプラインを channel で構築する。各段階のエラーを集約し、途中でキャンセルできるようにする
- バグ探し:以下のコードにあるバグをすべて見つけて修正する(ヒント:3 つ以上ある)
func process(items []string) map[string]int {
results := map[string]int{}
var wg sync.WaitGroup
for _, item := range items {
go func() {
results[item] = len(item)
wg.Done()
}()
wg.Add(1)
}
wg.Wait()
return results
}
次回予告
第6回では、ここまでの全要素を組み合わせて ポートフォリオとして公開できる API を完成させます。レイヤードアーキテクチャ、依存注入、エラーの層間変換、テスト戦略、README の書き方、デプロイまでを扱います。
Analyzegear