diff --git a/java-bigquerystorage/google-cloud-bigquerystorage/src/main/java/com/google/cloud/bigquery/storage/v1/ConnectionWorker.java b/java-bigquerystorage/google-cloud-bigquerystorage/src/main/java/com/google/cloud/bigquery/storage/v1/ConnectionWorker.java index 75aba8b3644b..113c4e5b2b8a 100644 --- a/java-bigquerystorage/google-cloud-bigquerystorage/src/main/java/com/google/cloud/bigquery/storage/v1/ConnectionWorker.java +++ b/java-bigquerystorage/google-cloud-bigquerystorage/src/main/java/com/google/cloud/bigquery/storage/v1/ConnectionWorker.java @@ -1156,6 +1156,8 @@ private void appendLoop() { boolean firstRequestForTableOrSchemaSwitch = true; while (!waitingQueueDrained()) { + long queuedRequestsCount = 0; + ClientStats.WindowStats windowStats = null; this.lock.lock(); try { hasMessageInWaitingQueue.await(100, TimeUnit.MILLISECONDS); @@ -1203,6 +1205,11 @@ private void appendLoop() { localQueue.addLast(requestWrapper); healthCheckMetrics.updateRequestsSent(requestWrapper.messageSize); } + if (!localQueue.isEmpty()) { + queuedRequestsCount = + inflightRequestQueue.size(); // waitingRequestQueue is empty so can be omitted here + windowStats = createWindowStats(); + } } catch (InterruptedException e) { log.warning( "Interrupted while waiting for message. Stream: " @@ -1289,6 +1296,13 @@ private void appendLoop() { } firstRequestForTableOrSchemaSwitch = false; + ClientStats clientStats = + createClientStats(wrapper.requestSendTimeStamp, queuedRequestsCount, windowStats); + windowStats = null; + if (!clientStats.equals(ClientStats.getDefaultInstance())) { + originalRequestBuilder.setClientStats(clientStats); + } + requestProfilerHook.startOperation( RequestProfiler.OperationName.RESPONSE_LATENCY, requestUniqueId); @@ -1307,6 +1321,69 @@ private void appendLoop() { cleanupConnectionAndRequests(/* avoidBlocking= */ false); } + private static ClientStats createClientStats( + @Nullable Instant requestSendTimeStamp, + long queuedRequestsCount, + @Nullable ClientStats.WindowStats windowStats) { + ClientStats.Builder clientStatsBuilder = ClientStats.newBuilder(); + ClientStats.RequestStats requestStats = + createRequestStats(requestSendTimeStamp, queuedRequestsCount); + if (!requestStats.equals(ClientStats.RequestStats.getDefaultInstance())) { + clientStatsBuilder.setRequestStats(requestStats); + } + if (windowStats != null && !windowStats.equals(ClientStats.WindowStats.getDefaultInstance())) { + clientStatsBuilder.setWindowStats(windowStats); + } + return clientStatsBuilder.build(); + } + + private static ClientStats.RequestStats createRequestStats( + @Nullable Instant requestSendTimeStamp, long queuedRequestsCount) { + ClientStats.RequestStats.Builder requestStatsBuilder = ClientStats.RequestStats.newBuilder(); + if (requestSendTimeStamp != null) { + requestStatsBuilder.setSendTimeMillis(requestSendTimeStamp.toEpochMilli()); + } + if (queuedRequestsCount > 0) { + requestStatsBuilder.setQueuedRequestsCount(queuedRequestsCount); + } + return requestStatsBuilder.build(); + } + + @GuardedBy("lock") + private ClientStats.WindowStats createWindowStats() { + ClientStats.WindowStats.Builder windowStatsBuilder = ClientStats.WindowStats.newBuilder(); + if (healthCheckMetrics.windowedMilliLatencyMax > 0) { + windowStatsBuilder.setMaxResponseLatencyMillis(healthCheckMetrics.windowedMilliLatencyMax); + } + long avgResponseLatencyMillis = + healthCheckMetrics.windowedResponsesAcked > 0 + ? healthCheckMetrics.windowedMilliLatencySum / healthCheckMetrics.windowedResponsesAcked + : 0; + if (avgResponseLatencyMillis > 0) { + windowStatsBuilder.setAvgResponseLatencyMillis(avgResponseLatencyMillis); + } + if (healthCheckMetrics.windowedMilliResponseWaitTimeMax > 0) { + windowStatsBuilder.setLongestWaitNoResponseMillis( + healthCheckMetrics.windowedMilliResponseWaitTimeMax); + } + if (healthCheckMetrics.windowedRequestsSent > 0) { + windowStatsBuilder.setRequestsSentCount(healthCheckMetrics.windowedRequestsSent); + } + if (healthCheckMetrics.windowedResponsesAcked > 0) { + windowStatsBuilder.setResponsesReceivedCount(healthCheckMetrics.windowedResponsesAcked); + } + if (healthCheckMetrics.windowedRequestsSentBytes > 0) { + windowStatsBuilder.setBytesSentCount(healthCheckMetrics.windowedRequestsSentBytes); + } + long windowStartTimeEpochMillis = healthCheckMetrics.healthCheckTimeStamp.toEpochMilli(); + windowStatsBuilder.setWindowStartTimeEpochMillis(windowStartTimeEpochMillis); + long windowMillis = Instant.now().toEpochMilli() - windowStartTimeEpochMillis; + if (windowMillis >= 0) { + windowStatsBuilder.setWindowMillis(windowMillis); + } + return windowStatsBuilder.build(); + } + @Nullable private AppendRowsSchema getSchema(AppendRowsRequest request) { if (request.hasProtoRows()) { diff --git a/java-bigquerystorage/google-cloud-bigquerystorage/src/test/java/com/google/cloud/bigquery/storage/v1/ConnectionWorkerTest.java b/java-bigquerystorage/google-cloud-bigquerystorage/src/test/java/com/google/cloud/bigquery/storage/v1/ConnectionWorkerTest.java index 6e4ee2642a6a..8d36b9e5e729 100644 --- a/java-bigquerystorage/google-cloud-bigquerystorage/src/test/java/com/google/cloud/bigquery/storage/v1/ConnectionWorkerTest.java +++ b/java-bigquerystorage/google-cloud-bigquerystorage/src/test/java/com/google/cloud/bigquery/storage/v1/ConnectionWorkerTest.java @@ -1612,6 +1612,17 @@ void testInflightRetryCountHealthMetricExactlyOnce() throws Exception { assertTrue(healthCheckFields.responseCodes.containsKey(Status.Code.OK.value())); assertEquals(3, healthCheckFields.responseCodes.get(Status.Code.OK.value())); assertEquals("projects/p1/datasets/d1/tables/t1/streams/s1", healthCheckFields.streamName); + + AppendRowsRequest secondRequest = testBigQueryWrite.getAppendRequests().get(1); + assertThat(secondRequest.getClientStats().getRequestStats().hasQueuedRequestsCount()).isTrue(); + assertThat(secondRequest.getClientStats().getRequestStats().getQueuedRequestsCount()) + .isEqualTo(2); + // When both requests are retried together in localQueue, only the first request in localQueue + // should have window_stats populated. + AppendRowsRequest retryFirstInLocalQueue = testBigQueryWrite.getAppendRequests().get(2); + AppendRowsRequest retrySecondInLocalQueue = testBigQueryWrite.getAppendRequests().get(3); + assertThat(retryFirstInLocalQueue.getClientStats().hasWindowStats()).isTrue(); + assertThat(retrySecondInLocalQueue.getClientStats().hasWindowStats()).isFalse(); } @Test @@ -1722,6 +1733,114 @@ void testEarliestSendTime_retryScenario() throws Exception { // Once exhausted, state should be null assertThat(connectionWorker.getEarliestSendTime()).isNull(); + + AppendRowsRequest initialRequest = testBigQueryWrite.getAppendRequests().get(0); + AppendRowsRequest retryRequest = testBigQueryWrite.getAppendRequests().get(1); + assertThat(initialRequest.getClientStats().getRequestStats().hasSendTimeMillis()).isTrue(); + assertThat(retryRequest.getClientStats().getRequestStats().hasSendTimeMillis()).isTrue(); + assertThat(retryRequest.getClientStats().getRequestStats().getSendTimeMillis()) + .isAtLeast(initialRequest.getClientStats().getRequestStats().getSendTimeMillis()); + assertThat(initialRequest.getClientStats().getRequestStats().hasQueuedRequestsCount()) + .isTrue(); + assertThat(initialRequest.getClientStats().getRequestStats().getQueuedRequestsCount()) + .isEqualTo(1); + assertThat(retryRequest.getClientStats().getRequestStats().hasQueuedRequestsCount()).isTrue(); + assertThat(retryRequest.getClientStats().getRequestStats().getQueuedRequestsCount()) + .isEqualTo(1); + } + } + + @Test + void testWindowStats() throws Exception { + long beforeCreationMillis = Instant.now().toEpochMilli(); + ProtoSchema schema1 = createProtoSchema("foo"); + try (StreamWriter sw1 = + StreamWriter.newBuilder(TEST_STREAM_1, client) + .setLocation("us") + .setWriterSchema(schema1) + .build(); + ConnectionWorker connectionWorker = createConnectionWorker()) { + long afterCreationMillis = Instant.now().toEpochMilli(); + int msecResponseDelay = 300; + testBigQueryWrite.setResponseSleep(Duration.ofMillis(msecResponseDelay)); + testBigQueryWrite.addResponse(createAppendResponse(0)); + testBigQueryWrite.addResponse(createAppendResponse(1)); + testBigQueryWrite.addResponse(createAppendResponse(2)); + + // First request: no prior responses or wait times have been recorded yet. + ApiFuture future1 = + sendTestMessage(connectionWorker, sw1, createFooProtoRows(new String[] {"0"}), 0); + future1.get(); + + AppendRowsRequest firstRequest = testBigQueryWrite.getAppendRequests().get(0); + long sizePerRequest = firstRequest.getProtoRows().getSerializedSize(); + assertThat(firstRequest.getClientStats().hasWindowStats()).isTrue(); + ClientStats.WindowStats firstWindowStats = firstRequest.getClientStats().getWindowStats(); + assertThat(firstWindowStats.hasMaxResponseLatencyMillis()).isFalse(); + assertThat(firstWindowStats.hasAvgResponseLatencyMillis()).isFalse(); + assertThat(firstWindowStats.hasLongestWaitNoResponseMillis()).isFalse(); + assertThat(firstWindowStats.hasRequestsSentCount()).isTrue(); + assertThat(firstWindowStats.getRequestsSentCount()).isEqualTo(1); + assertThat(firstWindowStats.hasResponsesReceivedCount()).isFalse(); + assertThat(firstWindowStats.hasBytesSentCount()).isTrue(); + assertThat(firstWindowStats.getBytesSentCount()).isEqualTo(sizePerRequest); + assertThat(firstWindowStats.hasWindowStartTimeEpochMillis()).isTrue(); + assertThat(firstWindowStats.getWindowStartTimeEpochMillis()).isAtLeast(beforeCreationMillis); + assertThat(firstWindowStats.getWindowStartTimeEpochMillis()).isAtMost(afterCreationMillis); + assertThat(firstWindowStats.hasWindowMillis()).isTrue(); + assertThat(firstWindowStats.getWindowMillis()).isAtLeast(0L); + + // Second request: sent after firstRequest waited in-flight (>100ms) and received a response + // (>=300ms latency), so all expected WindowStats fields are now non-zero and populated. + ApiFuture future2 = + sendTestMessage(connectionWorker, sw1, createFooProtoRows(new String[] {"1"}), 1); + future2.get(); + + AppendRowsRequest secondRequest = testBigQueryWrite.getAppendRequests().get(1); + assertThat(secondRequest.getClientStats().hasWindowStats()).isTrue(); + ClientStats.WindowStats secondWindowStats = secondRequest.getClientStats().getWindowStats(); + assertThat(secondWindowStats.hasMaxResponseLatencyMillis()).isTrue(); + assertThat(secondWindowStats.getMaxResponseLatencyMillis()).isAtLeast(msecResponseDelay); + assertThat(secondWindowStats.hasAvgResponseLatencyMillis()).isTrue(); + assertThat(secondWindowStats.getAvgResponseLatencyMillis()).isAtLeast(msecResponseDelay); + assertThat(secondWindowStats.hasLongestWaitNoResponseMillis()).isTrue(); + assertThat(secondWindowStats.getLongestWaitNoResponseMillis()).isGreaterThan(0L); + assertThat(secondWindowStats.hasRequestsSentCount()).isTrue(); + assertThat(secondWindowStats.getRequestsSentCount()).isEqualTo(2); + assertThat(secondWindowStats.hasResponsesReceivedCount()).isTrue(); + assertThat(secondWindowStats.getResponsesReceivedCount()).isEqualTo(1); + assertThat(secondWindowStats.hasBytesSentCount()).isTrue(); + assertThat(secondWindowStats.getBytesSentCount()).isEqualTo(2 * sizePerRequest); + assertThat(secondWindowStats.hasWindowStartTimeEpochMillis()).isTrue(); + assertThat(secondWindowStats.getWindowStartTimeEpochMillis()) + .isEqualTo(firstWindowStats.getWindowStartTimeEpochMillis()); + assertThat(secondWindowStats.hasWindowMillis()).isTrue(); + assertThat(secondWindowStats.getWindowMillis()).isAtLeast(msecResponseDelay); + + // Trigger a health check window rollover and verify WindowStats resets for the new window. + connectionWorker.setTestOnlyHealthCheckInterval(Duration.ofMillis(100)); + Thread.sleep(250); + + ApiFuture future3 = + sendTestMessage(connectionWorker, sw1, createFooProtoRows(new String[] {"2"}), 2); + future3.get(); + + AppendRowsRequest thirdRequest = testBigQueryWrite.getAppendRequests().get(2); + assertThat(thirdRequest.getClientStats().hasWindowStats()).isTrue(); + ClientStats.WindowStats thirdWindowStats = thirdRequest.getClientStats().getWindowStats(); + assertThat(thirdWindowStats.hasMaxResponseLatencyMillis()).isFalse(); + assertThat(thirdWindowStats.hasAvgResponseLatencyMillis()).isFalse(); + assertThat(thirdWindowStats.hasLongestWaitNoResponseMillis()).isFalse(); + assertThat(thirdWindowStats.hasRequestsSentCount()).isTrue(); + assertThat(thirdWindowStats.getRequestsSentCount()).isEqualTo(1); + assertThat(thirdWindowStats.hasResponsesReceivedCount()).isFalse(); + assertThat(thirdWindowStats.hasBytesSentCount()).isTrue(); + assertThat(thirdWindowStats.getBytesSentCount()).isEqualTo(sizePerRequest); + assertThat(thirdWindowStats.hasWindowStartTimeEpochMillis()).isTrue(); + assertThat(thirdWindowStats.getWindowStartTimeEpochMillis()) + .isGreaterThan(secondWindowStats.getWindowStartTimeEpochMillis()); + assertThat(thirdWindowStats.hasWindowMillis()).isTrue(); + assertThat(thirdWindowStats.getWindowMillis()).isAtLeast(0L); } } }