Skip to content

Commit 0361f51

Browse files
committed
Improve sharding for bounded pcolleciton, to better control concurrent connections. this will keep elements for same destination close to each other and shard them. For single table write it's same behaviour, for dynamic destination it will improve reduce amount of connections used
1 parent 3d68e9d commit 0361f51

1 file changed

Lines changed: 25 additions & 1 deletion

File tree

  • sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery

sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/bigquery/StorageApiLoads.java

Lines changed: 25 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@
2323
import com.google.cloud.bigquery.storage.v1.AppendRowsRequest;
2424
import java.io.IOException;
2525
import java.nio.ByteBuffer;
26+
import java.nio.charset.StandardCharsets;
2627
import java.util.Map;
2728
import java.util.concurrent.ThreadLocalRandom;
2829
import java.util.function.Predicate;
@@ -38,6 +39,8 @@
3839
import org.apache.beam.sdk.transforms.ParDo;
3940
import org.apache.beam.sdk.transforms.Redistribute;
4041
import org.apache.beam.sdk.transforms.SerializableFunction;
42+
import org.apache.beam.sdk.transforms.Values;
43+
import org.apache.beam.sdk.transforms.WithKeys;
4144
import org.apache.beam.sdk.transforms.errorhandling.BadRecord;
4245
import org.apache.beam.sdk.transforms.errorhandling.BadRecordRouter;
4346
import org.apache.beam.sdk.transforms.errorhandling.BadRecordRouter.ThrowingBadRecordRouter;
@@ -50,6 +53,7 @@
5053
import org.apache.beam.sdk.values.PCollectionList;
5154
import org.apache.beam.sdk.values.PCollectionTuple;
5255
import org.apache.beam.sdk.values.TupleTag;
56+
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.hash.Hashing;
5357
import org.joda.time.Duration;
5458

5559
/** This {@link PTransform} manages loads into BigQuery using the Storage API. */
@@ -379,12 +383,32 @@ public WriteResult expandUntriggered(
379383
PCollection<KV<DestinationT, StorageApiWritePayload>> successfulConvertedRows =
380384
convertMessagesResult.get(successfulConvertedRowsTag);
381385

382-
if (numShards > 0) {
386+
if (numShards > 0 && input.isBounded() == PCollection.IsBounded.UNBOUNDED) {
383387
successfulConvertedRows =
384388
successfulConvertedRows.apply(
385389
"ResdistibuteNumShards",
386390
Redistribute.<KV<DestinationT, StorageApiWritePayload>>arbitrarily()
387391
.withNumBuckets(numShards));
392+
} else if (numShards > 0 && input.isBounded() == PCollection.IsBounded.BOUNDED) {
393+
successfulConvertedRows =
394+
successfulConvertedRows
395+
.apply(
396+
"Add shard",
397+
WithKeys.of(
398+
(SerializableFunction<KV<DestinationT, StorageApiWritePayload>, Integer>)
399+
kv ->
400+
Math.abs(
401+
Hashing.murmur3_32_fixed()
402+
.hashString(
403+
dynamicDestinations
404+
.getTable(kv.getKey())
405+
.getShortTableUrn(),
406+
StandardCharsets.UTF_8)
407+
.asInt()
408+
^ ThreadLocalRandom.current().nextInt(numShards))
409+
% numShards))
410+
.apply("RedistributeNumShards", Redistribute.byKey())
411+
.apply("Remove shard", Values.create());
388412
}
389413

390414
PCollectionTuple writeRecordsResult =

0 commit comments

Comments
 (0)