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);
 
     }
 }

Reply via email to