This is an automated email from the ASF dual-hosted git repository. bertty pushed a commit to branch WAYANG-225 in repository https://gitbox.apache.org/repos/asf/incubator-wayang.git
commit acb33226a9c5d52ddb940165b6da7728d7e28a1b Author: Bertty Contreras-Rojas <[email protected]> AuthorDate: Mon May 2 14:36:00 2022 +0200 [WAYANG-#225] Correction in distinct Operator in platform Apache Flink Signed-off-by: bertty <[email protected]> --- .../java/org/apache/wayang/flink/compiler/KeySelectorDistinct.java | 3 ++- .../apache/wayang/flink/operators/FlinkDistinctOperatorTest.java | 7 +++---- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/wayang-platforms/wayang-flink/code/main/java/org/apache/wayang/flink/compiler/KeySelectorDistinct.java b/wayang-platforms/wayang-flink/code/main/java/org/apache/wayang/flink/compiler/KeySelectorDistinct.java index 732c7453..38d33f8e 100644 --- a/wayang-platforms/wayang-flink/code/main/java/org/apache/wayang/flink/compiler/KeySelectorDistinct.java +++ b/wayang-platforms/wayang-flink/code/main/java/org/apache/wayang/flink/compiler/KeySelectorDistinct.java @@ -18,6 +18,7 @@ package org.apache.wayang.flink.compiler; +import java.io.IOException; import org.apache.flink.api.java.functions.KeySelector; import java.io.ByteArrayOutputStream; @@ -36,7 +37,7 @@ public class KeySelectorDistinct<T> implements KeySelector<T, String>, Serializa ObjectOutputStream objStream = new ObjectOutputStream(b); objStream.writeObject(value); return Base64.getEncoder().encodeToString(b.toByteArray()); - }finally { + }catch (IOException e) { return ""; } } diff --git a/wayang-platforms/wayang-flink/code/test/java/org/apache/wayang/flink/operators/FlinkDistinctOperatorTest.java b/wayang-platforms/wayang-flink/code/test/java/org/apache/wayang/flink/operators/FlinkDistinctOperatorTest.java index 6edb68ea..80f749f7 100644 --- a/wayang-platforms/wayang-flink/code/test/java/org/apache/wayang/flink/operators/FlinkDistinctOperatorTest.java +++ b/wayang-platforms/wayang-flink/code/test/java/org/apache/wayang/flink/operators/FlinkDistinctOperatorTest.java @@ -54,11 +54,10 @@ public class FlinkDistinctOperatorTest extends FlinkOperatorTestBase{ // Verify the outcome. final List<Integer> result = ((DataSetChannel.Instance) outputs[0]).<Integer>provideDataSet().collect(); - for(Object e : result){ - System.out.println(e); - } + Assert.assertEquals(4, result.size()); - Assert.assertEquals(Arrays.asList(0, 1, 6, 2), result); + result.sort((a, b) -> a - b); + Assert.assertEquals(Arrays.asList(0, 1, 2, 6), result); } }
