B システム設計/インフラ — GKE Autopilot × OTel 1.x × DataDog Agent v7 OTLP native ingestion × Argo Workflows W3C TraceContext 伝播 × Unified Service Tagging(MOps Bad→Good)

2026-06-16 (Day 71) 火曜 B: システム設計/インフラ ★★★★☆ OTel 1.x / DataDog Agent v7 / OTLP gRPC GKE Autopilot / Argo Workflows / W3C TraceContext

概要

🔭

OTel 1.x + BatchSpanProcessor で DataDog OTLP に送信

print()logging.basicConfig の混在はトレースコンテキストなし・非構造化ログの典型パターン。TracerProvider + OTLPSpanExporter(port 4317)+ BatchSpanProcessor を初期化し、DataDog Agent v7 の OTLP native ingestion に送信する。SimpleSpanProcessor は同期エクスポートで本番不向き。

📡

Downward API で status.hostIP → OTLP エンドポイント動的解決

GKE Autopilot では hostNetwork: true が禁止されており、localhost:4317 は DataDog DaemonSet Agent に届かない。status.hostIP(Downward API)を環境変数 NODE_IP で取得し、OTEL_EXPORTER_OTLP_ENDPOINT=http://$(NODE_IP):4317 で同一ノードの Agent に接続する。

🔗

W3C TraceContext で Argo Workflows ステップ間を1トレースに統合

Argo Workflows の各ステップは別 Pod で動くため gRPC ヘッダー伝播が使えない。TraceContextTextMapPropagator.inject()traceparent 文字列を生成し、Argo の outputs.parametersinputs.parameters で引き継ぐことで、バッチ全体を1トレースとして DataDog APM で可視化できる。

🏷️

Unified Service Tagging で APM・Logs・SLO を自動相関

DD_SERVICE / DD_ENV / DD_VERSION を Pod ラベル(tags.datadoghq.com/*)+ コンテナ環境変数として設定することで DataDog サービスマップ・エラートラッキング・デプロイ追跡が有効になる。OTel Resource の deployment.environment も同じ値に統一する。

問題

ECサイト MOps チームでは、Argo Workflows 上のバッチ処理(メール配信・キャンペーン集計・注文同期)を GKE Autopilot で運用している。現状のコードと Kubernetes マニフェストには Observability に関する7つの問題 が潜んでいる。問題点を全て洗い出し、OTel 1.x(OpenTelemetry Stable)+ DataDog Agent v7(OTLP native ingestion)+ 構造化ログ(JSON)+ GKE Autopilot 対応マニフェスト を使った Bad→Good リファクタリングを行え。

制約・前提条件

  • OTel SDK: opentelemetry-sdk 1.x stable(Python 3.12)
  • DataDog Agent v7 を DaemonSet で GKE Autopilot にデプロイ済み
  • OTLP gRPC エンドポイント: http://$(OTEL_EXPORTER_OTLP_ENDPOINT):4317(DataDog Agent v7 OTLP native ingestion)
  • Unified Service Tagging(DD_SERVICE / DD_ENV / DD_VERSION)を全コンテナに設定すること
  • GKE Autopilot 制約: hostPID / hostNetwork / privileged は禁止
  • ログは Cloud Logging JSON 形式(jsonPayload)で出力し、trace_id / span_id を埋め込むこと
  • Argo Workflows ステップ間でコンテキスト伝播(W3C TraceContext)を行うこと
期待する回答形式: 問題点の列挙(番号付き)+ 改善後コード(Python OTel instrumentation)+ 改善後 Kubernetes マニフェスト + Argo Workflows テンプレート + 設計意図の説明

悪いコード (Before)

このコード・マニフェストには 7つの Observability 設計の問題 が隠れています。
bad_mops_batch.py — print 混在・localhost ハードコード・span.end() 漏れ・コンテキスト伝播なし
import logging
import os

from opentelemetry import trace
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import SimpleSpanProcessor  # 問題①
from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter

# 問題③: OTLP エンドポイントをハードコード
OTLP_ENDPOINT = "http://localhost:4317"

# 問題②: basicConfig では trace_id がログに入らない
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

def setup_tracer():
    provider = TracerProvider()
    exporter = OTLPSpanExporter(endpoint=OTLP_ENDPOINT, insecure=True)
    # 問題①: SimpleSpanProcessor は同期エクスポート(本番不適)
    provider.add_span_processor(SimpleSpanProcessor(exporter))
    trace.set_tracer_provider(provider)
    return trace.get_tracer("mops-email-batch")

def run_batch(batch_id: str):
    tracer = setup_tracer()
    # 問題④: 手動 span.end()(例外時に漏れる)
    span = tracer.start_span("campaign_email_batch")
    print(f"[INFO] batch started: {batch_id}")  # 問題⑤: print デバッグ

    # 問題⑦: Unified Service Tagging なし → DataDog APM でサービスが不明
    try:
        result = _send_emails(batch_id)
        logger.info(f"sent: {result}")   # trace_id なし
        span.end()  # 問題④: 例外が起きたら呼ばれない
    except Exception as e:
        print(f"[ERROR] {e}")   # 問題⑤: 非構造化エラーログ
        span.end()   # finally にしないと漏れる
        raise

    # 問題⑥: traceparent を次ステップに渡さない → 別トレースになる
    return {"batch_id": batch_id, "status": "ok"}

def _send_emails(batch_id: str) -> dict:
    return {"sent": 1000, "failed": 0}
問題点サマリー(7点)
1SimpleSpanProcessor(同期エクスポート) — リクエストごとに同期 OTLP 送信。BatchSpanProcessor に変更して非同期バッチ化する
2logging.basicConfig(trace_id 未挿入) — DataDog Logs ↔ APM の相関不可。structlog + OTel コンテキスト挿入プロセッサに変更
3localhost:4317 ハードコード — GKE Autopilot の別ノードでは DaemonSet Agent に届かない。Downward API status.hostIP で動的解決
4手動 span.end() — 例外経路で漏れる。with tracer.start_as_current_span(...) で自動化
5print() デバッグ混在 — Cloud Logging で検索・集計不可。structlog JSON ロガーに統一
6Argo Workflows ステップ間 TraceContext 伝播なし — 各ステップが別トレースになる。traceparentoutputs.parameters で引き継ぐ
7Unified Service Tagging なし — DataDog のサービスマップ・エラートラッキング・デプロイ追跡が機能しない。Pod ラベル + 環境変数に DD_SERVICE/DD_ENV/DD_VERSION を設定

ヒント(段階的開示)

ヒント1 — 方向性
現状コードは print()logging.basicConfig の混在でトレースコンテキストが一切ない。DataDog では DD_SERVICE / DD_ENV / DD_VERSION の3タグが Unified Service Tagging の必須条件。OTel SDK の TracerProvider を初期化し、OTLPSpanExporter を DataDog Agent の OTLP エンドポイント(port 4317)に向けること。GKE Autopilot では hostPID が使えないため DaemonSet の DataDog Agent から Pod メタデータを取得するには Downward APIstatus.hostIP)を活用する。
ヒント2 — アプローチ
  • 問題①: SimpleSpanProcessorBatchSpanProcessor(非同期バッチエクスポート)に変更
  • 問題②: structlog + trace.get_current_span().get_span_context()trace_id / span_id を JSON ログに自動挿入
  • 問題③: Downward API status.hostIPNODE_IP 環境変数で取得し、OTEL_EXPORTER_OTLP_ENDPOINT=http://$(NODE_IP):4317 に設定
  • 問題④: with tracer.start_as_current_span(...) as span: で例外時も end() + StatusCode.ERROR を自動設定
  • 問題⑤: print() を全て削除し structlog.get_logger().bind(...).info(...) に統一
  • 問題⑥: TraceContextTextMapPropagator.inject(carrier)traceparent 文字列を生成し、Argo の outputs.parameters に書き出す。次ステップで PROPAGATOR.extract(carrier) して復元
  • 問題⑦: Pod ラベル tags.datadoghq.com/service + env + version を設定し、コンテナ環境変数 DD_SERVICE / DD_ENV / DD_VERSION を Downward API で参照
ヒント3 — コードの骨格
# OTel TracerProvider の正しい初期化パターン
from opentelemetry import trace
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor
from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter
from opentelemetry.sdk.resources import Resource, SERVICE_NAME, SERVICE_VERSION

def init_tracer() -> trace.Tracer:
    resource = Resource.create({
        SERVICE_NAME: os.environ["DD_SERVICE"],
        SERVICE_VERSION: os.environ["DD_VERSION"],
        "deployment.environment": os.environ["DD_ENV"],
    })
    provider = TracerProvider(resource=resource)
    exporter = OTLPSpanExporter(
        endpoint=os.environ["OTEL_EXPORTER_OTLP_ENDPOINT"],
        insecure=True,
    )
    provider.add_span_processor(BatchSpanProcessor(exporter))  # 非同期バッチ
    trace.set_tracer_provider(provider)
    return trace.get_tracer(os.environ["DD_SERVICE"])

# W3C TraceContext 伝播(Argo Workflows ステップ間)
from opentelemetry.propagators.textmap import TraceContextTextMapPropagator

def extract_context(traceparent: str | None):
    if not traceparent:
        return trace.context_api.get_current()
    return TraceContextTextMapPropagator().extract({"traceparent": traceparent})

def inject_context() -> str:
    carrier: dict[str, str] = {}
    TraceContextTextMapPropagator().inject(carrier)
    return carrier.get("traceparent", "")

# コンテキストマネージャで span.end() を保証
with tracer.start_as_current_span("batch_step", context=parent_ctx) as span:
    try:
        do_work()
        span.set_status(StatusCode.OK)
    except Exception as e:
        span.record_exception(e)
        span.set_status(StatusCode.ERROR, str(e))
        raise

問題点分析(7点)

#問題点分類改善方法
1SimpleSpanProcessor(同期エクスポート)パフォーマンスBatchSpanProcessor で非同期バッチ化
2logging.basicConfig(trace_id 未挿入)可観測性structlog + OTel コンテキスト挿入プロセッサ
3localhost:4317 ハードコード信頼性Downward API status.hostIP で動的解決
4手動 span.end()(例外漏れ)エラー処理start_as_current_span コンテキストマネージャ
5print() デバッグ混在可観測性structlog JSON ロガーに統一
6Argo ステップ間 TraceContext なし可観測性W3C traceparent を outputs.parameters で伝達
7Unified Service Tagging なしDataDog 設定Pod ラベル + DD_SERVICE/DD_ENV/DD_VERSION

アーキテクチャ図 — Bad vs Good の変換フロー

Bad(変更前)— Observability なし Argo Step 1: fetch-segments Pod ① SimpleSpanProcessor(同期)② print() ③ localhost:4317 ⚠️ 別ノードの DD Agent に届かない / 同期エクスポートでレイテンシ増加 ④ span.end() 手動(例外で漏れ) ⑦ DD_SERVICE なし(サービス不明) × traceparent なし Argo Step 2: send-emails Pod ⑥ 前ステップの trace_id を引き継げない → 別 Trace 生成 ⚠️ DataDog APM で3ステップが別々のトレースに分散 ② logging.basicConfig → trace_id がログに入らない(相関不可) Argo Step 3: update-kpi Pod ⑦ DD_ENV / DD_VERSION もなし → DataDog エラートラッキング機能しない ⚠️ サービスマップ / デプロイ追跡が完全に無効 DataDog APM(Bad) 3つの孤立トレース / trace_id なしログ / サービス不明 バッチ全体の所要時間・エラー原因が追跡不可能 DataDog Agent DaemonSet(Bad) ⑥ hostPID: true / privileged: true — GKE Autopilot で禁止 ⚠️ DaemonSet が起動せず → Observability 全滅 OTLP port 4317 も設定なし → Pod からの OTLP 送信を受け付けない Good(変更後)— 統合 Observability Argo Step 1: fetch-segments Pod BatchSpanProcessor(非同期) / structlog JSON ログ(trace_id 自動挿入) OTEL_EXPORTER_OTLP_ENDPOINT=http://$(NODE_IP):4317(Downward API) with start_as_current_span() → 例外時も end() 保証 + record_exception → traceparent を /tmp/traceparent に書き出し(outputs.parameters) traceparent 引き継ぎ Argo Step 2: send-emails Pod extract_trace_context(traceparent) で前ステップの context を復元 同一 trace_id で子スパン生成 → DataDog APM で1フレームグラフに統合 structlog JSON: {"trace_id":"4bf9...", "span_id":"00f0...", "event":"emails_sent"} Argo Step 3: update-kpi Pod DD_SERVICE=mops-email-batch / DD_ENV=production / DD_VERSION=1.4.2 OTel Resource: SERVICE_NAME / SERVICE_VERSION / deployment.environment を統一 DataDog サービスマップ・エラートラッキング・デプロイ追跡が全て有効 DataDog APM(Good) 3ステップが1トレース(同一 trace_id)/ フレームグラフで全体可視化 Cloud Logging ↔ DataDog Logs 相関(logging.googleapis.com/trace) SLO: campaign_email_batch P99 < 10分 / エラーレート < 0.1% を監視 DataDog Agent DaemonSet(GKE Autopilot 対応) hostPID: false / privileged: false / allowPrivilegeEscalation: false DD_OTLP_CONFIG_RECEIVER_PROTOCOLS_GRPC_ENDPOINT: 0.0.0.0:4317 Downward API: DD_KUBERNETES_KUBELET_NODENAME=spec.nodeName 修正

模範解答

"""mops_batch_tracer.py — OTel 1.x + DataDog OTLP + structlog JSON + W3C TraceContext"""
from __future__ import annotations
import logging, os, sys
from typing import Final
import structlog
from opentelemetry import trace
from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter
from opentelemetry.propagators.textmap import TraceContextTextMapPropagator
from opentelemetry.sdk.resources import SERVICE_NAME, SERVICE_VERSION, Resource
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor  # 修正①: Batch
from opentelemetry.trace import StatusCode

# ── 名前付き定数 ──────────────────────────────
OTEL_ENDPOINT_ENV: Final[str] = "OTEL_EXPORTER_OTLP_ENDPOINT"
TRACEPARENT_KEY:   Final[str] = "traceparent"   # W3C TraceContext
PROPAGATOR: Final[TraceContextTextMapPropagator] = TraceContextTextMapPropagator()


# ── 修正②: structlog JSON ロガー(trace_id 自動挿入)──────────
def _configure_structlog() -> None:
    def _inject_trace_context(logger, method, event_dict: dict) -> dict:
        span = trace.get_current_span()
        ctx  = span.get_span_context()
        if ctx.is_valid:
            event_dict["logging.googleapis.com/trace"] = (
                f"projects/{os.environ.get('GCP_PROJECT_ID','')}"
                f"/traces/{ctx.trace_id:032x}"
            )
            event_dict["trace_id"] = f"{ctx.trace_id:032x}"
            event_dict["span_id"]  = f"{ctx.span_id:016x}"
        return event_dict

    structlog.configure(
        processors=[
            structlog.contextvars.merge_contextvars,
            structlog.processors.add_log_level,
            structlog.processors.TimeStamper(fmt="iso"),
            _inject_trace_context,            # OTel context 自動挿入
            structlog.processors.JSONRenderer(),   # Cloud Logging jsonPayload 互換
        ],
        wrapper_class=structlog.make_filtering_bound_logger(logging.INFO),
        context_class=dict,
        logger_factory=structlog.PrintLoggerFactory(sys.stdout),
    )


# ── TracerProvider 初期化 ────────────────────────────────
def init_tracer() -> trace.Tracer:
    """OTel TracerProvider を初期化し DataDog Agent OTLP エンドポイントに接続。

    Args:
        なし(環境変数 OTEL_EXPORTER_OTLP_ENDPOINT / DD_* を参照)

    Returns:
        Tracer: サービス名・バージョンに紐づいた OTel Tracer

    Raises:
        KeyError: 必須環境変数が未設定の場合
    """
    resource = Resource.create({
        SERVICE_NAME:    os.environ["DD_SERVICE"],    # 修正⑦: Unified Service Tagging
        SERVICE_VERSION: os.environ["DD_VERSION"],
        "deployment.environment": os.environ["DD_ENV"],
    })
    provider = TracerProvider(resource=resource)
    exporter = OTLPSpanExporter(
        endpoint=os.environ[OTEL_ENDPOINT_ENV],   # 修正③: ハードコード禁止
        insecure=True,
    )
    # 修正①: BatchSpanProcessor で非同期バッチエクスポート
    provider.add_span_processor(BatchSpanProcessor(exporter))
    trace.set_tracer_provider(provider)
    return trace.get_tracer(os.environ["DD_SERVICE"], os.environ["DD_VERSION"])


# ── 修正⑥: W3C TraceContext 伝播ユーティリティ ───────────────
def extract_trace_context(traceparent: str | None) -> trace.Context:
    """Argo Workflows outputs.parameters から W3C traceparent を復元。"""
    if not traceparent:
        return trace.context_api.get_current()
    return PROPAGATOR.extract({TRACEPARENT_KEY: traceparent})


def inject_trace_context() -> str:
    """現在スパンから W3C traceparent 文字列を生成(次ステップへの引き渡し用)。"""
    carrier: dict[str, str] = {}
    PROPAGATOR.inject(carrier)
    return carrier.get(TRACEPARENT_KEY, "")


# ── バッチ処理本体(改善後)────────────────────────────────
def run_campaign_email_batch(batch_id: str, traceparent: str | None = None) -> str:
    """キャンペーンメール送信バッチ(OTel 計装版)。

    Args:
        batch_id: バッチ識別子
        traceparent: Argo Workflows 前ステップからの W3C traceparent

    Returns:
        str: 次ステップへ渡す traceparent 文字列
    """
    _configure_structlog()
    tracer = init_tracer()
    logger = structlog.get_logger().bind(batch_id=batch_id)

    parent_ctx = extract_trace_context(traceparent)   # 修正⑥: コンテキスト復元

    # 修正④: start_as_current_span で span.end() を自動保証
    with tracer.start_as_current_span(
        "campaign_email_batch",
        context=parent_ctx,
        attributes={"batch.id": batch_id, "batch.type": "email"},
    ) as root_span:
        # 修正⑤: print() を structlog に統一(trace_id 自動挿入)
        logger.info("batch_started")
        try:
            with tracer.start_as_current_span("fetch_segments") as seg_span:
                count = 3   # 実際は BigQuery から取得
                seg_span.set_attribute("segment.count", count)
                logger.info("segments_fetched", count=count)

            with tracer.start_as_current_span("send_emails") as send_span:
                sent, failed = 1000, 2
                send_span.set_attribute("email.sent",   sent)
                send_span.set_attribute("email.failed", failed)
                if failed > 0:
                    send_span.set_status(StatusCode.ERROR, f"{failed}件の送信失敗")
                logger.info("emails_sent", sent=sent, failed=failed)

            with tracer.start_as_current_span("update_kpi"):
                logger.info("kpi_updated")

            root_span.set_status(StatusCode.OK)
            logger.info("batch_completed")

        except Exception as exc:
            root_span.record_exception(exc)     # DataDog エラートラッキングに自動反映
            root_span.set_status(StatusCode.ERROR, str(exc))
            logger.error("batch_failed", error=str(exc))
            raise

    return inject_trace_context()   # 修正⑥: 次ステップへ traceparent を渡す

DataDog Agent DaemonSet(GKE Autopilot 対応)

apiVersion: apps/v1
kind: DaemonSet
metadata:
  name: datadog-agent
  namespace: monitoring
spec:
  selector:
    matchLabels: {app: datadog-agent}
  template:
    metadata:
      labels: {app: datadog-agent}
    spec:
      serviceAccountName: datadog-agent
      # 修正⑥: GKE Autopilot では hostPID / hostNetwork 禁止
      hostPID:     false
      hostNetwork: false
      containers:
        - name: agent
          image: gcr.io/datadoghq/agent:7
          env:
            - {name: DD_API_KEY, valueFrom: {secretKeyRef: {name: datadog-secret, key: api-key}}}
            - {name: DD_SITE, value: "datadoghq.com"}
            # OTLP native ingestion(gRPC port 4317)
            - {name: DD_OTLP_CONFIG_RECEIVER_PROTOCOLS_GRPC_ENDPOINT, value: "0.0.0.0:4317"}
            - {name: DD_OTLP_CONFIG_RECEIVER_PROTOCOLS_HTTP_ENDPOINT, value: "0.0.0.0:4318"}
            # Autopilot: Downward API でノード名を取得
            - name: DD_KUBERNETES_KUBELET_NODENAME
              valueFrom:
                fieldRef:
                  fieldPath: spec.nodeName
            # Autopilot: system-probe は無効(特権不要)
            - {name: DD_SYSTEM_PROBE_ENABLED, value: "false"}
          ports:
            - {containerPort: 4317, protocol: TCP}   # OTLP gRPC
            - {containerPort: 4318, protocol: TCP}   # OTLP HTTP
          securityContext:
            allowPrivilegeEscalation: false   # 修正⑥: Autopilot 対応
            readOnlyRootFilesystem:   false
            runAsNonRoot:             false
            capabilities:
              add: ["SYS_PTRACE"]
              drop: ["NET_RAW"]
          resources:
            requests: {cpu: "200m", memory: "256Mi"}
            limits:   {cpu: "200m", memory: "256Mi"}
---
# Argo Workflows バッチ Pod — Unified Service Tagging + Downward API
apiVersion: v1
kind: Pod
metadata:
  name: campaign-email-batch
  labels:
    tags.datadoghq.com/service: "mops-email-batch"   # 修正⑦: Unified Service Tagging
    tags.datadoghq.com/env:     "production"
    tags.datadoghq.com/version: "1.4.2"
spec:
  serviceAccountName: mops-batch-sa
  containers:
    - name: batch
      image: asia-northeast1-docker.pkg.dev/my-project/mops/batch@sha256:${IMAGE_DIGEST}
      env:
        # 修正⑦: Downward API でラベルから DD_* を環境変数に展開
        - name: DD_SERVICE
          valueFrom: {fieldRef: {fieldPath: "metadata.labels['tags.datadoghq.com/service']"}}
        - name: DD_ENV
          valueFrom: {fieldRef: {fieldPath: "metadata.labels['tags.datadoghq.com/env']"}}
        - name: DD_VERSION
          valueFrom: {fieldRef: {fieldPath: "metadata.labels['tags.datadoghq.com/version']"}}
        # 修正③: status.hostIP で同一ノードの DD Agent OTLP エンドポイントを動的解決
        - name: NODE_IP
          valueFrom: {fieldRef: {fieldPath: status.hostIP}}
        - name: OTEL_EXPORTER_OTLP_ENDPOINT
          value: "http://$(NODE_IP):4317"
        - name: GCP_PROJECT_ID
          valueFrom: {configMapKeyRef: {name: mops-config, key: gcp_project_id}}
      securityContext:
        allowPrivilegeEscalation: false
        readOnlyRootFilesystem:   true
        runAsNonRoot:             true
        runAsUser:                1000
        capabilities: {drop: ["ALL"]}
      resources:
        requests: {cpu: "500m", memory: "512Mi"}
        limits:   {cpu: "500m", memory: "512Mi"}

ConfigMap — OTLP 設定

apiVersion: v1
kind: ConfigMap
metadata:
  name: mops-config
  namespace: mops
data:
  gcp_project_id: "my-gcp-project-id"

Argo Workflows DAG テンプレート(W3C TraceContext 伝播)

apiVersion: argoproj.io/v1alpha1
kind: Workflow
metadata:
  generateName: campaign-batch-
spec:
  entrypoint: campaign-pipeline
  templates:
    - name: campaign-pipeline
      dag:
        tasks:
          - name: fetch-segments
            template: batch-step
            arguments:
              parameters:
                - {name: step,        value: "fetch_segments"}
                - {name: traceparent, value: ""}   # 最初: 空(新規 trace 生成)

          - name: send-emails
            template: batch-step
            dependencies: [fetch-segments]
            arguments:
              parameters:
                - {name: step, value: "send_emails"}
                # 修正⑥: 前ステップの traceparent を引き継ぐ
                - name: traceparent
                  value: "{{tasks.fetch-segments.outputs.parameters.traceparent}}"

          - name: update-kpi
            template: batch-step
            dependencies: [send-emails]
            arguments:
              parameters:
                - {name: step, value: "update_kpi"}
                - name: traceparent
                  value: "{{tasks.send-emails.outputs.parameters.traceparent}}"

    - name: batch-step
      inputs:
        parameters:
          - {name: step}
          - {name: traceparent}
      outputs:
        parameters:
          # 次ステップへ traceparent を引き渡し(W3C TraceContext 伝播)
          - name: traceparent
            valueFrom:
              path: /tmp/traceparent   # Python コードが inject_trace_context() の結果を書く
      container:
        image: asia-northeast1-docker.pkg.dev/my-project/mops/batch@sha256:${IMAGE_DIGEST}
        command: ["python", "-m", "mops_batch_tracer"]
        args:
          - "--step={{inputs.parameters.step}}"
          - "--traceparent={{inputs.parameters.traceparent}}"
        envFrom:
          - configMapRef: {name: mops-config}
        env:
          - name: DD_SERVICE
            valueFrom: {fieldRef: {fieldPath: "metadata.labels['tags.datadoghq.com/service']"}}
          - name: DD_ENV
            valueFrom: {fieldRef: {fieldPath: "metadata.labels['tags.datadoghq.com/env']"}}
          - name: DD_VERSION
            valueFrom: {fieldRef: {fieldPath: "metadata.labels['tags.datadoghq.com/version']"}}
          - name: NODE_IP
            valueFrom: {fieldRef: {fieldPath: status.hostIP}}
          - name: OTEL_EXPORTER_OTLP_ENDPOINT
            value: "http://$(NODE_IP):4317"
問題修正内容効果
① SimpleSpanProcessorBatchSpanProcessor(非同期)レイテンシへの影響ゼロ・バックプレッシャー対応
② logging.basicConfigstructlog + trace_id 自動挿入Cloud Logging ↔ DataDog Logs 相関・JSON 検索
③ localhost:4317 ハードコードstatus.hostIP Downward APIGKE Autopilot 別ノードでも確実に DD Agent に到達
④ 手動 span.end()start_as_current_span CM例外時も必ず end() + ERROR ステータス記録
⑤ print() デバッグstructlog JSON 統一Cloud Logging でフィルタ・集計・アラート可能
⑥ TraceContext 伝播なしtraceparent outputs.parameters3ステップが1フレームグラフ / バッチ全体 SLO 監視
⑦ Unified Service Tagging なしPod ラベル + DD_SERVICE/ENV/VERSIONサービスマップ・エラートラッキング・デプロイ追跡

ポイント解説

1 BatchSpanProcessor vs SimpleSpanProcessor
SimpleSpanProcessor はスパンが完了するたびに同期でエクスポートし、OTLP gRPC の RTT(往復遅延)がそのままアプリケーションのレイテンシに加算される。BatchSpanProcessor はスパンをメモリにバッファリングし、設定した max_export_batch_size(デフォルト512)に達するか schedule_delay_millis(デフォルト5000ms)経過したら非同期でエクスポートする。DataDog Agent v7 の OTLP gRPC エンドポイント(port 4317)はローカルネットワーク内での受信なので通常 <1ms だが、本番では必ず BatchSpanProcessor を使うこと。
2 Downward API による status.hostIP の取得
GKE Autopilot では hostNetwork: true が禁止されているため localhost はコンテナ内部のループバックのみを指し、DaemonSet の DataDog Agent に届かない。status.hostIP(Downward API)を環境変数 NODE_IP で取得すると、Pod が実行されているノードのIPアドレスが得られる。DataDog Agent DaemonSet は各ノードで動作しているため、http://$(NODE_IP):4317 で同一ノードの Agent に確実に OTLP データを送信できる。$(NODE_IP) の展開は Kubernetes が自動で行う(環境変数参照の $(VAR_NAME) 構文)。
3 W3C TraceContext の traceparent フォーマット
traceparent は W3C TraceContext 仕様(RFC)で定義された文字列: 00-{32桁 trace_id}-{16桁 parent_span_id}-{8桁 flags}。例: 00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01。Argo Workflows のステップ間は gRPC/HTTP ヘッダーで伝播できないため、outputs.parameters(ファイル書き出し)→ inputs.parameters(引数受け取り)のパターンでバトンリレーする。TraceContextTextMapPropagator.inject(carrier) がキャリア dict に {"traceparent": "00-..."} を書き込み、extract(carrier) が復元する。
4 start_as_current_span のエラー自動記録
with tracer.start_as_current_span("step") as span: ブロック内で例外が発生すると、コンテキストマネージャは自動で span.set_status(StatusCode.ERROR) + span.end() を呼び出す。さらに span.record_exception(exc) を明示的に呼ぶと例外のスタックトレースが DataDog APM のエラーパネルと連携する。DataDog の dd_error.message / dd_error.type / dd_error.stack タグとして自動でマッピングされる。
5 Unified Service Tagging と OTel Resource の一致
DataDog の Unified Service Tagging は3タグ(service/env/version)を全てのテレメトリデータ(メトリクス・ログ・トレース)に付与することで相関を実現する。OTel Resource の SERVICE_NAME は DataDog の service タグに、deployment.environmentenv タグにマッピングされる。Pod ラベル(tags.datadoghq.com/*)+ コンテナ環境変数 + OTel Resource の3箇所を同一値に揃えることで DataDog のサービスマップ・SLO・デプロイ追跡が全て正しく機能する。
6 structlog + Cloud Logging トレース相関
Cloud Logging の jsonPayloadlogging.googleapis.com/trace フィールドを含めると、Cloud Logging コンソールでトレースビューアとログを相関表示できる。フォーマット: "projects/{PROJECT_ID}/traces/{trace_id}"。DataDog Logs でも trace_id フィールドを dd.trace_id と同じ値に設定することで APM ↔ Logs の自動リンクが有効になる。OTel の trace_id は 128bit(32桁16進)、DataDog は 64bit(符号なし整数)なので変換が必要な場合は dd-trace のドキュメントを参照。

実務への応用

  • DataDog SLO の設定: campaign_email_batch スパンの P99 レイテンシ < 10分 を SLO として設定し、超過時に PagerDuty / Slack アラートを発火させることでバッチ遅延を自動検知できる。DataDog APM の Monitor から Trace Analyticsresource_name:campaign_email_batch をフィルタして設定する。
  • opentelemetry-instrumentation-httpx: HTTPXClientInstrumentor().instrument() を追加すると、Cloud Run API 呼び出しや外部 API リクエストが自動でスパンとして記録される(zero-code instrumentation)。手動で tracer.start_as_current_span("http_request") を書く必要がなくなる。
  • DataDog Agent の OTLP ポート確認: kubectl exec -n monitoring datadog-agent-XXXXX -- agent status | grep otlp で OTLP receiver が正常に起動しているか確認できる。DD_OTLP_CONFIG_RECEIVER_PROTOCOLS_GRPC_ENDPOINT0.0.0.0:4317 に設定されていないと Pod からの送信を受け付けない。
  • MOps バッチの可視化例: Argo Workflows の fetch-segments → send-emails → update-kpi の3ステップが DataDog APM で1フレームグラフに統合され、「セグメント取得に2分・メール送信に7分・KPI更新に30秒」という内訳がひと目でわかる。エラー時は send-emails スパンが赤くなり、email.failed カスタム属性で失敗件数を即確認できる。

今日のまとめ

GKE Autopilot × Argo Workflows の Observability 設計チェックリストは「① BatchSpanProcessor(非同期)② structlog JSON + trace_id 挿入 ③ status.hostIP Downward API で OTLP エンドポイント動的解決 ④ start_as_current_spanend() 自動保証 ⑤ print() 全廃 ⑥ W3C traceparent を Argo outputs.parameters で伝達 ⑦ Unified Service Tagging(Pod ラベル + DD_SERVICE/ENV/VERSION)」の7点。

print() + localhost:4317 + SimpleSpanProcessor の組み合わせは GKE Autopilot では Observability が完全に機能しない典型的なアンチパターン。DataDog Agent の OTLP native ingestion は Prometheus や Jaeger を別途立てる必要がなく、既存の DataDog 環境に OTel を後付けで統合できる最も実用的な構成。

自己評価

自分の回答

気づき・メモ