You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用Java SDK及PipelineResult统计Google DataFlow各步骤执行时间

How to Measure Execution Time of Each Step in Apache Beam DataFlow Job (Java SDK, Beam 2.3.0)

Problem Statement

I'm running a DataFlow job on Google Cloud Platform using Apache Beam 2.3.0. The job has 5 repeated steps (looping 5 times with two transforms each), and I want to track how long each individual step takes to complete using the Java SDK. Here's my current code:

Pipeline pipeline = Pipeline.create(options); 
for(int i=0; i<5; i++) { 
    PCollection<String> csv = pipeline.apply(transform1); 
    csv.apply(transform2); 
} 
pipeline.run().waitUntilFinish();

How can I use PipelineResult to measure the execution time of each step?


Solution

To track execution time for each step, you'll need to give each transform a unique name (to distinguish repeated steps) and leverage Beam's built-in metrics system or DataFlow's job metrics exposed via PipelineResult. Here's how to do it step by step:

1. Assign Unique Names to Your Transforms

First, update your loop to give each transform1 and transform2 instance a unique name tied to the loop index. This ensures you can tell apart metrics from each iteration:

Pipeline pipeline = Pipeline.create(options); 
for(int i=0; i<5; i++) { 
    String stepPrefix = "Step-" + (i + 1); // Name each step clearly
    PCollection<String> csv = pipeline.apply(stepPrefix + "-Transform1", transform1); 
    csv.apply(stepPrefix + "-Transform2", transform2); 
} 
PipelineResult result = pipeline.run();
result.waitUntilFinish();

2. Extract Execution Time via PipelineResult Metrics

Once the job finishes, you can pull metrics directly from PipelineResult. DataFlow automatically collects execution time metrics for each transform, which you can filter and display:

// Retrieve all metrics from the completed job
MetricResults metricResults = result.metrics();

// Query and filter for execution time metrics
MetricQueryResults queryResults = metricResults.queryMetrics(
    MetricsFilter.builder()
        .addNameFilter(MetricsFilter.NameFilter.startsWith("execution_time"))
        .build()
);

// Iterate through results and print transform execution times
for (MetricResult<?> metric : queryResults.allMetrics()) {
    MetricName metricName = metric.getName();
    String transformFullName = metricName.getNamespace();
    // Convert milliseconds to seconds for readability
    double executionTimeSeconds = ((Double) metric.getAttempted()) / 1000;
    
    System.out.printf("Transform %s completed in %.2f seconds%n", 
                      transformFullName, executionTimeSeconds);
}

3. (Optional) Custom Timer Metrics for Fine-Grained Control

If you need more precise control over what's timed (e.g., only the processing logic inside a DoFn, not the entire transform setup/teardown), you can add custom timer metrics directly in your DoFns:

public class CustomTransformFn extends DoFn<String, String> {
    private final Timer executionTimer;

    // Pass a unique step identifier to distinguish metrics across loop iterations
    public CustomTransformFn(String stepId) {
        this.executionTimer = Metrics.timer(stepId, "processing_time");
    }

    @ProcessElement
    public void processElement(ProcessContext c) {
        // Start the timer before processing
        Timer.Context timerCtx = executionTimer.start();
        try {
            // Your actual processing logic here
            String input = c.element();
            // ... process input ...
            c.output(input);
        } finally {
            // Ensure the timer stops even if an exception occurs
            timerCtx.stop();
        }
    }
}

Then update your loop to use this custom DoFn with unique step IDs:

for(int i=0; i<5; i++) { 
    String stepId = "Step-" + (i + 1);
    PCollection<String> csv = pipeline.apply(
        stepId + "-CustomTransform1", 
        ParDo.of(new CustomTransformFn(stepId + "-Transform1"))
    );
    // Repeat for transform2 with its own custom DoFn
}

You can retrieve these custom metrics using the same MetricResults code from step 2—just adjust the filter to match your custom metric name (processing_time in this example).

Notes

  • For Beam 2.3.0, ensure you're using the correct MetricResults API (some method names may differ from newer versions).
  • If you're running on DataFlow, you can also view these metrics in the GCP DataFlow Console under the job's "Metrics" tab, but the above code lets you access them directly in your Java application.

内容的提问来源于stack exchange,提问作者Achaius

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.20 10:37:10