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 项目中:
- 导航到您的 LangSmith 实例。
- 选择您的项目。
- 在 **追踪** tab.
- 点击单个追踪以查看完整的 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
});
故障排除
追踪未显示
- **验证 API 密钥**:确保
LANGSMITH_API_KEY设置正确 - **检查端点**:确认您正在使用
https://api.smith.langchain.com/otel/v1/traces - **关闭时刷新**:调用
provider.shutdown()在应用程序退出前刷新待处理的 spans - **检查项目**:验证追踪是否发送到正确的项目(默认是
"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();
}