diff --git a/java-bigtable/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/EnhancedBigtableStub.java b/java-bigtable/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/EnhancedBigtableStub.java index 1f1a45c46b05..8b4fd04716c6 100644 --- a/java-bigtable/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/EnhancedBigtableStub.java +++ b/java-bigtable/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/EnhancedBigtableStub.java @@ -58,6 +58,7 @@ import com.google.bigtable.v2.ReadRowsResponse; import com.google.bigtable.v2.RowRange; import com.google.bigtable.v2.SampleRowKeysResponse; +import com.google.cloud.bigtable.Version; import com.google.cloud.bigtable.data.v2.internal.NameUtil; import com.google.cloud.bigtable.data.v2.internal.PrepareQueryRequest; import com.google.cloud.bigtable.data.v2.internal.PrepareResponse; @@ -116,6 +117,7 @@ import com.google.common.base.Functions; import com.google.common.base.MoreObjects; import com.google.common.base.Preconditions; +import com.google.common.collect.ImmutableList; import com.google.common.collect.ImmutableMap; import com.google.protobuf.ByteString; import io.grpc.MethodDescriptor; @@ -146,6 +148,7 @@ public class EnhancedBigtableStub implements AutoCloseable { private static final String CLIENT_NAME = "Bigtable"; private static final long FLOW_CONTROL_ADJUSTING_INTERVAL_MS = TimeUnit.SECONDS.toMillis(20); + private static final String BATCHER_API_CLIENT_TOKEN = "java-bigtable-batcher/" + Version.VERSION; private final ClientOperationSettings perOpSettings; private final BigtableClientContext bigtableClientContext; @@ -803,8 +806,7 @@ public Batcher newMutateRowsBatcher( perOpSettings.bulkMutateRowsSettings.getBatchingSettings(), bigtableClientContext.getClientContext().getExecutor(), bulkMutationFlowController, - MoreObjects.firstNonNull( - ctx, bigtableClientContext.getClientContext().getDefaultCallContext())); + batcherCallContext(ctx)); } /** @@ -835,8 +837,7 @@ public Batcher newMutateRowsBatcher( perOpSettings.bulkMutateRowsSettings.getBatchingSettings(), bigtableClientContext.getClientContext().getExecutor(), bulkMutationFlowController, - MoreObjects.firstNonNull( - ctx, bigtableClientContext.getClientContext().getDefaultCallContext())); + batcherCallContext(ctx)); } /** @@ -864,8 +865,7 @@ public Batcher newBulkReadRowsBatcher( perOpSettings.bulkReadRowsSettings.getBatchingSettings(), bigtableClientContext.getClientContext().getExecutor(), null, - MoreObjects.firstNonNull( - ctx, bigtableClientContext.getClientContext().getDefaultCallContext())); + batcherCallContext(ctx)); } /** @@ -1187,6 +1187,13 @@ private UnaryCallable createUserFacin bigtableClientContext.getClientContext().getDefaultCallContext()); } + private ApiCallContext batcherCallContext(@Nullable GrpcCallContext userCtx) { + return MoreObjects.firstNonNull( + userCtx, bigtableClientContext.getClientContext().getDefaultCallContext()) + .withExtraHeaders( + ImmutableMap.of("x-goog-api-client", ImmutableList.of(BATCHER_API_CLIENT_TOKEN))); + } + private Map composeRequestParams( String appProfileId, String tableName, String authorizedViewName) { if (tableName.isEmpty() && !authorizedViewName.isEmpty()) { diff --git a/java-bigtable/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/HeadersTest.java b/java-bigtable/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/HeadersTest.java index dcf22e1e19d4..609a59497d32 100644 --- a/java-bigtable/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/HeadersTest.java +++ b/java-bigtable/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/HeadersTest.java @@ -54,6 +54,7 @@ import com.google.cloud.bigtable.data.v2.models.RowMutationEntry; import com.google.cloud.bigtable.data.v2.models.TableId; import com.google.cloud.bigtable.data.v2.models.sql.PreparedStatement; +import com.google.protobuf.ByteString; import com.google.rpc.Status; import io.grpc.Metadata; import io.grpc.Server; @@ -111,7 +112,7 @@ public void setUp() throws Exception { HeaderProvider headerProvider = FixedHeaderProvider.create(TEST_FIXED_HEADER_STRING, "test_header_value"); - // Force immediate flush + // Force immediate flush for both batcher types settings .stubSettings() .setHeaderProvider(headerProvider) @@ -120,6 +121,13 @@ public void setUp() throws Exception { settings.stubSettings().bulkMutateRowsSettings().getBatchingSettings().toBuilder() .setElementCountThreshold(1L) .build()); + settings + .stubSettings() + .bulkReadRowsSettings() + .setBatchingSettings( + settings.stubSettings().bulkReadRowsSettings().getBatchingSettings().toBuilder() + .setElementCountThreshold(1L) + .build()); client = BigtableDataClient.create(settings.build()); } @@ -207,6 +215,40 @@ public void prepareQueryTest() { verifyHeaderSent(true); } + @Test + public void bulkMutationBatcherHasBatcherToken() throws Exception { + try (Batcher batcher = client.newBulkMutationBatcher(TABLE_ID)) { + batcher.add(RowMutationEntry.create("fake-key").deleteRow()); + } + Metadata metadata = sentMetadata.take(); + assertThat(hasApiClientToken(metadata, "java-bigtable-batcher")).isTrue(); + } + + @Test + public void bulkReadRowsBatcherHasBatcherToken() throws Exception { + try (Batcher batcher = client.newBulkReadRowsBatcher(TABLE_ID)) { + batcher.add(ByteString.copyFromUtf8("fake-key")); + } + Metadata metadata = sentMetadata.take(); + assertThat(hasApiClientToken(metadata, "java-bigtable-batcher")).isTrue(); + } + + @Test + public void regularRpcDoesNotHaveBatcherToken() throws Exception { + client.mutateRowAsync(RowMutation.create(TABLE_ID, "fake-key").deleteRow()).get(); + Metadata metadata = sentMetadata.take(); + assertThat(hasApiClientToken(metadata, "java-bigtable-batcher")).isFalse(); + } + + private static boolean hasApiClientToken(Metadata metadata, String token) { + Iterable values = metadata.getAll(API_CLIENT_HEADER_KEY); + if (values == null) return false; + for (String value : values) { + if (value.contains(token)) return true; + } + return false; + } + private void verifyHeaderSent() { verifyHeaderSent(false); }