概要
W3C Trace Context 伝播
inject(headers) が自動で traceparent を追加。受信側は extract(headers) でコンテキストを復元。
DataDog ログ相関
ログの dd.trace_id フィールドに 32桁16進数の trace_id を埋め込むことで、DataDog がログとトレースを自動相関。
Argo ステップ間伝播
各ステップはPodが独立するため、outputs.parameters で traceparent を次ステップの環境変数に渡す。
structlog JSON 出力
JSONRenderer() により1行JSONログ。DataDog Agent が自動収集・パース・相関表示する。
問題
Argo Workflows + Cloud Run + BigQuery を組み合わせた販促バッチ処理基盤で以下の問題が発生しています。
| 問題 | 現状 |
|---|---|
| バッチ状況確認 | GCPコンソールで毎回手動確認している |
| 遅延原因特定 | Cloud Runの処理遅延がどのサービスか特定に30分以上かかる |
| DataDogアラート | アラートを受け取るがどのトレースと紐づいているか不明 |
以下の3つを実現する可観測性設計を行ってください。
- 分散トレーシングの導入: OpenTelemetry SDK で
process_coupon_batchから Cloud Run への HTTP リクエストまでをE2Eでトレース - 構造化ログ: DataDog が Trace ID と自動相関できる形式でログを出力(
dd.trace_idを含む) - 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トレース伝播
模範解答
# 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 で受け取ったスパンと、ログの
DataDog は OTLP で受け取ったスパンと、ログの
dd.trace_id フィールドを照合してトレース↔ログを相関表示する。format(ctx.trace_id, "032x") で128bit整数を32桁16進数に変換することで、DataDog が期待する形式に合わせられる。
3
Argo Workflows でのステップ間伝播
K8s Job として独立したPodで実行されるArgoステップはプロセスをまたがるため、インメモリでのコンテキスト伝播は不可。
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や失敗パターンをトレースから即座に確認できる
今日のまとめ
分散トレーシングの本質は「コンテキスト伝播(
DataDog のログ相関は
traceparent)」と「スパン生成」の2つ。プロセス境界(HTTP / Argo WorkflowステップのPod)を越えるたびに inject → extract のペアでコンテキストを橋渡しすることで、複数サービスにまたがる処理を1本のトレースとして可視化できる。DataDog のログ相関は
dd.trace_id フィールドを32桁16進数で埋め込むだけ——シンプルだが見落とすと効果がゼロになる。