システム設計/インフラ — OpenTelemetry 分散トレーシング × DataDog 相関

2026-05-12 (Day 29) 火曜 B: システム設計/インフラ ★★★☆☆ OpenTelemetry SDK 1.x / DataDog Agent v7 Argo Workflows トレースコンテキスト伝播

概要

🔍

W3C Trace Context 伝播

inject(headers) が自動で traceparent を追加。受信側は extract(headers) でコンテキストを復元。

🐶

DataDog ログ相関

ログの dd.trace_id フィールドに 32桁16進数の trace_id を埋め込むことで、DataDog がログとトレースを自動相関。

⚙️

Argo ステップ間伝播

各ステップはPodが独立するため、outputs.parameterstraceparent を次ステップの環境変数に渡す。

🌊

structlog JSON 出力

JSONRenderer() により1行JSONログ。DataDog Agent が自動収集・パース・相関表示する。

問題

Argo Workflows + Cloud Run + BigQuery を組み合わせた販促バッチ処理基盤で以下の問題が発生しています。

問題現状
バッチ状況確認GCPコンソールで毎回手動確認している
遅延原因特定Cloud Runの処理遅延がどのサービスか特定に30分以上かかる
DataDogアラートアラートを受け取るがどのトレースと紐づいているか不明

以下の3つを実現する可観測性設計を行ってください。

  1. 分散トレーシングの導入: OpenTelemetry SDK で process_coupon_batch から Cloud Run への HTTP リクエストまでをE2Eでトレース
  2. 構造化ログ: DataDog が Trace ID と自動相関できる形式でログを出力(dd.trace_id を含む)
  3. Argo Workflows への伝播: Workflow 全体にトレースコンテキスト(traceparent)を伝播させる設計方針を示す
期待する回答形式: コード(Python)+ 設計説明文(YAML/構成図は任意)

現状のコード(問題あり)

トレーシングなし・構造化ログなし・全例外を握り潰す設計。
# batch_processor.py(現状コード)
import requests
import logging

logging.basicConfig(level=logging.DEBUG)
logger = logging.getLogger(__name__)

def process_coupon_batch(coupon_ids: list[str]) -> dict:
    results = {}
    for coupon_id in coupon_ids:
        try:
            logger.debug(f"Processing {coupon_id}")
            resp = requests.post(
                "https://coupon-service.example.com/apply",
                json={"coupon_id": coupon_id},
                timeout=5
            )
            resp.raise_for_status()
            results[coupon_id] = "ok"
            logger.debug(f"Done {coupon_id}")
        except Exception as e:
            logger.error(f"Error: {e}")
            results[coupon_id] = "error"
    return results

ヒント(段階的開示)

ヒント1 — 方向性
OpenTelemetry の「トレース → スパン → コンテキスト伝播」の3層構造を理解することがカギ。ログとトレースの相関は、ログ出力時に現在のスパンの trace_id を取得して埋め込むことで実現できる。DataDog は OTLP 形式で受け取った後、dd.trace_id フィールドでログと自動相関する。
ヒント2 — アプローチ(構成要素)
  • TracerProvider + OTLPSpanExporter → DataDog Agent(4317番ポート)へエクスポート
  • @tracer.start_as_current_span() でスパンを開始
  • HTTP リクエスト時は inject(headers) で W3C Trace Context ヘッダーを伝播
  • ログ出力時は get_current_span().get_span_context() から trace_id を取得
  • Argo Workflows での伝播は、前ステップの出力を traceparent として次ステップの環境変数に渡す
ヒント3 — Python コードの骨格
from opentelemetry import trace
from opentelemetry.propagate import inject
from opentelemetry.trace import get_current_span

# スパンの作成と構造化ログ
with tracer.start_as_current_span("process_coupon_batch") as span:
    ctx = get_current_span().get_span_context()
    # trace_id を16進数に変換してログに埋め込む
    trace_id = format(ctx.trace_id, "032x")
    logger.info("processing", dd={"trace_id": trace_id}, ...)

    # HTTPリクエスト時はヘッダーにコンテキストを注入
    headers = {}
    inject(headers)  # traceparent, tracestate を自動で追加
    resp = requests.post(url, headers=headers, ...)

トレース構造図 — E2Eトレース伝播

Argo Workflow Step: extract output: traceparent env: TRACEPARENT batch_processor.py process_coupon_batch Root Span apply_coupon (×N) Child Span × クーポン数 inject(headers) → traceparent HTTP POST traceparent header coupon-service Cloud Run extract(headers) → ctx DataDog(OTLP ingestion via Agent :4317) Trace View スパンのウォーターフォール Log Correlation dd.trace_id で自動相関 TRACEPARENT env → extract() で復元

模範解答

# batch_processor.py(改善後)
"""販促クーポンバッチ処理モジュール。

OpenTelemetry による分散トレーシングと構造化ログ(DataDog 相関対応)を実装。
"""
from __future__ import annotations

import os

import requests
import structlog
from opentelemetry import trace
from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter
from opentelemetry.propagate import inject
from opentelemetry.sdk.resources import Resource
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor
from opentelemetry.trace import get_current_span

OTEL_ENDPOINT = os.getenv("OTEL_EXPORTER_OTLP_ENDPOINT", "http://datadog-agent:4317")
COUPON_SERVICE_URL = os.getenv("COUPON_SERVICE_URL", "https://coupon-service.example.com")
REQUEST_TIMEOUT_SEC = 5
SERVICE_NAME = "batch-processor"


def _setup_tracer_provider() -> TracerProvider:
    resource = Resource.create({"service.name": SERVICE_NAME})
    exporter = OTLPSpanExporter(endpoint=OTEL_ENDPOINT, insecure=True)
    provider = TracerProvider(resource=resource)
    provider.add_span_processor(BatchSpanProcessor(exporter))
    trace.set_tracer_provider(provider)
    return provider


_provider = _setup_tracer_provider()
_tracer = trace.get_tracer(SERVICE_NAME)


def _get_logger() -> structlog.BoundLogger:
    structlog.configure(
        processors=[
            structlog.processors.add_log_level,
            structlog.processors.TimeStamper(fmt="iso"),
            structlog.processors.JSONRenderer(),
        ]
    )
    return structlog.get_logger()


logger = _get_logger()


def _current_trace_id() -> str:
    """DataDog 形式の trace_id(32桁16進数)を返す。"""
    ctx = get_current_span().get_span_context()
    if not ctx.is_valid:
        return "0" * 32
    return format(ctx.trace_id, "032x")


def _current_span_id() -> str:
    ctx = get_current_span().get_span_context()
    if not ctx.is_valid:
        return "0" * 16
    return format(ctx.span_id, "016x")


def _post_with_trace(url: str, payload: dict) -> requests.Response:
    """W3C Trace Context ヘッダーを注入して HTTP POST を送信する。"""
    headers: dict[str, str] = {"Content-Type": "application/json"}
    inject(headers)  # traceparent / tracestate を自動追加
    resp = requests.post(url, json=payload, headers=headers, timeout=REQUEST_TIMEOUT_SEC)
    resp.raise_for_status()
    return resp


def process_coupon_batch(coupon_ids: list[str]) -> dict[str, str]:
    results: dict[str, str] = {}

    with _tracer.start_as_current_span("process_coupon_batch") as batch_span:
        logger.info(
            "batch_started",
            dd={"trace_id": _current_trace_id(), "span_id": _current_span_id()},
            total=len(coupon_ids),
        )

        for coupon_id in coupon_ids:
            with _tracer.start_as_current_span(
                "apply_coupon",
                attributes={"coupon.id": coupon_id},
            ) as coupon_span:
                try:
                    _post_with_trace(f"{COUPON_SERVICE_URL}/apply", {"coupon_id": coupon_id})
                    results[coupon_id] = "ok"
                    coupon_span.set_attribute("coupon.result", "ok")
                    logger.info(
                        "coupon_applied",
                        dd={"trace_id": _current_trace_id(), "span_id": _current_span_id()},
                        coupon_id=coupon_id,
                    )
                except requests.HTTPError as exc:
                    results[coupon_id] = "error"
                    coupon_span.set_status(trace.StatusCode.ERROR, str(exc))
                    coupon_span.record_exception(exc)
                    logger.error(
                        "coupon_apply_failed",
                        dd={"trace_id": _current_trace_id(), "span_id": _current_span_id()},
                        coupon_id=coupon_id,
                    )
                except requests.Timeout:
                    results[coupon_id] = "error"
                    coupon_span.set_status(trace.StatusCode.ERROR, "timeout")
                    logger.error("coupon_apply_timeout", coupon_id=coupon_id)

        success_count = sum(1 for v in results.values() if v == "ok")
        batch_span.set_attribute("batch.success_count", success_count)
        batch_span.set_attribute("batch.error_count", len(coupon_ids) - success_count)

    return results
# argo-workflow-otel.yaml
apiVersion: argoproj.io/v1alpha1
kind: Workflow
metadata:
  name: coupon-batch-workflow
spec:
  entrypoint: coupon-pipeline
  arguments:
    parameters:
      - name: traceparent
        value: ""  # 外部から注入可能(CI/CDから呼び出す場合)

  templates:
    - name: coupon-pipeline
      steps:
        - - name: extract-coupons
            template: extract
            arguments:
              parameters:
                - name: traceparent
                  value: "{{workflow.parameters.traceparent}}"

        - - name: apply-coupons
            template: apply
            arguments:
              parameters:
                - name: traceparent
                  # 前ステップの出力 traceparent を次ステップへ伝播
                  value: "{{steps.extract-coupons.outputs.parameters.traceparent}}"

    - name: extract
      inputs:
        parameters:
          - name: traceparent
      outputs:
        parameters:
          - name: traceparent
            valueFrom:
              path: /tmp/traceparent.txt  # このステップが生成したtraceparentを出力
      container:
        image: gcr.io/myproject/batch-processor:latest
        command: [python, extract_coupons.py]
        env:
          - name: TRACEPARENT
            value: "{{inputs.parameters.traceparent}}"
          - name: OTEL_EXPORTER_OTLP_ENDPOINT
            value: "http://datadog-agent:4317"
          - name: OTEL_SERVICE_NAME
            value: "coupon-batch-extract"

    - name: apply
      inputs:
        parameters:
          - name: traceparent
      container:
        image: gcr.io/myproject/batch-processor:latest
        command: [python, batch_processor_entrypoint.py]
        env:
          - name: TRACEPARENT
            value: "{{inputs.parameters.traceparent}}"
          - name: OTEL_EXPORTER_OTLP_ENDPOINT
            value: "http://datadog-agent:4317"
# batch_processor_entrypoint.py
"""Argo Workflows ステップのエントリーポイント。

前ステップから渡された TRACEPARENT 環境変数を使って
トレースコンテキストを復元し、子スパンとして処理を継続する。
"""
import os
from opentelemetry.propagate import extract
from opentelemetry import trace
from batch_processor import _tracer, process_coupon_batch

def main() -> None:
    # 環境変数から W3C Trace Context を復元
    carrier = {"traceparent": os.getenv("TRACEPARENT", "")}
    ctx = extract(carrier)  # 親スパンのコンテキストを復元

    # 復元したコンテキストを親として子スパンを開始
    with _tracer.start_as_current_span(
        "argo-apply-step",
        context=ctx,  # 前ステップのスパンの子になる
    ):
        coupon_ids = _load_coupon_ids()
        process_coupon_batch(coupon_ids)


def _load_coupon_ids() -> list[str]:
    """処理対象クーポンIDをファイルから読み込む。"""
    with open("/tmp/coupon_ids.txt") as f:
        return [line.strip() for line in f if line.strip()]


if __name__ == "__main__":
    main()

ポイント解説

1 W3C Trace Context(traceparent)の伝播
inject(headers) が自動で traceparent: 00-{trace_id}-{span_id}-01 を追加する。受信側は extract(headers) でコンテキストを復元し、同一トレースの子スパンとして継続できる。Cloud Run間の呼び出しが1つのウォーターフォールトレースとして可視化される。
2 DataDog ログ相関の仕組み
DataDog は OTLP で受け取ったスパンと、ログの dd.trace_id フィールドを照合してトレース↔ログを相関表示する。format(ctx.trace_id, "032x") で128bit整数を32桁16進数に変換することで、DataDog が期待する形式に合わせられる。
3 Argo Workflows でのステップ間伝播
K8s Job として独立したPodで実行されるArgoステップはプロセスをまたがるため、インメモリでのコンテキスト伝播は不可。outputs.parameters でファイルから traceparent を取り出し、次ステップの環境変数に渡すことで伝播を実現する。
4 スパンのエラーステータス記録
span.set_status(StatusCode.ERROR) + span.record_exception(exc) を組み合わせることで、DataDog のトレースビューでエラースパンが赤色表示され、例外スタックトレースも添付される。

実務への応用

  • 障害特定の高速化: DataDog のトレースビューで「どのクーポンIDで、何秒目に、どのCloud Runサービスがタイムアウトしたか」を即座に特定できる(現状の30分 → 2〜3分)
  • Argo Workflow の可視化: Argo UI の成功/失敗表示に加えて、DataDog のサービスマップで batch-extract → batch-apply → coupon-service の依存関係とレイテンシが可視化される
  • アラート品質の向上: DataDog モニターで trace.error_rate > 5% を検知した際に、関連するクーポンIDや失敗パターンをトレースから即座に確認できる

今日のまとめ

分散トレーシングの本質は「コンテキスト伝播(traceparent)」と「スパン生成」の2つ。プロセス境界(HTTP / Argo WorkflowステップのPod)を越えるたびに injectextract のペアでコンテキストを橋渡しすることで、複数サービスにまたがる処理を1本のトレースとして可視化できる。

DataDog のログ相関は dd.trace_id フィールドを32桁16進数で埋め込むだけ——シンプルだが見落とすと効果がゼロになる。

自己評価

自分の回答

気づき・メモ