在分布式流式传输流水线中,仅使用汇总指标可能难以确定是哪个特定阶段、外部服务调用或工作器 shuffle 导致了延迟峰值。通过在 Dataflow 中使用 OpenTelemetry 分布式跟踪,您可以在单个 Trace ID 下跨流水线转换、网络边界和集成服务端到端地跟踪各个元素。这种可见性有助于您发现性能瓶颈、优化流水线执行费用,以及排查复杂的流式工作负载。
您可以使用以下信息在 Apache Beam 流水线中启用、配置和使用 OpenTelemetry 分布式跟踪。
使用场景
Dataflow 中的分布式跟踪在以下场景中特别有用:
- 流式工作负载:跟踪每个元素的处理时间,并确定会引入延迟或延迟峰值的阶段。
- 用户代码中的外部服务调用跟踪:当流水线在用户代码(例如在
DoFn中)中调用外部数据库、微服务或生成式 AI API 时,测量执行时长。例如,您可以测量对外部数据库(例如 Spanner 或 Bigtable)、微服务(使用 HTTP 或 gRPC)或生成式 AI API(例如 Gemini)的调用所花费的时间。如需了解详情,请参阅 在 DoFn 中写入自定义轨迹和与客户端库集成。 - 成本和性能优化:在实际负载下,显示
PTransform实例中特定操作的持续时间,以识别低效的序列化、缓慢的查询或工作线程争用。 - 复杂的智能体和工作流流水线:深入了解具有多个分支、迭代智能体循环和联接的流水线,以确定哪个转换或分支耗时最长。
限制
Dataflow 中的分布式跟踪具有以下限制:
- 此功能仅适用于流处理流水线,并已针对流处理流水线进行设计和验证。
- 如需使用此功能,您必须使用 Java 版 Beam SDK(版本 2.76.0 或更高版本)。
- 跨阶段上下文传播需要 Streaming Java Runner(之前称为 Runner v1)。可移植 Runner(之前称为 Runner v2)会在 shuffle 边界重置跟踪记录上下文。
- Beam I/O 连接器中的内置 OpenTelemetry 跟踪功能仅支持以下连接器:Pub/Sub、Apache Kafka 和 Spanner 变更数据流。对于 Apache Kafka,只有在直接使用
KafkaIO连接器时才支持跟踪。使用 Kafka 的托管式 I/O 时,不支持此功能。 其他 I/O 连接器(例如BigtableIO)没有内置的连接器跟踪功能;如需跟踪与这些服务相关的操作,请在用户代码(例如在DoFn中)内插桩您的调用。 - 将多个元素合并为一个元素的操作不会传播单个跟踪记录上下文。这包括
GroupByKey、CoGroupByKey和Combine等操作。 - 在设置计时器 (
Timer.set()) 和执行@OnTimer回调方法之间,上下文不会传播。
前提条件
如需在 Dataflow 中使用分布式跟踪,请确保您的环境满足以下要求:
- 您的工作必须满足以下要求:
- 是流处理流水线
- 使用 Streaming Java Runner(之前称为 Runner v1)
- 使用 Apache Beam Java SDK 2.76.0 版或更高版本
如果您将轨迹导出到 Cloud Trace,请在 Google Cloud 项目中启用 Telemetry API (
telemetry.googleapis.com)。启用 API 所需的角色
如需启用 API,您需要拥有
serviceusage.services.enable权限。如果您创建了项目,则可能已经通过 Owner 角色 (roles/owner) 拥有此权限。否则,您可以通过 Service Usage Admin 角色 (roles/serviceusage.serviceUsageAdmin) 获取此权限。 了解如何授予角色。工作器服务账号必须具有 Cloud Trace Agent (
roles/cloudtrace.agent) 角色,该角色包含cloudtrace.traces.patch权限。
核心概念
分布式跟踪依赖于以下 OpenTelemetry 概念:
- 跟踪记录
- 完整交易或工作流在分布式系统中的流转表示形式。一个轨迹由一个或多个跨度组成。
- span
- 轨迹的基本构建块,表示单个工作单元。例如,转换执行、方法调用或 RPC 请求。
- 上下文
- 跨线程和 API 边界传播的状态。这包括跟踪记录 ID 和 Span ID。
- 传播者
- 用于在网络边界之间注入和提取上下文表示的机制。例如,W3C 跟踪记录上下文。
- 出口商
- 一种将收集到的 span 数据发送到跟踪后端(例如 Cloud Trace 或自行托管的 OpenTelemetry 收集器)的组件。
- 采样器
- 一种通过确定记录和导出哪些轨迹来控制轨迹量的机制。例如,1% 的抽样率。
启用 OpenTelemetry 跟踪
如需使用默认设置启用 OpenTelemetry 跟踪,请在提交作业时配置所需的流水线选项和实验。默认设置会以 1% 的抽样率导出到 Cloud Trace。
如需修改默认值,请参阅设置和修改跟踪属性。
在运行作业时配置流水线选项和实验。
非 Flex 模板作业
提交作业时,请传递以下实验标志:
--experiments=enable_otel_defaults,element_metadata_supported,disable_portable_worker
Flex 模板作业
在
gcloud命令中使用--additional-experiments标志来传递实验:gcloud dataflow flex-template run JOB_NAME \ --template-file-gcs-location=gs://BUCKET_NAME/templates/TEMPLATE_NAME.json \ --additional-experiments="enable_otel_defaults,element_metadata_supported,disable_portable_worker" \ --region=REGION替换以下内容:
- JOB_NAME:Dataflow 作业的名称。
- BUCKET_NAME:包含模板的 Cloud Storage 存储桶。
- TEMPLATE_NAME:Flex 模板文件的名称。
- REGION:您要运行作业的 Google Cloud 区域。
这些实验配置了以下行为:
enable_otel_defaults:配置标准 OpenTelemetry 默认值,包括 Cloud Trace 导出器和 1% 的采样率。element_metadata_supported:支持在流水线阶段之间序列化和传播元数据(例如跟踪记录上下文)。disable_portable_worker:强制使用 Dataflow Streaming Java Runner(之前称为 Runner v1)。
如果您使用的是 Java 版 Beam SDK 2.76.0 或 2.77.0,请将所需的 OpenTelemetry 依赖项添加到 build 配置中。
由于版本 2.76.0 和 2.77.0 中存在问题,您必须将以下依赖项添加到
pom.xml文件中,才能导出轨迹。如需了解详情,请参阅管理流水线依赖项。<dependency> <groupId>io.opentelemetry</groupId> <artifactId>opentelemetry-sdk-extension-autoconfigure</artifactId> <scope>compile</scope> </dependency> <dependency> <groupId>io.opentelemetry</groupId> <artifactId>opentelemetry-exporter-otlp</artifactId> <version>1.62.0</version> <scope>compile</scope> </dependency> <dependency> <groupId>org.apache.beam</groupId> <artifactId>beam-sdks-java-extensions-opentelemetry-gcp-auth-extension</artifactId> <scope>compile</scope> </dependency>在流水线代码中,针对受支持的 I/O 连接器显式启用 OpenTelemetry 跟踪:
Pub/Sub
如需在从 Pub/Sub 读取数据或向 Pub/Sub 写入数据时启用跟踪,请执行以下操作:
// Reading from Pub/Sub pipeline.apply("ReadFromPubSub", PubsubIO.readMessages() .withEnableOpenTelemetryTracing() .fromSubscription("projects/PROJECT_ID/subscriptions/SUBSCRIPTION_ID")); // Writing to Pub/Sub pipeline.apply("WriteToPubSub", PubsubIO.writeStrings() .withEnableOpenTelemetryTracing() .to("projects/PROJECT_ID/topics/TOPIC_ID"));
替换以下内容:
- PROJECT_ID:您的 Google Cloud 项目的 ID。
- SUBSCRIPTION_ID:您的 Pub/Sub 订阅的名称。
- TOPIC_ID:您的 Pub/Sub 主题的名称。
上下文传播详情:
- 上游发布者:从外部应用发布消息时,将 W3C 跟踪上下文注入到 Pub/Sub 消息属性中。如需了解详情,请参阅 Pub/Sub OpenTelemetry 跟踪。
- 属性:Beam 需要并填充以下属性:
googclient_traceparent:携带 W3Ctraceparent标头。googclient_tracestate:包含 W3Ctracestate标头。
- 如需了解标头格式,请参阅 W3C 跟踪记录上下文规范。
Apache Kafka
如需在从 Apache Kafka 读取数据或向 Apache Kafka 写入数据时启用跟踪,请执行以下操作:
// Reading from Kafka pipeline.apply("ReadFromKafka", KafkaIO.<String, String>read() .withEnableOpenTelemetryTracing() // ... additional configuration ); // Writing to Kafka pipeline.apply("WriteToKafka", KafkaIO.<String, String>write() .withEnableOpenTelemetryTracing() // ... additional configuration );
上下文传播详情:
KafkaIO使用标准 Kafka 记录标头来传播上下文。- 它会预期并注入符合 W3C 跟踪记录上下文规范的标头:
traceparent:标识传入的跟踪记录上下文。tracestate:提供其他特定于供应商的路由和过滤元数据。
Spanner
如需在从 Spanner 变更数据流读取数据时启用跟踪,请执行以下操作:
pipeline.apply("ReadChangeStream", SpannerIO.readChangeStream() .withSpannerConfig(spannerConfig) .withEnableOpenTelemetryTracing(true) .withChangeStreamName("CHANGE_STREAM_NAME"));
替换以下内容:
- CHANGE_STREAM_NAME:Spanner 变更数据流的名称。
Span 创建详细信息:
- 启用跟踪会为变更数据捕获 (CDC) 读取器读取的每个数据库突变记录启动新的轨迹和 span。
- 如果 Spanner 读取器轮询周期未返回任何新记录,也可能会生成可在 Cloud Trace 界面中过滤掉的 trace span。
设置和修改跟踪属性
指定 enable_otel_defaults 时,Dataflow 会应用以下默认属性:
| 属性 | 默认值 | 说明 |
|---|---|---|
otel.traces.exporter |
otlp |
使用 OpenTelemetry 协议 (OTLP) 导出跟踪记录。 |
otel.exporter.otlp.endpoint |
https://telemetry.googleapis.com |
以 Cloud Trace OTLP 接收器为目标。 |
google.cloud.project |
您的项目 ID | 存储轨迹的 Google Cloud 项目。 |
otel.traces.sampler.arg |
0.01 (1%) |
记录每 100 个轨迹中的 1 个。 |
otel.service.name |
options.getAppName() |
与 span 关联的服务名称,派生自流水线 appName 选项。 |
如需自定义 OpenTelemetry 配置,请将以英文分号分隔的属性列表传递给 --openTelemetryProperties 流水线选项。
以下示例明确配置了所有默认属性:
--openTelemetryProperties="otel.traces.exporter=otlp;otel.exporter.otlp.endpoint=https://telemetry.googleapis.com;google.cloud.project=PROJECT_ID;otel.traces.sampler.arg=0.01;otel.service.name=CUSTOM_SERVICE_NAME;otel.java.global-autoconfigure.enabled=true"
设置自定义服务名称
默认情况下,跟踪服务名称派生自流水线 appName 选项 (--appName=YourPipelineName)。您可以通过设置 otel.service.name 属性来替换此值:
--openTelemetryProperties="otel.service.name=CUSTOM_SERVICE_NAME"
更改抽样率和媒体资源
设置较高的采样率可为您提供更可靠的跟踪信息。不过,较高的采样率(例如 1.0 或 always_on)会大幅增加发送到 Cloud Trace 的跟踪数据量,从而增加您的 Google Cloud 费用。仅在短期调试时使用高采样率,并针对生产工作负载降低采样率。
以下示例明确配置了抽样率。在这种情况下,您必须提供完整的实参列表:
--openTelemetryProperties="otel.traces.exporter=otlp;otel.exporter.otlp.endpoint=https://telemetry.googleapis.com;google.cloud.project=PROJECT_ID;otel.traces.sampler.arg=0.02;otel.service.name=CUSTOM_SERVICE_NAME;otel.java.global-autoconfigure.enabled=true"
导出到第三方后端
如需将跟踪记录发送到外部托管式跟踪后端或自行托管的 OpenTelemetry 收集器,而不是 Cloud Trace,请执行以下操作:
- 省略以下实验:
--experiments=enable_otel_defaults。 使用
--openTelemetryProperties手动配置导出器属性:--openTelemetryProperties="otel.traces.exporter=otlp;otel.exporter.otlp.endpoint=http://COLLECTOR_HOST:4317;otel.traces.sampler=parentbased_always_on"
Google Cloud 身份验证扩展
Beam 运行时包含一个 Google Cloud 身份验证扩展程序,该扩展程序会自动对发送到 Google API(例如 telemetry.googleapis.com)的 OpenTelemetry 协议 (OTLP) 请求进行身份验证。
- 使用默认设置时:系统会自动启用 Google Cloud 身份验证。
- 导出到非 Google 端点时:如果省略
enable_otel_defaults,Beam 会通过设置系统属性google.otel.auth.target.signals=none停用 Google Cloud 身份验证扩展程序。 - 手动配置:您可以将
google.otel.auth.target.signals系统属性设置为none(停用)或trace(仅针对轨迹启用)。
在 DoFn 中写入自定义轨迹
您可以在用户代码(例如 DoFn 内)中创建自定义 span,以衡量特定操作。这包括对外部 API 的调用或计算密集型逻辑。
DoFn 实现示例
import io.opentelemetry.api.trace.Span;
import io.opentelemetry.api.trace.StatusCode;
import io.opentelemetry.api.trace.Tracer;
import io.opentelemetry.context.Scope;
import org.apache.beam.sdk.options.PipelineOptions;
import org.apache.beam.sdk.options.SdkHarnessOptions;
import org.apache.beam.sdk.transforms.DoFn;
public class InvestigateTransactionDoFn extends DoFn<Transaction, RiskAssessment> {
private transient Tracer tracer;
@Setup
public void setup(PipelineOptions options) {
// Obtain the Tracer instance from SdkHarnessOptions
this.tracer = options
.as(SdkHarnessOptions.class)
.getOpenTelemetry()
.getTracer("org.apache.beam.examples.adk.aml");
}
@ProcessElement
public void processElement(@Element Transaction tx, OutputReceiver<RiskAssessment> out) {
// Create and start a custom span as a child of the active context
Span span = tracer.spanBuilder("InvestigateTransactionDoFn:Analyze")
.setAttribute("transaction.id", tx.getId())
.setAttribute("transaction.amount", tx.getAmount())
.startSpan();
// Make the span current in the execution thread
try (Scope scope = span.makeCurrent()) {
// Perform your logic (such as calling an external service)
RiskAssessment assessment = analyzeTransactionWithLLM(tx);
out.output(assessment);
} catch (Exception e) {
// Record exception details and set error status on the span
span.recordException(e);
span.setStatus(StatusCode.ERROR, e.getMessage());
throw e;
} finally {
// Always end the span
span.end();
}
}
private RiskAssessment analyzeTransactionWithLLM(Transaction tx) {
// Application analysis logic
return new RiskAssessment();
}
}
与客户端库集成
当 DoFn 调用也支持 OpenTelemetry 的外部库(例如 Google Cloud 客户端库)时,请配置客户端以使用流水线的 OpenTelemetry 实例,以便子 span 直接链接到流水线跟踪记录。
示例 Spanner Java 客户端
如需在 DoFn 内执行手动 Spanner 查询并将这些 span 关联到活跃的流水线轨迹,请使用流水线的 OpenTelemetry 实例初始化 SpannerOptions:
import com.google.cloud.spanner.DatabaseClient;
import com.google.cloud.spanner.DatabaseId;
import com.google.cloud.spanner.ResultSet;
import com.google.cloud.spanner.Spanner;
import com.google.cloud.spanner.SpannerOptions;
import com.google.cloud.spanner.Statement;
import org.apache.beam.sdk.options.PipelineOptions;
import org.apache.beam.sdk.options.SdkHarnessOptions;
import org.apache.beam.sdk.transforms.DoFn;
public class SpannerEnrichmentDoFn extends DoFn<String, String> {
private transient Spanner spanner;
private transient DatabaseClient dbClient;
private final String instanceId;
private final String databaseId;
public SpannerEnrichmentDoFn(String instanceId, String databaseId) {
this.instanceId = instanceId;
this.databaseId = databaseId;
}
@Setup
public void setup(PipelineOptions pipelineOptions) {
// Build SpannerOptions leveraging the pipeline's active OpenTelemetry instance
SpannerOptions options =
SpannerOptions.newBuilder()
.setOpenTelemetry(pipelineOptions
.as(SdkHarnessOptions.class)
.getOpenTelemetry())
.setEnableEndToEndTracing(true)
.setEnableExtendedTracing(true)
.build();
this.spanner = options.getService();
DatabaseId db = DatabaseId.of(options.getProjectId(), instanceId, databaseId);
this.dbClient = spanner.getDatabaseClient(db);
}
@ProcessElement
public void processElement(@Element String inputId, OutputReceiver<String> out) {
Statement statement = Statement.newBuilder(
"SELECT Details FROM AccountTable WHERE AccountId = @id")
.bind("id").to(inputId)
.build();
// Spanner queries automatically generate child spans linked to the active trace
try (ResultSet rs = dbClient.singleUse().executeQuery(statement)) {
while (rs.next()) {
out.output(rs.getString("Details"));
}
}
}
@Teardown
public void tearDown() {
if (spanner != null) {
spanner.close();
}
}
}
跨阶段的上下文传播
在分布式数据流水线中,跨工作器边界和 shuffle 操作传递跟踪上下文对于实现端到端的可观测性至关重要。
当您将 Redistribute.arbitrarily() 插入流水线时,活跃的跟踪上下文会序列化为元素元数据,通过网络传输,并在接收工作器虚拟机上反序列化。这样可确保在各个阶段保持跟踪的连续性。
以下示例展示了跟踪记录上下文如何在流水线阶段之间传播:
- 提取上下文:
KafkaIO.read()读取传入的记录并从消息标头中提取跟踪记录上下文。 - 处理并创建 span:
DoFn A处理有效跟踪 span 中的元素。 - 在工作器之间传播:
Redistribute.arbitrarily()将有效跟踪记录上下文序列化为元素元数据,并在网络中对元素进行 shuffle。 - 继续跟踪记录:
DoFn B在工作器虚拟机上接收元素,反序列化跟踪记录上下文,并在同一跟踪记录 ID 下继续处理。
在 Cloud Trace 中查看轨迹
启用跟踪后,在 Google Cloud 控制台中查看流水线跟踪记录:
在 Google Cloud 控制台中,前往 Trace 探索器页面。
在过滤条件框中,按
Service: SERVICE_NAME进行过滤。 此 SERVICE_NAME 与流水线appName或配置的otel.service.name相匹配。或者,搜索特定的轨迹 ID。
从结果列表中选择一条轨迹,即可查看其详细时间轴。
时间轴图表会显示:
- Runner span:由 Beam 运行时生成的 span。
例如,
PubSubIO.Read和ProcessElement。 - 自定义跨度:在
DoFn实现中创建的跨度。例如InvestigateTransactionDoFn:Analyze。 客户端库 span:由集成客户端库生成的子 span。例如,Spanner RPC 查询。

- Runner span:由 Beam 运行时生成的 span。
例如,
后续步骤
- 如需详细了解如何探索和分析分布式跟踪记录,请参阅 Cloud Trace 文档。
- 使用 Pub/Sub OpenTelemetry 跟踪跟踪从提取到处理的消息。
- 使用 Dataflow 监控界面监控流水线吞吐量、CPU 利用率和执行详情。
- 使用 Dataflow 流水线日志查看操作日志并诊断错误。
- 通过排查作业缓慢或卡住的问题和检测并解决流水线瓶颈来调查并解决性能问题。