Skip to content
Merged
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 @@ -163,7 +163,22 @@ public class BigtableOptionsFactory {
public static final String MAX_ELAPSED_BACKOFF_MILLIS_KEY =
"google.bigtable.grpc.retry.max.elapsed.backoff.ms";

/** Key to set the amount of time to wait when reading a partial row. */
/**
* Key to set how long a read may go without receiving a response before the stream is cancelled
* and retried. This is the gap between consecutive responses, reset every time the server sends
* something; it is not a deadline for the attempt as a whole. Raise it for scans that can
* legitimately go a long time without producing a row, such as a filtered scan over a large
* table. Defaults to 5 minutes.
*
* <p>An attempt is separately bounded by {@link #BIGTABLE_READ_RPC_ATTEMPT_TIMEOUT_MS_KEY}, 10
* minutes by default, so raising this beyond 10 minutes accomplishes nothing on its own — the
* attempt deadline would fire first. Raise that key too if you need a longer gap than that.
*
* <p>This can only raise the watchdog, never lower it: values below the 5 minute default are
* ignored. The key used to be handed to gax as the rpc timeout, where gax discarded it because
* the read path always put a timeout on the call context, so it never had any effect. Clamping it
* to the default keeps a previously ignored short value from suddenly cancelling reads.
*/
public static final String READ_PARTIAL_ROW_TIMEOUT_MS =
"google.bigtable.grpc.read.partial.row.timeout.ms";

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -338,10 +338,8 @@ private BigtableDataSettings buildBigtableDataSettings(ClientOperationTimeouts c
// Configure metrics
configureMetricsBridge(dataBuilder);

// Configure RPCs - this happens in two parts:
// - most of the timeouts are defined here
// - attempt timeouts for readRows is set in DataClientVeneerApi to workaround lack of attempt
// timeouts for streaming RPCs
// Configure RPCs. All of the timeouts are defined here; DataClientVeneerApi only sets the
// gRPC deadline for the overall operation.
// Complex RPC method settings
configureBulkMutationSettings(
dataBuilder.stubSettings().bulkMutateRowsSettings(),
Expand Down Expand Up @@ -748,8 +746,9 @@ private void configureReadRowsSettings(
OperationTimeouts operationTimeouts) {

// Configure retries
// NOTE: that similar but not the same as unary retry settings: per attempt timeouts don't
// exist, instead we use READ_PARTIAL_ROW_TIMEOUT_MS as the intra-row timeout
// NOTE: similar but not the same as unary retry settings: responseTimeout is the watchdog
// wait timeout, separate from the per attempt timeout, and the attempt deadline goes on the
// retry settings rather than the ApiCallContext. See below.
if (!configuration.getBoolean(ENABLE_GRPC_RETRIES_KEY, true)) {
// user explicitly disabled retries, treat it as a non-idempotent method
readRowsSettings.setRetryableCodes(Collections.emptySet());
Expand Down Expand Up @@ -780,16 +779,32 @@ private void configureReadRowsSettings(
configuration.getInt(MAX_SCAN_TIMEOUT_RETRIES, MAX_CONSECUTIVE_SCAN_ATTEMPTS));
}

// Per response timeouts (note: gax maps rpcTimeouts to response timeouts for streaming rpcs)
// The watchdog wait timeout: how long the stream may go without receiving a response before
// it is cancelled and retried. Reset by every response, so it bounds the gap between them
// rather than the attempt as a whole.
//
// This used to be handed to gax as the rpcTimeout, where it was silently discarded because
// the ApiCallContext already carried a timeout, so the key never had any effect. Honoring it
// outright would newly cancel reads for anyone who had set a short value, so it is only
// allowed to raise the watchdog, never lower it. Values below the client default are ignored.
if (operationTimeouts.getResponseTimeout().isPresent()) {
Duration defaultWaitTimeout = readRowsSettings.getWaitTimeout();
Duration responseTimeout = operationTimeouts.getResponseTimeout().get();
if (defaultWaitTimeout == null || responseTimeout.compareTo(defaultWaitTimeout) > 0) {
readRowsSettings.setWaitTimeout(responseTimeout);
}
Comment on lines +791 to +795

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

medium

If defaultWaitTimeout is null, the code currently sets the wait timeout to responseTimeout without any lower-bound check. However, the client default watchdog timeout is 5 minutes. If responseTimeout is less than 5 minutes, this would lower the watchdog timeout, violating the "raise-only" contract described in the Javadoc and PR description. To prevent this, we should default to a minimum of 5 minutes when defaultWaitTimeout is null.

      Duration defaultWaitTimeout = readRowsSettings.getWaitTimeout();
      Duration responseTimeout = operationTimeouts.getResponseTimeout().get();
      Duration minWaitTimeout = defaultWaitTimeout != null ? defaultWaitTimeout : Duration.ofMinutes(5);
      if (responseTimeout.compareTo(minWaitTimeout) > 0) {
        readRowsSettings.setWaitTimeout(responseTimeout);
      }

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Default timeout is 5 minutes.

}

// Attempt timeouts. gax applies rpcTimeout as the deadline for a single attempt, but only if
// the ApiCallContext doesn't already carry a timeout of its own (see
// ServerStreamingAttemptCallable#call), so DataClientVeneerApi deliberately leaves it unset.
if (operationTimeouts.getAttemptTimeout().isPresent()) {
readRowsSettings
.retrySettings()
.setInitialRpcTimeout(operationTimeouts.getResponseTimeout().get())
.setMaxRpcTimeout(operationTimeouts.getResponseTimeout().get());
.setInitialRpcTimeout(operationTimeouts.getAttemptTimeout().get())
.setMaxRpcTimeout(operationTimeouts.getAttemptTimeout().get());
}

// Attempt timeouts are set in DataClientVeneerApi

// overall timeout
if (operationTimeouts.getOperationTimeout().isPresent()) {
readRowsSettings
Expand Down Expand Up @@ -1004,9 +1019,9 @@ static class OperationTimeouts {
new OperationTimeouts(
Optional.<Duration>absent(), Optional.<Duration>absent(), Optional.<Duration>absent());

// responseTimeouts are only relevant to streaming RPCs, they limit the amount of timeout a
// stream will wait for the next response message. This is synonymous with attemptTimeouts in
// unary RPCs since they receive a single response (so its ignored).
// responseTimeouts are only relevant to streaming RPCs, they limit how long a stream will
// wait for the next response message. Unary RPCs receive a single response, so it's ignored
// there.
private final Optional<Duration> responseTimeout;
private final Optional<Duration> attemptTimeout;
private final Optional<Duration> operationTimeout;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,6 @@
import com.google.cloud.bigtable.data.v2.models.Row;
import com.google.cloud.bigtable.data.v2.models.RowMutation;
import com.google.cloud.bigtable.hbase.adapters.Adapters;
import com.google.cloud.bigtable.hbase.util.Logger;
import com.google.cloud.bigtable.hbase.wrappers.BulkMutationWrapper;
import com.google.cloud.bigtable.hbase.wrappers.BulkReadWrapper;
import com.google.cloud.bigtable.hbase.wrappers.DataClientWrapper;
Expand All @@ -60,14 +59,11 @@
import org.apache.hadoop.hbase.client.AbstractClientScanner;
import org.apache.hadoop.hbase.client.Result;
import org.apache.hadoop.hbase.client.ResultScanner;
import org.threeten.bp.Duration;

/** For internal use only - public for technical reasons. */
@InternalApi("For internal usage only")
public class DataClientVeneerApi implements DataClientWrapper {

private final Logger LOG = new Logger(DataClientVeneerApi.class);

private static final RowResultAdapter RESULT_ADAPTER = new RowResultAdapter();

private final BigtableDataClient delegate;
Expand Down Expand Up @@ -174,12 +170,16 @@ private ApiCallContext createReadRowCallContext() {
GrpcCallContext ctx = GrpcCallContext.createDefault();
OperationTimeouts callSettings = clientOperationTimeouts.getUnaryTimeouts();

if (callSettings.getAttemptTimeout().isPresent()) {
ctx = ctx.withTimeout(callSettings.getAttemptTimeout().get());
}
// TODO: remove this after fixing it in veneer/gax
// If the attempt timeout was overridden, it disables overall timeout limiting
// Fix it by settings the underlying grpc deadline
// NOTE: the attempt timeout is deliberately not set here. gax applies the retry settings'
// rpcTimeout as the attempt deadline, but only when the context has no timeout of its own,
// so setting one here would suppress it. See BigtableHBaseVeneerSettings, which puts the
// attempt timeout on the retry settings instead.

// Kept for backward compatibility: this context is built fresh per call, so the grpc deadline
// below bounds a single readRow. gax now clamps each attempt's rpcTimeout to the time left on
// the total timeout (ExponentialRetryAlgorithm#createNextAttempt), and readRowSettings already
// carries that total timeout, so this deadline is very likely redundant. Dropping it isn't
// provably behavior neutral though, so it stays.
if (callSettings.getOperationTimeout().isPresent()) {
ctx =
ctx.withCallOptions(
Expand All @@ -191,10 +191,17 @@ private ApiCallContext createReadRowCallContext() {
return ctx;
}

// Support 2 bigtable-hbase features not directly available in veneer:
// - per attempt deadlines - vener doesn't implement deadlines for attempts. To workaround this,
// the timeouts are set per call in the ApiCallContext. However this creates a separate issue of
// over running the operation deadline, so gRPC deadline is also set.
// Kept for backward compatibility: veneer has no operation deadline for streaming RPCs, so the
// grpc deadline below stands in for one. Most callers build a context per call, but
// PaginatedRowResultScanner holds onto the one it is handed and passes it to every segment
// fetch, and Deadline.after() is absolute, so on that path the deadline bounds the scanner's
// whole lifetime rather than a single ReadRows. Removing it in favor of gax's per operation
// total timeout would hand each segment its own fresh budget, which is a real behavior change.
//
// The attempt deadline is deliberately *not* set here. gax applies the retry settings'
// rpcTimeout as the attempt deadline, but only when the context carries no timeout of its own
// (ServerStreamingAttemptCallable#call), so setting one here would suppress it. The attempt
// timeout goes on the retry settings in BigtableHBaseVeneerSettings instead.
private GrpcCallContext createScanCallContext() {
GrpcCallContext ctx = GrpcCallContext.createDefault();
OperationTimeouts callSettings = clientOperationTimeouts.getScanTimeouts();
Expand All @@ -206,11 +213,6 @@ private GrpcCallContext createScanCallContext() {
Deadline.after(
callSettings.getOperationTimeout().get().toMillis(), TimeUnit.MILLISECONDS)));
}
if (callSettings.getAttemptTimeout().isPresent()) {
Duration attemptTimeout = callSettings.getAttemptTimeout().get();
LOG.info("effective attempt timeout for scan is %s", attemptTimeout);
ctx = ctx.withTimeout(attemptTimeout);
}

return ctx;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -55,10 +55,13 @@
import com.google.api.gax.core.NoCredentialsProvider;
import com.google.api.gax.grpc.InstantiatingGrpcChannelProvider;
import com.google.api.gax.retrying.RetrySettings;
import com.google.api.gax.rpc.ServerStreamingCallSettings;
import com.google.api.gax.rpc.UnaryCallSettings;
import com.google.auth.Credentials;
import com.google.cloud.bigtable.admin.v2.BigtableTableAdminSettings;
import com.google.cloud.bigtable.data.v2.BigtableDataSettings;
import com.google.cloud.bigtable.data.v2.models.Query;
import com.google.cloud.bigtable.data.v2.models.Row;
import com.google.cloud.bigtable.hbase.BigtableConfiguration;
import com.google.cloud.bigtable.hbase.BigtableHBaseVersion;
import com.google.cloud.bigtable.hbase.BigtableOptionsFactory;
Expand Down Expand Up @@ -294,10 +297,15 @@ public void testTimeoutBeingPassed() throws IOException {
RetrySettings readRowsRetrySettings =
dataSettings.getStubSettings().readRowsSettings().getRetrySettings();
assertEquals(initialElapsedMs, readRowsRetrySettings.getInitialRetryDelay().toMillis());
assertEquals(perRowTimeoutMs, readRowsRetrySettings.getInitialRpcTimeout().toMillis());
assertEquals(perRowTimeoutMs, readRowsRetrySettings.getMaxRpcTimeout().toMillis());
assertEquals(
readRowStreamAttemptTimeout, readRowsRetrySettings.getInitialRpcTimeout().toMillis());
assertEquals(readRowStreamAttemptTimeout, readRowsRetrySettings.getMaxRpcTimeout().toMillis());
assertEquals(maxAttempt, readRowsRetrySettings.getMaxAttempts());
assertEquals(readRowStreamTimeout, readRowsRetrySettings.getTotalTimeout().toMillis());
// The per row timeout is the watchdog wait timeout, not an attempt deadline. 1001ms is below
// the 5 minute client default, and the key can only raise the watchdog, so it is ignored.
assertEquals(
Duration.ofMinutes(5), dataSettings.getStubSettings().readRowsSettings().getWaitTimeout());

RetrySettings sampleRowKeysRetrySettings =
dataSettings.getStubSettings().sampleRowKeysSettings().getRetrySettings();
Expand All @@ -307,6 +315,100 @@ public void testTimeoutBeingPassed() throws IOException {
assertEquals(rpcAttemptTimeoutMs, sampleRowKeysRetrySettings.getMaxRpcTimeout().toMillis());
}

@Test
public void testReadRowsWaitTimeout() throws IOException {
BigtableDataSettings defaultSettings =
((BigtableHBaseVeneerSettings) BigtableHBaseVeneerSettings.create(configuration))
.getDataSettings();

// READ_PARTIAL_ROW_TIMEOUT_MS is the watchdog wait timeout: what cancels a sparse scan that
// goes 5 minutes without a row. Its default matches the veneer default, so wiring it through
// does not change the out of the box behavior. See also
// testPartialRowTimeoutBelowTheDefaultIsIgnored.
assertEquals(
Duration.ofMinutes(5),
defaultSettings.getStubSettings().readRowsSettings().getWaitTimeout());

configuration.set(READ_PARTIAL_ROW_TIMEOUT_MS, "540000");
BigtableDataSettings dataSettings =
((BigtableHBaseVeneerSettings) BigtableHBaseVeneerSettings.create(configuration))
.getDataSettings();

assertEquals(
Duration.ofMinutes(9), dataSettings.getStubSettings().readRowsSettings().getWaitTimeout());
// The wait timeout is independent of the idle timeout and of the attempt deadline.
assertEquals(
defaultSettings.getStubSettings().readRowsSettings().getIdleTimeout(),
dataSettings.getStubSettings().readRowsSettings().getIdleTimeout());
assertEquals(
defaultSettings.getStubSettings().readRowsSettings().getRetrySettings().toString(),
dataSettings.getStubSettings().readRowsSettings().getRetrySettings().toString());
}

@Test
public void testReadRowsAttemptTimeoutIsOnTheRetrySettings() throws IOException {
// The attempt deadline has to live on the retry settings rather than the ApiCallContext: gax
// only applies rpcTimeout when the context carries no timeout of its own. See
// TestReadRowsTimeoutSemantics.
BigtableDataSettings defaultSettings =
((BigtableHBaseVeneerSettings) BigtableHBaseVeneerSettings.create(configuration))
.getDataSettings();
assertEquals(
Duration.ofMinutes(10),
defaultSettings.getStubSettings().readRowsSettings().getRetrySettings().getMaxRpcTimeout());

configuration.set(BIGTABLE_READ_RPC_ATTEMPT_TIMEOUT_MS_KEY, "900000");
BigtableHBaseVeneerSettings settings =
(BigtableHBaseVeneerSettings) BigtableHBaseVeneerSettings.create(configuration);

ServerStreamingCallSettings<Query, Row> readRows =
settings.getDataSettings().getStubSettings().readRowsSettings();
assertEquals(Duration.ofMinutes(15), readRows.getRetrySettings().getInitialRpcTimeout());
assertEquals(Duration.ofMinutes(15), readRows.getRetrySettings().getMaxRpcTimeout());
assertEquals(
Optional.of(Duration.ofMinutes(15)),
settings.getClientTimeouts().getScanTimeouts().getAttemptTimeout());
}

@Test
public void testReadRowsWaitTimeoutBeyondTheAttemptTimeoutNeedsBothRaised() throws IOException {
// A 15 minute wait timeout can never fire against the default 10 minute attempt deadline.
configuration.set(READ_PARTIAL_ROW_TIMEOUT_MS, "900000");
BigtableHBaseVeneerSettings settings =
(BigtableHBaseVeneerSettings) BigtableHBaseVeneerSettings.create(configuration);

ServerStreamingCallSettings<Query, Row> readRows =
settings.getDataSettings().getStubSettings().readRowsSettings();
assertEquals(Duration.ofMinutes(15), readRows.getWaitTimeout());
assertEquals(Duration.ofMinutes(10), readRows.getRetrySettings().getMaxRpcTimeout());

configuration.set(BIGTABLE_READ_RPC_ATTEMPT_TIMEOUT_MS_KEY, "900000");
readRows =
((BigtableHBaseVeneerSettings) BigtableHBaseVeneerSettings.create(configuration))
.getDataSettings()
.getStubSettings()
.readRowsSettings();
assertEquals(Duration.ofMinutes(15), readRows.getWaitTimeout());
assertEquals(Duration.ofMinutes(15), readRows.getRetrySettings().getMaxRpcTimeout());
}

@Test
public void testPartialRowTimeoutBelowTheDefaultIsIgnored() throws IOException {
// The key is raise only. It never reached the wire before (gax discarded it because the call
// context carried a timeout), so honoring a short value now would newly cancel reads that
// used to survive. It is still parsed into the client timeouts, just not applied.
configuration.set(READ_PARTIAL_ROW_TIMEOUT_MS, "1000");
BigtableHBaseVeneerSettings settings =
(BigtableHBaseVeneerSettings) BigtableHBaseVeneerSettings.create(configuration);

assertEquals(
Optional.of(Duration.ofMillis(1000)),
settings.getClientTimeouts().getScanTimeouts().getResponseTimeout());
assertEquals(
Duration.ofMinutes(5),
settings.getDataSettings().getStubSettings().readRowsSettings().getWaitTimeout());
}

@Test
public void testWhenRetriesAreDisabled() throws IOException {
configuration.setBoolean(ENABLE_GRPC_RETRIES_KEY, false);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,11 @@ public static CloudBigtableScanConfiguration buildExportConfig(ExportOptions opt
BigtableOptionsFactory.BIGTABLE_READ_RPC_ATTEMPT_TIMEOUT_MS_KEY,
options.getBigtableReadRpcAttemptTimeoutMs());
}
if (options.getBigtableReadPartialRowTimeoutMs() != null) {
configBuilder.withConfiguration(
BigtableOptionsFactory.READ_PARTIAL_ROW_TIMEOUT_MS,
options.getBigtableReadPartialRowTimeoutMs());
}
if (options.getBigtableMaxAttempts() != null) {
configBuilder.withConfiguration(
BigtableOptionsFactory.MAX_SCAN_TIMEOUT_RETRIES, options.getBigtableMaxAttempts());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -197,6 +197,19 @@ public interface ExportOptions extends GcpOptions, GcsOptions {

@SuppressWarnings("unused")
void setBigtableMaxAttempts(ValueProvider<String> maxAttempts);

@Description(
"How long a scan may go without receiving a response, in milliseconds, before it is "
+ "cancelled and retried. This is the gap between consecutive responses, not a "
+ "deadline for the attempt, so raise it for a filtered scan that can traverse a lot "
+ "of non-matching rows between results. Defaults to 300000 (5 minutes); lower "
+ "values are ignored. A single attempt is separately capped by "
+ "--bigtableReadRpcAttemptTimeoutMs (10 minutes by default), so raise that as well "
+ "if you need a gap longer than that.")
ValueProvider<String> getBigtableReadPartialRowTimeoutMs();

@SuppressWarnings("unused")
void setBigtableReadPartialRowTimeoutMs(ValueProvider<String> partialRowTimeoutMs);
}

public static void main(String[] args) {
Expand Down
Loading