以编程方式使用文档

Temporal 是一个持久的执行平台,使开发者能够构建弹性的分布式应用程序。本指南将向您展示如何使用 OpenTelemetry 在 LangSmith 中追踪 Temporal 工作流和活动。

LangSmith 支持 OpenTelemetry (OTEL) 追踪数据接收,可与 Temporal 的原生 OpenTelemetry 拦截器无缝集成。这使您能够对工作流执行、活动及其中的任何 LLM 调用进行完整的分布式追踪。

前提条件

  • - A LangSmith 账户 和 API 密钥
  • - 正在运行的 Temporal 服务器(本地或云端)
  • - 您所用语言的 OpenTelemetry SDK

环境变量

为所有实现设置以下环境变量:

变量必填描述
LANGSMITH_API_KEY您在设置中的 LangSmith API 密钥。
LANGSMITH_PROJECT项目名称(默认为 "default").

设置追踪

Go

Go 使用 langsmith-go SDK 和 Temporal 的 OpenTelemetry 拦截器来自动追踪工作流和活动。

Install

安装 LangSmith Go SDK、Temporal SDK 和 OpenTelemetry 拦截器:

go get github.com/langchain-ai/langsmith-go@v0.1.0-alpha.7
go get go.temporal.io/sdk
go get go.temporal.io/sdk/contrib/opentelemetry

Initialize tracer

初始化 LangSmith 追踪器,创建 Temporal 的 OpenTelemetry 拦截器,并将其注册到 Temporal 客户端和工作线程中:

package main

	"context"
	"log"

	"github.com/langchain-ai/langsmith-go"
	"go.temporal.io/sdk/client"
	"go.temporal.io/sdk/contrib/opentelemetry"
	"go.temporal.io/sdk/interceptor"
	"go.temporal.io/sdk/worker"
)

func main() {
	ctx := context.Background()

	// Initialize LangSmith tracer (reads LANGSMITH_API_KEY and LANGSMITH_PROJECT)
	ls, err := langsmith.NewTracer(
		langsmith.WithServiceName("temporal-worker"),
	)
	if err != nil {
		log.Fatal("Failed to initialize LangSmith tracer:", err)
	}
	defer ls.Shutdown(ctx)

	// Create Temporal tracing interceptor
	tracer := ls.Tracer("temporal-app")
	tracingInterceptor, err := opentelemetry.NewTracingInterceptor(
		opentelemetry.TracerOptions{Tracer: tracer},
	)
	if err != nil {
		log.Fatal("Failed to create tracing interceptor:", err)
	}

	// Create Temporal client with tracing
	c, err := client.Dial(client.Options{
		Interceptors: []interceptor.ClientInterceptor{tracingInterceptor},
	})
	if err != nil {
		log.Fatal("Failed to create Temporal client:", err)
	}
	defer c.Close()

	// Create worker with tracing (uses same client)
	w := worker.New(c, "my-task-queue", worker.Options{})
	w.RegisterWorkflow(MyWorkflow)
	w.RegisterActivity(MyActivity)

	// Start worker
	if err := w.Run(worker.InterruptCh()); err != nil {
		log.Fatal("Worker failed:", err)
	}
}

Define workflow and activity

定义一个执行活动的工作流。该活动演示如何为 LangSmith 可视性添加自定义跨度属性:

package main

	"context"
	"fmt"
	"time"

	"go.opentelemetry.io/otel/attribute"
	"go.opentelemetry.io/otel/trace"
	"go.temporal.io/sdk/activity"
	"go.temporal.io/sdk/workflow"
)

// MyWorkflow executes an activity
func MyWorkflow(ctx workflow.Context, input string) (string, error) {
	ao := workflow.ActivityOptions{
		StartToCloseTimeout: 10 * time.Second,
	}
	ctx = workflow.WithActivityOptions(ctx, ao)

	var result string
	err := workflow.ExecuteActivity(ctx, MyActivity, input).Get(ctx, &result)
	return result, err
}

// MyActivity processes input with custom span attributes
func MyActivity(ctx context.Context, input string) (string, error) {
	logger := activity.GetLogger(ctx)
	logger.Info("Processing", "input", input)

	// Get the span created by Temporal's interceptor
	span := trace.SpanFromContext(ctx)

	// Add Gen AI attributes for LangSmith visibility
	span.SetAttributes(
		attribute.String("gen_ai.prompt", input),
		attribute.String("gen_ai.operation.name", "chat"),
	)

	result := fmt.Sprintf("Processed: %s", input)

	// Set completion attribute
	span.SetAttributes(
		attribute.String("gen_ai.completion", result),
	)

	return result, nil
}

Execute workflow

在单独的客户端应用程序中,初始化追踪器并执行工作流:

// In a separate function or client application
func executeWorkflow() {
	ctx := context.Background()

	// Initialize tracer for client
	ls, err := langsmith.NewTracer(
		langsmith.WithServiceName("temporal-client"),
	)
	if err != nil {
		log.Fatal(err)
	}
	defer ls.Shutdown(ctx)

	// Create client with tracing
	tracer := ls.Tracer("temporal-app")
	tracingInterceptor, err := opentelemetry.NewTracingInterceptor(
		opentelemetry.TracerOptions{Tracer: tracer},
	)
	if err != nil {
		log.Fatal(err)
	}

	c, err := client.Dial(client.Options{
		Interceptors: []interceptor.ClientInterceptor{tracingInterceptor},
	})
	if err != nil {
		log.Fatal(err)
	}
	defer c.Close()

	// Execute workflow
	workflowOptions := client.StartWorkflowOptions{
		ID:        "my-workflow-1",
		TaskQueue: "my-task-queue",
	}

	we, err := c.ExecuteWorkflow(ctx, workflowOptions, MyWorkflow, "Hello World")
	if err != nil {
		log.Fatal(err)
	}

	var result string
	if err := we.Get(ctx, &result); err != nil {
		log.Fatal(err)
	}

	log.Printf("Workflow result: %s", result)
}

Python

Python 使用 temporalio SDK 和 OpenTelemetry 拦截器,通过 OTLP 将追踪数据导出到 LangSmith。

Install

安装 Temporal SDK、LangSmith SDK 和 OpenTelemetry 包:

pip install temporalio
pip install langsmith
pip install opentelemetry-sdk
pip install opentelemetry-exporter-otlp-proto-http

Initialize tracer

创建一个 OpenTelemetry TracerProvider 并配置 OTLP 导出器以将追踪数据发送到 LangSmith:

from datetime import timedelta

from opentelemetry import trace
from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter
from opentelemetry.sdk.resources import Resource, SERVICE_NAME
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor

from temporalio import activity, workflow
from temporalio.client import Client
from temporalio.contrib.opentelemetry import TracingInterceptor
from temporalio.worker import Worker


def init_tracer_provider() -> TracerProvider:
    """Initialize OpenTelemetry with LangSmith exporter."""

    # Create OTLP exporter for LangSmith
    exporter = OTLPSpanExporter(
        endpoint="https://api.smith.langchain.com/otel/v1/traces",
        headers={
            "x-api-key": os.environ.get("LANGSMITH_API_KEY", ""),
            "Langsmith-Project": os.environ.get("LANGSMITH_PROJECT", "default"),
        },
    )

    # Create TracerProvider with resource attributes
    resource = Resource.create({
        SERVICE_NAME: "temporal-worker",
    })

    provider = TracerProvider(resource=resource)
    provider.add_span_processor(BatchSpanProcessor(exporter))

    # Set as global provider
    trace.set_tracer_provider(provider)

    return provider

Define workflow and activity

定义一个工作流类和一个活动函数。该活动演示如何为 LangSmith 可视性添加自定义跨度属性:

@activity.defn
async def process_activity(input: str) -> str:
    """Activity that processes input with custom span attributes."""
    activity.logger.info(f"Processing: {input}")

    # Get current span and add Gen AI attributes
    span = trace.get_current_span()
    span.set_attribute("gen_ai.prompt", input)
    span.set_attribute("gen_ai.operation.name", "chat")

    result = f"Processed: {input}"

    span.set_attribute("gen_ai.completion", result)

    return result


@workflow.defn
class MyWorkflow:
    @workflow.run
    async def run(self, input: str) -> str:
        return await workflow.execute_activity(
            process_activity,
            input,
            start_to_close_timeout=timedelta(seconds=10),
        )

Run worker

使用 TracingInterceptor 创建 Temporal 客户端并启动工作线程:

async def main():
    # Initialize tracing
    provider = init_tracer_provider()

    try:
        # Create Temporal client with tracing interceptor
        client = await Client.connect(
            "localhost:7233",
            interceptors=[TracingInterceptor()],
        )

        # Run worker
        worker = Worker(
            client,
            task_queue="my-task-queue",
            workflows=[MyWorkflow],
            activities=[process_activity],
        )

        print("Starting worker...")
        await worker.run()

    finally:
        # Shutdown tracer provider to flush traces
        provider.shutdown()


if __name__ == "__main__":
    asyncio.run(main())

Execute workflow

在单独的脚本中,使用追踪拦截器连接到 Temporal 并执行工作流:

from temporalio.client import Client
from temporalio.contrib.opentelemetry import TracingInterceptor

# Import the same tracer setup
from worker import init_tracer_provider


async def main():
    provider = init_tracer_provider()

    try:
        client = await Client.connect(
            "localhost:7233",
            interceptors=[TracingInterceptor()],
        )

        # Execute workflow
        result = await client.execute_workflow(
            MyWorkflow.run,
            "Hello World",
            id="my-workflow-1",
            task_queue="my-task-queue",
        )
        print(f"Workflow result: {result}")

    finally:
        provider.shutdown()


if __name__ == "__main__":
    asyncio.run(main())

TypeScript / JavaScript

TypeScript 使用 @temporalio/sdk 和 OpenTelemetry 拦截器将追踪数据发送到 LangSmith。

Install

安装 Temporal SDK、OpenTelemetry 拦截器和追踪包:

npm install @temporalio/client @temporalio/worker @temporalio/activity @temporalio/workflow
npm install @temporalio/interceptors-opentelemetry
npm install @opentelemetry/sdk-node @opentelemetry/sdk-trace-node
npm install @opentelemetry/exporter-trace-otlp-http
npm install @opentelemetry/resources @opentelemetry/semantic-conventions

Initialize tracer

创建一个 NodeTracerProvider 并配置 OTLP 导出器以将追踪数据发送到 LangSmith:

  // Create OTLP exporter for LangSmith
  const exporter = new OTLPTraceExporter({
    url: 'https://api.smith.langchain.com/otel/v1/traces',
    headers: {
      'x-api-key': process.env.LANGSMITH_API_KEY || '',
      'Langsmith-Project': process.env.LANGSMITH_PROJECT || 'default',
    },
  });

  // Create TracerProvider
  const provider = new NodeTracerProvider({
    resource: new Resource({
      [ATTR_SERVICE_NAME]: 'temporal-worker',
    }),
  });

  provider.addSpanProcessor(new BatchSpanProcessor(exporter));
  provider.register();

  return provider;
}

Define workflow

定义一个使用超时配置代理活动的工作流:

const { processActivity } = proxyActivities<typeof activities>({
  startToCloseTimeout: '10 seconds',
});

  return await processActivity(input);
}

Define activity

定义一个演示如何为 LangSmith 可视性添加自定义跨度属性的活动:

  log.info('Processing', { input });

  // Get current span and add Gen AI attributes
  const span = trace.getActiveSpan();
  span?.setAttribute('gen_ai.prompt', input);
  span?.setAttribute('gen_ai.operation.name', 'chat');

  const result = `Processed: ${input}`;

  span?.setAttribute('gen_ai.completion', result);

  return result;
}

Run worker

创建一个带有 OpenTelemetry 拦截器的 worker,用于活动和用于工作流跨度的工作线程导出器:

  makeWorkflowExporter,
  OpenTelemetryActivityInboundInterceptor,
} from '@temporalio/interceptors-opentelemetry';




async function run() {
  const provider = initTracerProvider();

  try {
    const connection = await NativeConnection.connect({
      address: 'localhost:7233',
    });

    const worker = await Worker.create({
      connection,
      namespace: 'default',
      taskQueue: 'my-task-queue',
      workflowsPath: require.resolve('./workflows'),
      activities,
      sinks: {
        exporter: makeWorkflowExporter(
          trace.getTracer('temporal-app'),
          new Resource({ [ATTR_SERVICE_NAME]: 'temporal-worker' })
        ),
      },
      interceptors: {
        activity: [() => ({ inbound: new OpenTelemetryActivityInboundInterceptor() })],
      },
    });

    console.log('Starting worker...');
    await worker.run();
  } finally {
    await provider.shutdown();
  }
}

run().catch((err) => {
  console.error(err);
  process.exit(1);
});

Execute workflow

在单独的客户端文件中,连接到 Temporal 并执行工作流:

async function run() {
  // Initialize tracing
  const provider = initTracerProvider();

  try {
    const connection = await Connection.connect({ address: 'localhost:7233' });
    const client = new Client({ connection });

    const result = await client.workflow.execute('myWorkflow', {
      taskQueue: 'my-task-queue',
      workflowId: 'my-workflow-1',
      args: ['Hello World'],
    });

    console.log('Workflow result:', result);
  } finally {
    await provider.shutdown();
  }
}

run().catch(console.error);

在 LangSmith 中查看追踪

配置完成后,追踪将出现在您的 LangSmith 项目中:

  1. 导航到您的 LangSmith 实例。
  2. 选择您的项目。
  3. 在 **追踪** tab.
  4. 点击单个追踪以查看完整的 span 层级。

配置选项

设置自定义服务名称

设置自定义服务名称以区分不同的 Temporal workers 或服务:

ls, err := langsmith.NewTracer(
    langsmith.WithServiceName("my-temporal-worker"),
)
resource = Resource.create({
    SERVICE_NAME: "my-temporal-worker",
})
const provider = new NodeTracerProvider({
  resource: new Resource({
    [ATTR_SERVICE_NAME]: 'my-temporal-worker',
  }),
});

添加自定义 span 属性

添加自定义属性以丰富您的追踪:

span := trace.SpanFromContext(ctx)
span.SetAttributes(
    attribute.String("user.id", userID),
    attribute.String("workflow.version", "v2"),
)
from opentelemetry import trace

span = trace.get_current_span()
span.set_attribute("user.id", user_id)
span.set_attribute("workflow.version", "v2")
const span = trace.getActiveSpan();
span?.setAttribute('user.id', userId);
span?.setAttribute('workflow.version', 'v2');

配置采样

对于高容量工作流,配置采样以减少追踪量:

// Note: langsmith.NewTracer() uses default sampling
// For custom sampling, use the TracerProvider directly
tp := sdktrace.NewTracerProvider(
    sdktrace.WithBatcher(exporter),
    sdktrace.WithSampler(sdktrace.TraceIDRatioBased(0.1)), // 10% sampling
)
from opentelemetry.sdk.trace.sampling import TraceIdRatioBased

provider = TracerProvider(
    resource=resource,
    sampler=TraceIdRatioBased(0.1),  # 10% sampling
)
const provider = new NodeTracerProvider({
  resource: resource,
  sampler: new TraceIdRatioBasedSampler(0.1), // 10% sampling
});

故障排除

追踪未显示

  1. **验证 API 密钥**:确保 LANGSMITH_API_KEY 设置正确
  2. **检查端点**:确认您正在使用 https://api.smith.langchain.com/otel/v1/traces
  3. **关闭时刷新**:调用 provider.shutdown() 在应用程序退出前刷新待处理的 spans
  4. **检查项目**:验证追踪是否发送到正确的项目(默认是 "default")

缺少 activity spans

确保追踪拦截器同时在客户端和 worker 上配置: - **客户端**:需要拦截器来启动工作流 - **Worker**:需要拦截器来执行 activities

上下文传播问题

验证传播器配置是否正确: - **Go**: langsmith.NewTracer() 自动配置传播器 - **Python/TypeScript**:确保 OpenTelemetry SDK 已正确初始化并配置了 trace 传播器

Worker 关闭挂起

如果追踪未刷新,请确保使用适当的超时调用 shutdown 方法:

defer ls.Shutdown(context.Background())
finally:
    provider.shutdown()
finally {
  await provider.shutdown();
}

后续步骤

其他资源