shunping commented on code in PR #40142:
URL: https://github.com/apache/beam/pull/40142#discussion_r4073867044
##########
sdks/java/extensions/google-cloud-platform-core/src/main/java/org/apache/beam/sdk/extensions/gcp/util/GcsUtilV1.java:
##########
@@ -615,11 +639,25 @@ SeekableByteChannel open(GcsPath path,
GoogleCloudStorageReadOptions readOptions
ServiceCallMetric serviceCallMetric =
new ServiceCallMetric(MonitoringInfoConstants.Urns.API_REQUEST_COUNT,
baseLabels);
try {
+ GoogleCloudStorage gcpStorage = this.googleCloudStorage;
+ MetricsContainer container = null;
+ if (gcsCountersOptions.getPerformanceMetricsEnabled()) {
+ container = MetricsEnvironment.getCurrentContainer();
+ if (container != null) {
+ HttpRequestInitializer scopedInitializer =
+ Transport.withMetricsContainer(this.httpRequestInitializer,
container, false);
+ gcpStorage =
Review Comment:
`MetricsEnvironment.getCurrentContainer()` is backed by a `ThreadLocal`,
which the runner only sets on the DoFn thread for the active step.
Calling `MetricsEnvironment.getCurrentContainer()` inside `Transport` when
HTTP requests execute doesn't work because:
- gcsio executes HTTP requests on background thread pools, where Beam's
`ThreadLocal` is not propagated and `getCurrentContainer()` returns `null`. By
capturing container in `open()`/`create()` on the DoFn thread, we can make sure
background threads can increment the counters directly without relying on the
`ThreadLocal` and avoid metric loss.
- `GoogleCloudStorageImpl` fixes its `HttpRequestInitializer` at
construction time. Since each step has its own `MetricsContainer`, caching
`GoogleCloudStorage` per `MetricsContainer` ensures we create at most one
client per step rather than one per file.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]