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 @@ -20,6 +20,7 @@
import static com.google.cloud.firestore.telemetry.TraceUtil.ATTRIBUTE_KEY_ATTEMPT;

import com.google.api.core.ApiFuture;
import com.google.api.core.BetaApi;
import com.google.api.core.InternalExtensionOnly;
import com.google.api.core.SettableApiFuture;
import com.google.api.gax.rpc.ResponseObserver;
Expand All @@ -35,6 +36,7 @@
import com.google.cloud.firestore.telemetry.TraceUtil.Scope;
import com.google.cloud.firestore.v1.FirestoreSettings;
import com.google.common.collect.ImmutableMap;
import com.google.firestore.v1.RequestOptions;
import com.google.firestore.v1.RunAggregationQueryRequest;
import com.google.firestore.v1.RunAggregationQueryResponse;
import com.google.firestore.v1.RunQueryRequest;
Expand Down Expand Up @@ -100,7 +102,20 @@ Pipeline pipeline() {
*/
@Nonnull
public ApiFuture<AggregateQuerySnapshot> get() {
return get(null, null);
return get(null, null, null, null);
}

/**
* Executes this query with execution options.
*
* @param executionOptions Options for executing the request.
* @return An {@link ApiFuture} that will be resolved with the results of the query.
*/
@BetaApi
@Nonnull
public ApiFuture<AggregateQuerySnapshot> get(
@Nonnull FirestoreExecutionOptions executionOptions) {
return get(null, null, executionOptions.getExplainOptions(), executionOptions);
}

/**
Expand All @@ -113,6 +128,35 @@ public ApiFuture<AggregateQuerySnapshot> get() {
*/
@Nonnull
public ApiFuture<ExplainResults<AggregateQuerySnapshot>> explain(ExplainOptions options) {
return explain(options, null);
}

/**
* Plans and optionally executes this query with execution options.
*
* @param executionOptions Options for executing the request.
* @return An ApiFuture that will be resolved with the planner information, statistics from the
* query execution (if any), and the query results (if any).
*/
@BetaApi
@Nonnull
public ApiFuture<ExplainResults<AggregateQuerySnapshot>> explain(
@Nonnull FirestoreExecutionOptions executionOptions) {
return explain(executionOptions.getExplainOptions(), executionOptions);
}

/**
* Plans and optionally executes this query with explain options and execution options.
*
* @param options The options for explain.
* @param executionOptions Options for executing the request.
* @return An ApiFuture that will be resolved with the planner information, statistics from the
* query execution (if any), and the query results (if any).
*/
@BetaApi
@Nonnull
public ApiFuture<ExplainResults<AggregateQuerySnapshot>> explain(
@Nullable ExplainOptions options, @Nullable FirestoreExecutionOptions executionOptions) {
TraceUtil.Span span =
getTraceUtil().startSpan(TelemetryConstants.METHOD_NAME_AGGREGATION_QUERY_GET);

Expand All @@ -125,8 +169,11 @@ public ApiFuture<ExplainResults<AggregateQuerySnapshot>> explain(ExplainOptions
/* transactionId= */ null,
/* readTime= */ null,
/* startTimeNanos= */ query.rpcContext.getClock().nanoTime(),
/* explainOptions= */ options,
metricsContext);
/* explainOptions= */ options != null
? options
: (executionOptions != null ? executionOptions.getExplainOptions() : null),
metricsContext,
executionOptions);
runQuery(responseDeliverer, /* attempt */ 0);
ApiFuture<ExplainResults<AggregateQuerySnapshot>> result = responseDeliverer.getFuture();
span.endAtFuture(result);
Expand All @@ -140,6 +187,15 @@ public ApiFuture<ExplainResults<AggregateQuerySnapshot>> explain(ExplainOptions
@Nonnull
ApiFuture<AggregateQuerySnapshot> get(
@Nullable final ByteString transactionId, @Nullable com.google.protobuf.Timestamp readTime) {
return get(transactionId, readTime, null, null);
}

@Nonnull
ApiFuture<AggregateQuerySnapshot> get(
@Nullable final ByteString transactionId,
@Nullable com.google.protobuf.Timestamp readTime,
@Nullable ExplainOptions explainOptions,
@Nullable FirestoreExecutionOptions executionOptions) {
TraceUtil.Span span =
getTraceUtil()
.startSpan(
Expand All @@ -159,7 +215,8 @@ ApiFuture<AggregateQuerySnapshot> get(
transactionId,
readTime,
/* startTimeNanos= */ query.rpcContext.getClock().nanoTime(),
metricsContext);
metricsContext,
executionOptions);
runQuery(responseDeliverer, /* attempt= */ 0);
ApiFuture<AggregateQuerySnapshot> result = responseDeliverer.getFuture();
span.endAtFuture(result);
Expand All @@ -175,7 +232,8 @@ private <T> void runQuery(ResponseDeliverer<T> responseDeliverer, int attempt) {
toProto(
responseDeliverer.getTransactionId(),
responseDeliverer.getReadTime(),
responseDeliverer.getExplainOptions());
responseDeliverer.getExplainOptions(),
responseDeliverer.getExecutionOptions());
AggregateQueryResponseObserver<T> responseObserver =
new AggregateQueryResponseObserver<T>(responseDeliverer, attempt);
ServerStreamingCallable<RunAggregationQueryRequest, RunAggregationQueryResponse> callable =
Expand All @@ -197,16 +255,24 @@ private abstract static class ResponseDeliverer<T> {
private final long startTimeNanos;
private final SettableApiFuture<T> future = SettableApiFuture.create();
private MetricsContext metricsContext;
private final @Nullable FirestoreExecutionOptions executionOptions;

ResponseDeliverer(
@Nullable ByteString transactionId,
@Nullable com.google.protobuf.Timestamp readTime,
long startTimeNanos,
MetricsContext metricsContext) {
MetricsContext metricsContext,
@Nullable FirestoreExecutionOptions executionOptions) {
this.transactionId = transactionId;
this.readTime = readTime;
this.startTimeNanos = startTimeNanos;
this.metricsContext = metricsContext;
this.executionOptions = executionOptions;
}

@Nullable
FirestoreExecutionOptions getExecutionOptions() {
return executionOptions;
}

@Nullable
Expand Down Expand Up @@ -265,8 +331,9 @@ private class AggregateQueryResponseDeliverer extends ResponseDeliverer<Aggregat
@Nullable ByteString transactionId,
@Nullable com.google.protobuf.Timestamp readTime,
long startTimeNanos,
MetricsContext metricsContext) {
super(transactionId, readTime, startTimeNanos, metricsContext);
MetricsContext metricsContext,
@Nullable FirestoreExecutionOptions executionOptions) {
super(transactionId, readTime, startTimeNanos, metricsContext, executionOptions);
}

@Override
Expand All @@ -293,8 +360,9 @@ private final class AggregateQueryExplainResponseDeliverer
@Nullable com.google.protobuf.Timestamp readTime,
long startTimeNanos,
@Nullable ExplainOptions explainOptions,
MetricsContext metricsContext) {
super(transactionId, readTime, startTimeNanos, metricsContext);
MetricsContext metricsContext,
@Nullable FirestoreExecutionOptions executionOptions) {
super(transactionId, readTime, startTimeNanos, metricsContext, executionOptions);
this.explainOptions = explainOptions;
}

Expand Down Expand Up @@ -437,14 +505,24 @@ public void onComplete() {
*/
@Nonnull
public RunAggregationQueryRequest toProto() {
return toProto(/* transactionId= */ null, /* readTime= */ null, /* explainOptions= */ null);
return toProto(
/* transactionId= */ null, /* readTime= */ null, /* explainOptions= */ null, null);
}

@Nonnull
RunAggregationQueryRequest toProto(
@Nullable final ByteString transactionId,
@Nullable final com.google.protobuf.Timestamp readTime,
@Nullable ExplainOptions explainOptions) {
return toProto(transactionId, readTime, explainOptions, null);
}

@Nonnull
RunAggregationQueryRequest toProto(
@Nullable final ByteString transactionId,
@Nullable final com.google.protobuf.Timestamp readTime,
@Nullable ExplainOptions explainOptions,
@Nullable FirestoreExecutionOptions executionOptions) {
RunQueryRequest runQueryRequest = query.toProto();

RunAggregationQueryRequest.Builder request = RunAggregationQueryRequest.newBuilder();
Expand All @@ -460,6 +538,13 @@ RunAggregationQueryRequest toProto(
request.setExplainOptions(explainOptions.toProto());
}

RequestOptions requestOptions =
RequestOptionsHelper.createRequestOptions(
query.rpcContext.getFirestore().getOptions(), executionOptions);
if (!requestOptions.equals(RequestOptions.getDefaultInstance())) {
request.setRequestOptions(requestOptions);
}

StructuredAggregationQuery.Builder structuredAggregationQuery =
request.getStructuredAggregationQueryBuilder();
structuredAggregationQuery.setStructuredQuery(runQueryRequest.getStructuredQuery());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@
import com.google.common.util.concurrent.MoreExecutors;
import com.google.firestore.v1.BatchWriteRequest;
import com.google.firestore.v1.BatchWriteResponse;
import com.google.firestore.v1.RequestOptions;
import io.grpc.Status;
import java.util.ArrayList;
import java.util.List;
Expand Down Expand Up @@ -107,6 +108,12 @@ private BatchWriteRequest buildBatchWriteRequest() {
BatchWriteRequest.Builder builder = BatchWriteRequest.newBuilder();
builder.setDatabase(firestore.getDatabaseName());
forEachWrite(builder::addWrites);
RequestOptions requestOptions =
RequestOptionsHelper.createRequestOptions(
firestore.getOptions(), (FirestoreExecutionOptions) null);
if (!requestOptions.equals(RequestOptions.getDefaultInstance())) {
builder.setRequestOptions(requestOptions);
}
return builder.build();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@

import com.google.api.core.ApiFuture;
import com.google.api.core.ApiFutures;
import com.google.api.core.BetaApi;
import com.google.api.gax.rpc.ApiException;
import com.google.api.gax.rpc.ApiExceptions;
import com.google.api.gax.rpc.ApiStreamObserver;
Expand All @@ -32,6 +33,7 @@
import com.google.common.util.concurrent.MoreExecutors;
import com.google.firestore.v1.Cursor;
import com.google.firestore.v1.PartitionQueryRequest;
import com.google.firestore.v1.RequestOptions;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
Expand Down Expand Up @@ -72,11 +74,27 @@ public class CollectionGroup extends Query {
*/
public void getPartitions(
long desiredPartitionCount, final ApiStreamObserver<QueryPartition> observer) {
getPartitions(desiredPartitionCount, observer, null);
}

/**
* Partitions a query by returning partition cursors that can be used to run the query in
* parallel, with execution options.
*
* @param desiredPartitionCount The desired maximum number of partition points.
* @param observer a stream observer that receives the result of the Partition request.
* @param executionOptions Options for executing the request.
*/
@BetaApi
public void getPartitions(
long desiredPartitionCount,
final ApiStreamObserver<QueryPartition> observer,
@Nullable FirestoreExecutionOptions executionOptions) {
if (desiredPartitionCount == 1) {
// Short circuit if the user only requested a single partition.
observer.onNext(new QueryPartition(partitionQuery, null, null));
} else {
PartitionQueryRequest request = buildRequest(desiredPartitionCount);
PartitionQueryRequest request = buildRequest(desiredPartitionCount, executionOptions);

final PartitionQueryPagedResponse response;
try {
Expand All @@ -100,12 +118,24 @@ public void getPartitions(
}

public ApiFuture<List<QueryPartition>> getPartitions(long desiredPartitionCount) {
return getPartitions(desiredPartitionCount, (FirestoreExecutionOptions) null);
}

/**
* Partitions a query by returning partition cursors with execution options.
*
* @param desiredPartitionCount The desired maximum number of partition points.
* @param executionOptions Options for executing the request.
*/
@BetaApi
public ApiFuture<List<QueryPartition>> getPartitions(
long desiredPartitionCount, @Nullable FirestoreExecutionOptions executionOptions) {
if (desiredPartitionCount == 1) {
// Short circuit if the user only requested a single partition.
return ApiFutures.immediateFuture(
Collections.singletonList(new QueryPartition(partitionQuery, null, null)));
} else {
PartitionQueryRequest request = buildRequest(desiredPartitionCount);
PartitionQueryRequest request = buildRequest(desiredPartitionCount, executionOptions);

TraceUtil.Span span =
rpcContext
Expand Down Expand Up @@ -152,7 +182,8 @@ public ApiFuture<List<QueryPartition>> getPartitions(long desiredPartitionCount)
}
}

private PartitionQueryRequest buildRequest(long desiredPartitionCount) {
private PartitionQueryRequest buildRequest(
long desiredPartitionCount, @Nullable FirestoreExecutionOptions executionOptions) {
Preconditions.checkArgument(
desiredPartitionCount > 0, "Desired partition count must be one or greater");

Expand All @@ -163,6 +194,12 @@ private PartitionQueryRequest buildRequest(long desiredPartitionCount) {
// Since we are always returning an extra partition (with en empty endBefore cursor), we
// reduce the desired partition count by one.
request.setPartitionCount(desiredPartitionCount - 1);
RequestOptions requestOptions =
RequestOptionsHelper.createRequestOptions(
rpcContext.getFirestore().getOptions(), executionOptions);
if (!requestOptions.equals(RequestOptions.getDefaultInstance())) {
request.setRequestOptions(requestOptions);
}
return request.build();
}

Expand Down
Loading
Loading