This is an automated email from the ASF dual-hosted git repository.
Amar3tto pushed a commit to branch fix-java-postcommit
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/fix-java-postcommit by this
push:
new d847c2e3272 Normalize output
d847c2e3272 is described below
commit d847c2e32727a5c850b7c6b1551e97e329f1cab5
Author: Vitaly Terentyev <[email protected]>
AuthorDate: Thu Sep 10 13:43:47 2026 +0400
Normalize output
---
.../extensions/gcp/transforms/GcpGroupByKeyIT.java | 45 +++++++++++++++++++---
1 file changed, 39 insertions(+), 6 deletions(-)
diff --git
a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/transforms/GcpGroupByKeyIT.java
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/transforms/GcpGroupByKeyIT.java
index e8d3466cf09..b99ecf887ad 100644
---
a/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/transforms/GcpGroupByKeyIT.java
+++
b/sdks/java/extensions/google-cloud-platform-core/src/test/java/org/apache/beam/sdk/extensions/gcp/transforms/GcpGroupByKeyIT.java
@@ -28,7 +28,9 @@ import com.google.cloud.secretmanager.v1.SecretPayload;
import com.google.protobuf.ByteString;
import java.io.IOException;
import java.security.SecureRandom;
+import java.util.ArrayList;
import java.util.Arrays;
+import java.util.Collections;
import java.util.List;
import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.coders.KvCoder;
@@ -41,10 +43,13 @@ import org.apache.beam.sdk.testing.TestPipeline;
import org.apache.beam.sdk.transforms.Combine;
import org.apache.beam.sdk.transforms.Create;
import org.apache.beam.sdk.transforms.GroupByKey;
+import org.apache.beam.sdk.transforms.MapElements;
import org.apache.beam.sdk.transforms.Redistribute;
import org.apache.beam.sdk.transforms.Sum;
import org.apache.beam.sdk.values.KV;
import org.apache.beam.sdk.values.PCollection;
+import org.apache.beam.sdk.values.TypeDescriptor;
+import org.apache.beam.sdk.values.TypeDescriptors;
import org.junit.AfterClass;
import org.junit.BeforeClass;
import org.junit.Rule;
@@ -141,6 +146,22 @@ public class GcpGroupByKeyIT {
}
}
+ private static <K, V extends Comparable<? super V>> PCollection<KV<K,
List<V>>> sortGroupedValues(
+ PCollection<KV<K, Iterable<V>>> input,
+ TypeDescriptor<K> keyType,
+ TypeDescriptor<V> valueType) {
+
+ return input.apply(
+ MapElements.into(TypeDescriptors.kvs(keyType,
TypeDescriptors.lists(valueType)))
+ .via(
+ kv -> {
+ List<V> values = new ArrayList<>();
+ kv.getValue().forEach(values::add);
+ Collections.sort(values);
+ return KV.of(kv.getKey(), values);
+ }));
+ }
+
@Test
public void testGroupByKeyWithValidGcpSecretOption() throws Exception {
if (gcpSecretVersionName == null) {
@@ -165,9 +186,13 @@ public class GcpGroupByKeyIT {
Create.of(ungroupedPairs)
.withCoder(KvCoder.of(StringUtf8Coder.of(),
VarIntCoder.of())));
- PCollection<KV<String, Iterable<Integer>>> output =
input.apply(GroupByKey.create());
+ PCollection<KV<String, List<Integer>>> normalizedOutput =
+ sortGroupedValues(
+ input.apply(GroupByKey.create()),
+ TypeDescriptors.strings(),
+ TypeDescriptors.integers());
- PAssert.that(output)
+ PAssert.that(normalizedOutput)
.containsInAnyOrder(
KV.of("k1", Arrays.asList(3, 4)),
KV.of("k5", Arrays.asList(Integer.MAX_VALUE, Integer.MIN_VALUE)),
@@ -201,9 +226,13 @@ public class GcpGroupByKeyIT {
Create.of(ungroupedPairs)
.withCoder(KvCoder.of(StringUtf8Coder.of(),
VarIntCoder.of())));
- PCollection<KV<String, Iterable<Integer>>> output =
input.apply(GroupByKey.create());
+ PCollection<KV<String, List<Integer>>> normalizedOutput =
+ sortGroupedValues(
+ input.apply(GroupByKey.create()),
+ TypeDescriptors.strings(),
+ TypeDescriptors.integers());
- PAssert.that(output)
+ PAssert.that(normalizedOutput)
.containsInAnyOrder(
KV.of("k1", Arrays.asList(3, 4)),
KV.of("k5", Arrays.asList(Integer.MAX_VALUE, Integer.MIN_VALUE)),
@@ -240,9 +269,13 @@ public class GcpGroupByKeyIT {
Create.of(ungroupedPairs)
.withCoder(KvCoder.of(StringUtf8Coder.of(),
VarIntCoder.of())));
- PCollection<KV<String, Iterable<Integer>>> output =
input.apply(GroupByKey.create());
+ PCollection<KV<String, List<Integer>>> normalizedOutput =
+ sortGroupedValues(
+ input.apply(GroupByKey.create()),
+ TypeDescriptors.strings(),
+ TypeDescriptors.integers());
- PAssert.that(output)
+ PAssert.that(normalizedOutput)
.containsInAnyOrder(
KV.of("k1", Arrays.asList(3, 4)),
KV.of("k5", Arrays.asList(Integer.MAX_VALUE, Integer.MIN_VALUE)),