diff --git a/flink/v1.13/flink/src/test/java/org/apache/iceberg/flink/sink/TestFlinkIcebergSink.java b/flink/v1.13/flink/src/test/java/org/apache/iceberg/flink/sink/TestFlinkIcebergSink.java index 99e8b7519a1d..04b085051fec 100644 --- a/flink/v1.13/flink/src/test/java/org/apache/iceberg/flink/sink/TestFlinkIcebergSink.java +++ b/flink/v1.13/flink/src/test/java/org/apache/iceberg/flink/sink/TestFlinkIcebergSink.java @@ -49,7 +49,6 @@ import org.junit.Assert; import org.junit.Before; import org.junit.ClassRule; -import org.junit.Ignore; import org.junit.Test; import org.junit.rules.TemporaryFolder; import org.junit.runner.RunWith; @@ -265,7 +264,6 @@ public void testShuffleByPartitionWithSchema() throws Exception { } @Test - @Ignore // Ignored as one DAG completing first can cause an infinite checkpoint loop in the other and CI timeouts public void testTwoSinksInDisjointedDAG() throws Exception { Map props = ImmutableMap.of(TableProperties.DEFAULT_FILE_FORMAT, format.name()); @@ -287,7 +285,7 @@ public void testTwoSinksInDisjointedDAG() throws Exception { env.getConfig().disableAutoGeneratedUIDs(); List leftRows = createRows("left-"); - DataStream leftStream = env.addSource(createBoundedSource(leftRows), ROW_TYPE_INFO) + DataStream leftStream = env.fromCollection(leftRows, ROW_TYPE_INFO) .name("leftCustomSource") .uid("leftCustomSource"); FlinkSink.forRow(leftStream, SimpleDataUtil.FLINK_SCHEMA) @@ -299,7 +297,7 @@ public void testTwoSinksInDisjointedDAG() throws Exception { .append(); List rightRows = createRows("right-"); - DataStream rightStream = env.addSource(createBoundedSource(rightRows), ROW_TYPE_INFO) + DataStream rightStream = env.fromCollection(rightRows, ROW_TYPE_INFO) .name("rightCustomSource") .uid("rightCustomSource"); FlinkSink.forRow(rightStream, SimpleDataUtil.FLINK_SCHEMA)