Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down Expand Up @@ -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: "
Expand Down Expand Up @@ -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);

Expand All @@ -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()) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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<AppendRowsResponse> 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<AppendRowsResponse> 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<AppendRowsResponse> 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);
}
}
}
Loading