You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@flink.apache.org by tr...@apache.org on 2018/04/02 15:13:19 UTC
[2/5] flink git commit: [FLINK-9121] [flip6] Remove Flip6 prefixes
and other references
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/ResourceManagerHATest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/ResourceManagerHATest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/ResourceManagerHATest.java
index a3d17f2..7843d7d 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/ResourceManagerHATest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/ResourceManagerHATest.java
@@ -31,7 +31,7 @@ import org.apache.flink.runtime.rpc.RpcService;
import org.apache.flink.runtime.rpc.TestingRpcService;
import org.apache.flink.runtime.testingUtils.TestingUtils;
import org.apache.flink.runtime.util.TestingFatalErrorHandler;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.junit.Assert;
import org.junit.Test;
@@ -45,7 +45,7 @@ import static org.mockito.Mockito.mock;
/**
* resourceManager HA test, including grant leadership and revoke leadership
*/
-@Category(Flip6.class)
+@Category(New.class)
public class ResourceManagerHATest extends TestLogger {
@Test
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/ResourceManagerJobMasterTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/ResourceManagerJobMasterTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/ResourceManagerJobMasterTest.java
index 7e51a2c..312adaf 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/ResourceManagerJobMasterTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/ResourceManagerJobMasterTest.java
@@ -43,7 +43,7 @@ import org.apache.flink.runtime.registration.RegistrationResponse;
import org.apache.flink.runtime.rpc.exceptions.FencingTokenException;
import org.apache.flink.runtime.testingUtils.TestingUtils;
import org.apache.flink.runtime.util.TestingFatalErrorHandler;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.ExceptionUtils;
import org.apache.flink.util.TestLogger;
import org.junit.After;
@@ -60,7 +60,7 @@ import static org.junit.Assert.assertTrue;
import static org.junit.Assert.fail;
import static org.mockito.Mockito.*;
-@Category(Flip6.class)
+@Category(New.class)
public class ResourceManagerJobMasterTest extends TestLogger {
private TestingRpcService rpcService;
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/ResourceManagerTaskExecutorTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/ResourceManagerTaskExecutorTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/ResourceManagerTaskExecutorTest.java
index e4b93d3..d49ba57 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/ResourceManagerTaskExecutorTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/ResourceManagerTaskExecutorTest.java
@@ -40,7 +40,7 @@ import org.apache.flink.runtime.taskexecutor.TaskExecutorGateway;
import org.apache.flink.runtime.taskexecutor.TaskExecutorRegistrationSuccess;
import org.apache.flink.runtime.testingUtils.TestingUtils;
import org.apache.flink.runtime.util.TestingFatalErrorHandler;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.ExceptionUtils;
import org.apache.flink.util.TestLogger;
import org.junit.After;
@@ -61,7 +61,7 @@ import static org.junit.Assert.fail;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
-@Category(Flip6.class)
+@Category(New.class)
public class ResourceManagerTaskExecutorTest extends TestLogger {
private final Time timeout = Time.seconds(10L);
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/ResourceManagerTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/ResourceManagerTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/ResourceManagerTest.java
index f5fa899..7dab685 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/ResourceManagerTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/ResourceManagerTest.java
@@ -35,7 +35,7 @@ import org.apache.flink.runtime.taskexecutor.TaskExecutorGateway;
import org.apache.flink.runtime.taskexecutor.TestingTaskExecutorGateway;
import org.apache.flink.runtime.testingUtils.TestingUtils;
import org.apache.flink.runtime.util.TestingFatalErrorHandler;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.junit.After;
@@ -48,7 +48,7 @@ import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
-@Category(Flip6.class)
+@Category(New.class)
public class ResourceManagerTest extends TestLogger {
private TestingRpcService rpcService;
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/slotmanager/SlotManagerTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/slotmanager/SlotManagerTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/slotmanager/SlotManagerTest.java
index 90ed164..f504d77 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/slotmanager/SlotManagerTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/slotmanager/SlotManagerTest.java
@@ -39,7 +39,7 @@ import org.apache.flink.runtime.taskexecutor.TaskExecutorGateway;
import org.apache.flink.runtime.taskexecutor.TestingTaskExecutorGateway;
import org.apache.flink.runtime.taskexecutor.exceptions.SlotAllocationException;
import org.apache.flink.runtime.testingUtils.TestingUtils;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.junit.Test;
@@ -74,7 +74,7 @@ import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
-@Category(Flip6.class)
+@Category(New.class)
public class SlotManagerTest extends TestLogger {
/**
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/slotmanager/SlotProtocolTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/slotmanager/SlotProtocolTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/slotmanager/SlotProtocolTest.java
index d47ac33..b3e5e91 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/slotmanager/SlotProtocolTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/resourcemanager/slotmanager/SlotProtocolTest.java
@@ -33,7 +33,7 @@ import org.apache.flink.runtime.taskexecutor.SlotReport;
import org.apache.flink.runtime.taskexecutor.SlotStatus;
import org.apache.flink.runtime.taskexecutor.TaskExecutorGateway;
import org.apache.flink.runtime.testingUtils.TestingUtils;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.ExecutorUtils;
import org.apache.flink.util.TestLogger;
@@ -54,7 +54,7 @@ import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.timeout;
import static org.mockito.Mockito.verify;
-@Category(Flip6.class)
+@Category(New.class)
public class SlotProtocolTest extends TestLogger {
private static final long timeout = 10000L;
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rest/RestServerEndpointITCase.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/RestServerEndpointITCase.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/RestServerEndpointITCase.java
index e510798..88fdeb8 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/RestServerEndpointITCase.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/RestServerEndpointITCase.java
@@ -43,7 +43,7 @@ import org.apache.flink.runtime.rpc.RpcUtils;
import org.apache.flink.runtime.testingUtils.TestingUtils;
import org.apache.flink.runtime.webmonitor.RestfulGateway;
import org.apache.flink.runtime.webmonitor.retriever.GatewayRetriever;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.ExceptionUtils;
import org.apache.flink.util.Preconditions;
import org.apache.flink.util.TestLogger;
@@ -94,7 +94,7 @@ import static org.mockito.Mockito.when;
/**
* IT cases for {@link RestClient} and {@link RestServerEndpoint}.
*/
-@Category(Flip6.class)
+@Category(New.class)
public class RestServerEndpointITCase extends TestLogger {
private static final JobID PATH_JOB_ID = new JobID();
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rest/RestServerEndpointTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/RestServerEndpointTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/RestServerEndpointTest.java
index de32883..bdb3f5d 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/RestServerEndpointTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/RestServerEndpointTest.java
@@ -18,7 +18,7 @@
package org.apache.flink.runtime.rest;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.junit.Assume;
@@ -44,7 +44,7 @@ import static org.junit.Assert.fail;
/**
* Test cases for the {@link RestServerEndpoint}.
*/
-@Category(Flip6.class)
+@Category(New.class)
public class RestServerEndpointTest extends TestLogger {
@Rule
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/BlobServerPortHandlerTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/BlobServerPortHandlerTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/BlobServerPortHandlerTest.java
index 0c2c2ed..6c81313 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/BlobServerPortHandlerTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/BlobServerPortHandlerTest.java
@@ -28,7 +28,7 @@ import org.apache.flink.runtime.rest.messages.EmptyMessageParameters;
import org.apache.flink.runtime.rest.messages.EmptyRequestBody;
import org.apache.flink.runtime.rpc.RpcUtils;
import org.apache.flink.runtime.webmonitor.retriever.GatewayRetriever;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.apache.flink.shaded.netty4.io.netty.handler.codec.http.HttpResponseStatus;
@@ -48,7 +48,7 @@ import static org.mockito.Mockito.when;
/**
* Tests for the {@link BlobServerPortHandler}.
*/
-@Category(Flip6.class)
+@Category(New.class)
public class BlobServerPortHandlerTest extends TestLogger {
private static final int PORT = 64;
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/JobExecutionResultHandlerTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/JobExecutionResultHandlerTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/JobExecutionResultHandlerTest.java
index d69c5d5..e32e16a 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/JobExecutionResultHandlerTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/JobExecutionResultHandlerTest.java
@@ -33,7 +33,7 @@ import org.apache.flink.runtime.rest.messages.JobMessageParameters;
import org.apache.flink.runtime.rest.messages.job.JobExecutionResultResponseBody;
import org.apache.flink.runtime.rest.messages.queue.QueueStatus;
import org.apache.flink.runtime.webmonitor.TestingRestfulGateway;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.ExceptionUtils;
import org.apache.flink.util.TestLogger;
@@ -57,7 +57,7 @@ import static org.junit.Assert.fail;
/**
* Tests for {@link JobExecutionResultHandler}.
*/
-@Category(Flip6.class)
+@Category(New.class)
public class JobExecutionResultHandlerTest extends TestLogger {
private static final JobID TEST_JOB_ID = new JobID();
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/JobSubmitHandlerTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/JobSubmitHandlerTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/JobSubmitHandlerTest.java
index 2428d38..ac09186 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/JobSubmitHandlerTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/JobSubmitHandlerTest.java
@@ -28,7 +28,7 @@ import org.apache.flink.runtime.rest.messages.EmptyMessageParameters;
import org.apache.flink.runtime.rest.messages.job.JobSubmitRequestBody;
import org.apache.flink.runtime.rpc.RpcUtils;
import org.apache.flink.runtime.webmonitor.retriever.GatewayRetriever;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.apache.flink.shaded.netty4.io.netty.handler.codec.http.HttpResponseStatus;
@@ -47,7 +47,7 @@ import static org.mockito.Mockito.when;
/**
* Tests for the {@link JobSubmitHandler}.
*/
-@Category(Flip6.class)
+@Category(New.class)
public class JobSubmitHandlerTest extends TestLogger {
@Test
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/metrics/AbstractMetricsHandlerTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/metrics/AbstractMetricsHandlerTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/metrics/AbstractMetricsHandlerTest.java
index 80a8759..0d018bc 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/metrics/AbstractMetricsHandlerTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/metrics/AbstractMetricsHandlerTest.java
@@ -35,7 +35,7 @@ import org.apache.flink.runtime.rest.messages.job.metrics.Metric;
import org.apache.flink.runtime.rest.messages.job.metrics.MetricCollectionResponseBody;
import org.apache.flink.runtime.rest.messages.job.metrics.MetricsFilterParameter;
import org.apache.flink.runtime.webmonitor.retriever.GatewayRetriever;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.junit.Before;
@@ -61,7 +61,7 @@ import static org.mockito.Mockito.when;
/**
* Tests for {@link AbstractMetricsHandler}.
*/
-@Category(Flip6.class)
+@Category(New.class)
public class AbstractMetricsHandlerTest extends TestLogger {
private static final String TEST_METRIC_NAME = "test_counter";
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/metrics/MetricsHandlerTestBase.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/metrics/MetricsHandlerTestBase.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/metrics/MetricsHandlerTestBase.java
index 92ee185..bf37cc9 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/metrics/MetricsHandlerTestBase.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/handler/job/metrics/MetricsHandlerTestBase.java
@@ -29,7 +29,7 @@ import org.apache.flink.runtime.rest.messages.EmptyRequestBody;
import org.apache.flink.runtime.rest.messages.job.metrics.Metric;
import org.apache.flink.runtime.rest.messages.job.metrics.MetricCollectionResponseBody;
import org.apache.flink.runtime.webmonitor.retriever.GatewayRetriever;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.junit.Before;
@@ -51,7 +51,7 @@ import static org.mockito.Mockito.when;
/**
* Unit test base class for subclasses of {@link AbstractMetricsHandler}.
*/
-@Category(Flip6.class)
+@Category(New.class)
public abstract class MetricsHandlerTestBase<T extends AbstractMetricsHandler> extends TestLogger {
private static final String TEST_METRIC_NAME = "test_counter";
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/MessageParametersTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/MessageParametersTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/MessageParametersTest.java
index 8d73231..5bc9b60 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/MessageParametersTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/MessageParametersTest.java
@@ -19,7 +19,7 @@
package org.apache.flink.runtime.rest.messages;
import org.apache.flink.api.common.JobID;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.junit.Assert;
@@ -32,7 +32,7 @@ import java.util.Collections;
/**
* Tests for {@link MessageParameters}.
*/
-@Category(Flip6.class)
+@Category(New.class)
public class MessageParametersTest extends TestLogger {
@Test
public void testResolveUrl() {
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/RestRequestMarshallingTestBase.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/RestRequestMarshallingTestBase.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/RestRequestMarshallingTestBase.java
index 4dc81e8..c5ae563 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/RestRequestMarshallingTestBase.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/RestRequestMarshallingTestBase.java
@@ -19,7 +19,7 @@
package org.apache.flink.runtime.rest.messages;
import org.apache.flink.runtime.rest.util.RestMapperUtils;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper;
@@ -31,7 +31,7 @@ import org.junit.experimental.categories.Category;
/**
* Test base for verifying that marshalling / unmarshalling REST {@link RequestBody}s work properly.
*/
-@Category(Flip6.class)
+@Category(New.class)
public abstract class RestRequestMarshallingTestBase<R extends RequestBody> extends TestLogger {
/**
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/RestResponseMarshallingTestBase.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/RestResponseMarshallingTestBase.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/RestResponseMarshallingTestBase.java
index 7442dcb..fa7ec4c 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/RestResponseMarshallingTestBase.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/RestResponseMarshallingTestBase.java
@@ -19,7 +19,7 @@
package org.apache.flink.runtime.rest.messages;
import org.apache.flink.runtime.rest.util.RestMapperUtils;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.JavaType;
@@ -35,7 +35,7 @@ import java.util.Collections;
/**
* Test base for verifying that marshalling / unmarshalling REST {@link ResponseBody}s work properly.
*/
-@Category(Flip6.class)
+@Category(New.class)
public abstract class RestResponseMarshallingTestBase<R extends ResponseBody> extends TestLogger {
/**
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/AbstractMetricsHeadersTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/AbstractMetricsHeadersTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/AbstractMetricsHeadersTest.java
index 5228e0e..0ea2d37 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/AbstractMetricsHeadersTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/AbstractMetricsHeadersTest.java
@@ -21,7 +21,7 @@ package org.apache.flink.runtime.rest.messages.job.metrics;
import org.apache.flink.runtime.rest.HttpMethodWrapper;
import org.apache.flink.runtime.rest.messages.EmptyMessageParameters;
import org.apache.flink.runtime.rest.messages.EmptyRequestBody;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.apache.flink.shaded.netty4.io.netty.handler.codec.http.HttpResponseStatus;
@@ -36,7 +36,7 @@ import static org.junit.Assert.assertThat;
/**
* Tests for {@link AbstractMetricsHeaders}.
*/
-@Category(Flip6.class)
+@Category(New.class)
public class AbstractMetricsHeadersTest extends TestLogger {
private AbstractMetricsHeaders<EmptyMessageParameters> metricsHandlerHeaders;
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/JobManagerMetricsHeadersTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/JobManagerMetricsHeadersTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/JobManagerMetricsHeadersTest.java
index e87d3de..74c7603 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/JobManagerMetricsHeadersTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/JobManagerMetricsHeadersTest.java
@@ -18,7 +18,7 @@
package org.apache.flink.runtime.rest.messages.job.metrics;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.junit.Test;
@@ -31,7 +31,7 @@ import static org.junit.Assert.assertThat;
/**
* Tests for {@link JobManagerMetricsHeaders}.
*/
-@Category(Flip6.class)
+@Category(New.class)
public class JobManagerMetricsHeadersTest extends TestLogger {
private final JobManagerMetricsHeaders jobManagerMetricsHeaders =
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/JobMetricsHeadersTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/JobMetricsHeadersTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/JobMetricsHeadersTest.java
index 515c7c4..5811442 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/JobMetricsHeadersTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/JobMetricsHeadersTest.java
@@ -19,7 +19,7 @@
package org.apache.flink.runtime.rest.messages.job.metrics;
import org.apache.flink.runtime.rest.messages.JobIDPathParameter;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.junit.Test;
@@ -32,7 +32,7 @@ import static org.junit.Assert.assertThat;
/**
* Tests for {@link JobMetricsHeaders}.
*/
-@Category(Flip6.class)
+@Category(New.class)
public class JobMetricsHeadersTest extends TestLogger {
private final JobMetricsHeaders jobMetricsHeaders = JobMetricsHeaders.getInstance();
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/JobVertexMetricsHeadersTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/JobVertexMetricsHeadersTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/JobVertexMetricsHeadersTest.java
index f20abdb..5dd2567 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/JobVertexMetricsHeadersTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/JobVertexMetricsHeadersTest.java
@@ -20,7 +20,7 @@ package org.apache.flink.runtime.rest.messages.job.metrics;
import org.apache.flink.runtime.rest.messages.JobIDPathParameter;
import org.apache.flink.runtime.rest.messages.JobVertexIdPathParameter;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.junit.Test;
@@ -33,7 +33,7 @@ import static org.junit.Assert.assertThat;
/**
* Tests for {@link JobVertexMetricsHeaders}.
*/
-@Category(Flip6.class)
+@Category(New.class)
public class JobVertexMetricsHeadersTest extends TestLogger {
private final JobVertexMetricsHeaders jobVertexMetricsHeaders = JobVertexMetricsHeaders
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/MetricsFilterParameterTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/MetricsFilterParameterTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/MetricsFilterParameterTest.java
index b13cb01..f30132d 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/MetricsFilterParameterTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/MetricsFilterParameterTest.java
@@ -18,7 +18,7 @@
package org.apache.flink.runtime.rest.messages.job.metrics;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.junit.Before;
@@ -32,7 +32,7 @@ import static org.junit.Assert.assertThat;
/**
* Tests for {@link MetricsFilterParameter}.
*/
-@Category(Flip6.class)
+@Category(New.class)
public class MetricsFilterParameterTest extends TestLogger {
private MetricsFilterParameter metricsFilterParameter;
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/SubtaskMetricsHeadersTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/SubtaskMetricsHeadersTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/SubtaskMetricsHeadersTest.java
index 0f82465..f99c423 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/SubtaskMetricsHeadersTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/SubtaskMetricsHeadersTest.java
@@ -21,7 +21,7 @@ package org.apache.flink.runtime.rest.messages.job.metrics;
import org.apache.flink.runtime.rest.messages.JobIDPathParameter;
import org.apache.flink.runtime.rest.messages.JobVertexIdPathParameter;
import org.apache.flink.runtime.rest.messages.SubtaskIndexPathParameter;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.junit.Test;
@@ -34,7 +34,7 @@ import static org.junit.Assert.assertThat;
/**
* Tests for {@link SubtaskMetricsHeaders}.
*/
-@Category(Flip6.class)
+@Category(New.class)
public class SubtaskMetricsHeadersTest extends TestLogger {
private final SubtaskMetricsHeaders subtaskMetricsHeaders = SubtaskMetricsHeaders.getInstance();
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/TaskManagerMetricsHeadersTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/TaskManagerMetricsHeadersTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/TaskManagerMetricsHeadersTest.java
index 477e9f8..66b0d8e 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/TaskManagerMetricsHeadersTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/job/metrics/TaskManagerMetricsHeadersTest.java
@@ -19,7 +19,7 @@
package org.apache.flink.runtime.rest.messages.job.metrics;
import org.apache.flink.runtime.rest.messages.taskmanager.TaskManagerIdPathParameter;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.junit.Test;
@@ -32,7 +32,7 @@ import static org.junit.Assert.assertThat;
/**
* Tests for {@link TaskManagerMetricsHeaders}.
*/
-@Category(Flip6.class)
+@Category(New.class)
public class TaskManagerMetricsHeadersTest extends TestLogger {
private final TaskManagerMetricsHeaders taskManagerMetricsHeaders =
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/taskmanager/TaskManagerIdPathParameterTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/taskmanager/TaskManagerIdPathParameterTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/taskmanager/TaskManagerIdPathParameterTest.java
index 379fe4c..9c62915 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/taskmanager/TaskManagerIdPathParameterTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rest/messages/taskmanager/TaskManagerIdPathParameterTest.java
@@ -19,7 +19,7 @@
package org.apache.flink.runtime.rest.messages.taskmanager;
import org.apache.flink.runtime.clusterframework.types.ResourceID;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.junit.Before;
@@ -33,7 +33,7 @@ import static org.junit.Assert.assertTrue;
/**
* Tests for {@link TaskManagerIdPathParameter}.
*/
-@Category(Flip6.class)
+@Category(New.class)
public class TaskManagerIdPathParameterTest extends TestLogger {
private TaskManagerIdPathParameter taskManagerIdPathParameter;
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/AsyncCallsTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/AsyncCallsTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/AsyncCallsTest.java
index 1f9d9e3..5f5ba44 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/AsyncCallsTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/AsyncCallsTest.java
@@ -25,7 +25,7 @@ import org.apache.flink.runtime.concurrent.FutureUtils;
import org.apache.flink.runtime.messages.Acknowledge;
import org.apache.flink.runtime.rpc.akka.AkkaRpcService;
import org.apache.flink.runtime.rpc.exceptions.FencingTokenException;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.ExceptionUtils;
import org.apache.flink.util.TestLogger;
@@ -52,7 +52,7 @@ import static org.junit.Assert.assertThat;
import static org.junit.Assert.assertTrue;
import static org.junit.Assert.fail;
-@Category(Flip6.class)
+@Category(New.class)
public class AsyncCallsTest extends TestLogger {
// ------------------------------------------------------------------------
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/FencedRpcEndpointTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/FencedRpcEndpointTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/FencedRpcEndpointTest.java
index f488308..00ff3d7 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/FencedRpcEndpointTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/FencedRpcEndpointTest.java
@@ -23,7 +23,7 @@ import org.apache.flink.core.testutils.OneShotLatch;
import org.apache.flink.runtime.messages.Acknowledge;
import org.apache.flink.runtime.rpc.exceptions.FencingTokenException;
import org.apache.flink.runtime.rpc.exceptions.RpcException;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.ExceptionUtils;
import org.apache.flink.util.FlinkException;
import org.apache.flink.util.TestLogger;
@@ -46,7 +46,7 @@ import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertTrue;
import static org.junit.Assert.fail;
-@Category(Flip6.class)
+@Category(New.class)
public class FencedRpcEndpointTest extends TestLogger {
private static final Time timeout = Time.seconds(10L);
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/RpcConnectionTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/RpcConnectionTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/RpcConnectionTest.java
index 017c1f5..78137e8 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/RpcConnectionTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/RpcConnectionTest.java
@@ -27,7 +27,7 @@ import org.apache.flink.runtime.concurrent.FutureUtils;
import org.apache.flink.runtime.rpc.akka.AkkaRpcService;
import org.apache.flink.runtime.rpc.exceptions.RpcConnectionException;
import org.apache.flink.runtime.taskexecutor.TaskExecutorGateway;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import akka.actor.Terminated;
@@ -49,7 +49,7 @@ import static org.junit.Assert.*;
* This test validates that the RPC service gives a good message when it cannot
* connect to an RpcEndpoint.
*/
-@Category(Flip6.class)
+@Category(New.class)
public class RpcConnectionTest extends TestLogger {
@Test
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/RpcEndpointTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/RpcEndpointTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/RpcEndpointTest.java
index b5add60..1ca3949 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/RpcEndpointTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/RpcEndpointTest.java
@@ -22,7 +22,7 @@ import org.apache.flink.api.common.time.Time;
import org.apache.flink.runtime.akka.AkkaUtils;
import org.apache.flink.runtime.concurrent.FutureUtils;
import org.apache.flink.runtime.rpc.akka.AkkaRpcService;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import akka.actor.ActorSystem;
@@ -42,7 +42,7 @@ import static org.junit.Assert.fail;
/**
* Tests for the RpcEndpoint and its self gateways.
*/
-@Category(Flip6.class)
+@Category(New.class)
public class RpcEndpointTest extends TestLogger {
private static final Time TIMEOUT = Time.seconds(10L);
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/akka/AkkaRpcActorTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/akka/AkkaRpcActorTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/akka/AkkaRpcActorTest.java
index 2530bce..a92235c 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/akka/AkkaRpcActorTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/akka/AkkaRpcActorTest.java
@@ -28,7 +28,7 @@ import org.apache.flink.runtime.rpc.RpcUtils;
import org.apache.flink.runtime.rpc.TestingRpcService;
import org.apache.flink.runtime.rpc.akka.exceptions.AkkaRpcException;
import org.apache.flink.runtime.rpc.exceptions.RpcConnectionException;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.FlinkException;
import org.apache.flink.util.Preconditions;
import org.apache.flink.util.TestLogger;
@@ -51,7 +51,7 @@ import static org.junit.Assert.assertThat;
import static org.junit.Assert.assertTrue;
import static org.junit.Assert.fail;
-@Category(Flip6.class)
+@Category(New.class)
public class AkkaRpcActorTest extends TestLogger {
// ------------------------------------------------------------------------
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/akka/AkkaRpcServiceTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/akka/AkkaRpcServiceTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/akka/AkkaRpcServiceTest.java
index d92e496..d73ee40 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/akka/AkkaRpcServiceTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/akka/AkkaRpcServiceTest.java
@@ -23,7 +23,7 @@ import org.apache.flink.core.testutils.OneShotLatch;
import org.apache.flink.runtime.akka.AkkaUtils;
import org.apache.flink.runtime.concurrent.FutureUtils;
import org.apache.flink.runtime.concurrent.ScheduledExecutor;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import akka.actor.ActorSystem;
@@ -46,7 +46,7 @@ import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertTrue;
import static org.junit.Assert.fail;
-@Category(Flip6.class)
+@Category(New.class)
public class AkkaRpcServiceTest extends TestLogger {
// ------------------------------------------------------------------------
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/akka/MainThreadValidationTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/akka/MainThreadValidationTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/akka/MainThreadValidationTest.java
index 6dacdfd..34f8eb8 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/akka/MainThreadValidationTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/akka/MainThreadValidationTest.java
@@ -25,7 +25,7 @@ import org.apache.flink.runtime.rpc.RpcEndpoint;
import org.apache.flink.runtime.rpc.RpcGateway;
import org.apache.flink.runtime.rpc.RpcService;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.junit.Test;
import org.junit.experimental.categories.Category;
@@ -34,7 +34,7 @@ import java.util.concurrent.CompletableFuture;
import static org.junit.Assert.assertTrue;
-@Category(Flip6.class)
+@Category(New.class)
public class MainThreadValidationTest extends TestLogger {
@Test
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/akka/MessageSerializationTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/akka/MessageSerializationTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/akka/MessageSerializationTest.java
index 6006850..02c0094 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/akka/MessageSerializationTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/rpc/akka/MessageSerializationTest.java
@@ -24,7 +24,7 @@ import org.apache.flink.runtime.concurrent.FutureUtils;
import org.apache.flink.runtime.rpc.RpcEndpoint;
import org.apache.flink.runtime.rpc.RpcGateway;
import org.apache.flink.runtime.rpc.RpcService;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import akka.actor.ActorSystem;
@@ -51,7 +51,7 @@ import static org.junit.Assert.fail;
/**
* Tests that akka rpc invocation messages are properly serialized and errors reported
*/
-@Category(Flip6.class)
+@Category(New.class)
public class MessageSerializationTest extends TestLogger {
private static ActorSystem actorSystem1;
private static ActorSystem actorSystem2;
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/NetworkBufferCalculationTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/NetworkBufferCalculationTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/NetworkBufferCalculationTest.java
index 8dcd642..2116c2f 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/NetworkBufferCalculationTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/NetworkBufferCalculationTest.java
@@ -22,7 +22,7 @@ import org.apache.flink.configuration.TaskManagerOptions;
import org.apache.flink.core.memory.MemoryType;
import org.apache.flink.runtime.state.LocalRecoveryConfig;
import org.apache.flink.runtime.taskmanager.NetworkEnvironmentConfiguration;
-import org.apache.flink.testutils.category.OldAndFlip6;
+import org.apache.flink.testutils.category.LegacyAndNew;
import org.apache.flink.util.TestLogger;
import org.junit.Test;
@@ -35,7 +35,7 @@ import static org.junit.Assert.assertEquals;
/**
* Tests the network buffer calculation from heap size.
*/
-@Category(OldAndFlip6.class)
+@Category(LegacyAndNew.class)
public class NetworkBufferCalculationTest extends TestLogger {
/**
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorITCase.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorITCase.java b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorITCase.java
index fc8337f..885d99f 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorITCase.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorITCase.java
@@ -58,7 +58,7 @@ import org.apache.flink.runtime.taskexecutor.slot.TimerService;
import org.apache.flink.runtime.taskmanager.TaskManagerLocation;
import org.apache.flink.runtime.testingUtils.TestingUtils;
import org.apache.flink.runtime.util.TestingFatalErrorHandler;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.hamcrest.Matchers;
@@ -86,7 +86,7 @@ import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
-@Category(Flip6.class)
+@Category(New.class)
public class TaskExecutorITCase extends TestLogger {
private final Time timeout = Time.seconds(10L);
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorTest.java
index 7aae287..465619e 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorTest.java
@@ -88,7 +88,7 @@ import org.apache.flink.runtime.taskmanager.TaskManagerActions;
import org.apache.flink.runtime.taskmanager.TaskManagerLocation;
import org.apache.flink.runtime.testingUtils.TestingUtils;
import org.apache.flink.runtime.util.TestingFatalErrorHandler;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.FlinkException;
import org.apache.flink.util.SerializedValue;
import org.apache.flink.util.TestLogger;
@@ -139,7 +139,7 @@ import static org.mockito.Mockito.timeout;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
-@Category(Flip6.class)
+@Category(New.class)
public class TaskExecutorTest extends TestLogger {
@Rule
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskManagerServicesConfigurationTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskManagerServicesConfigurationTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskManagerServicesConfigurationTest.java
index 2699a05..9cf76f2 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskManagerServicesConfigurationTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskManagerServicesConfigurationTest.java
@@ -20,7 +20,7 @@ package org.apache.flink.runtime.taskexecutor;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.configuration.TaskManagerOptions;
-import org.apache.flink.testutils.category.OldAndFlip6;
+import org.apache.flink.testutils.category.LegacyAndNew;
import org.apache.flink.util.TestLogger;
import org.junit.Test;
@@ -31,7 +31,7 @@ import static org.junit.Assert.*;
/**
* Unit test for {@link TaskManagerServicesConfiguration}.
*/
-@Category(OldAndFlip6.class)
+@Category(LegacyAndNew.class)
public class TaskManagerServicesConfigurationTest extends TestLogger {
/**
* Verifies that {@link TaskManagerServicesConfiguration#hasNewNetworkBufConf(Configuration)}
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskManagerServicesTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskManagerServicesTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskManagerServicesTest.java
index f6e7b07..097fc61 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskManagerServicesTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskManagerServicesTest.java
@@ -20,7 +20,7 @@ package org.apache.flink.runtime.taskexecutor;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.configuration.TaskManagerOptions;
-import org.apache.flink.testutils.category.OldAndFlip6;
+import org.apache.flink.testutils.category.LegacyAndNew;
import org.apache.flink.util.TestLogger;
import org.junit.Test;
@@ -35,7 +35,7 @@ import static org.junit.Assert.fail;
/**
* Unit test for {@link TaskManagerServices}.
*/
-@Category(OldAndFlip6.class)
+@Category(LegacyAndNew.class)
public class TaskManagerServicesTest extends TestLogger {
/**
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TestingTaskExecutorGateway.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TestingTaskExecutorGateway.java b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TestingTaskExecutorGateway.java
index a22cc6f..43b0be2 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TestingTaskExecutorGateway.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TestingTaskExecutorGateway.java
@@ -34,7 +34,7 @@ import org.apache.flink.runtime.jobmaster.JobMasterId;
import org.apache.flink.runtime.messages.Acknowledge;
import org.apache.flink.runtime.messages.StackTraceSampleResponse;
import org.apache.flink.runtime.resourcemanager.ResourceManagerId;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.Preconditions;
import org.junit.experimental.categories.Category;
@@ -45,7 +45,7 @@ import java.util.function.Consumer;
/**
* Simple {@link TaskExecutorGateway} implementation for testing purposes.
*/
-@Category(Flip6.class)
+@Category(New.class)
public class TestingTaskExecutorGateway implements TaskExecutorGateway {
private final String address;
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/slot/TimerServiceTest.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/slot/TimerServiceTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/slot/TimerServiceTest.java
index ddaa0a6..642c790 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/slot/TimerServiceTest.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/slot/TimerServiceTest.java
@@ -19,7 +19,7 @@
package org.apache.flink.runtime.taskexecutor.slot;
import org.apache.flink.runtime.clusterframework.types.AllocationID;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.TestLogger;
import org.junit.Test;
@@ -39,7 +39,7 @@ import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
-@Category(Flip6.class)
+@Category(New.class)
public class TimerServiceTest extends TestLogger {
/**
* Test all timeouts registered can be unregistered
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-runtime/src/test/java/org/apache/flink/runtime/taskmanager/TaskCancelAsyncProducerConsumerITCase.java
----------------------------------------------------------------------
diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/taskmanager/TaskCancelAsyncProducerConsumerITCase.java b/flink-runtime/src/test/java/org/apache/flink/runtime/taskmanager/TaskCancelAsyncProducerConsumerITCase.java
index 4b73b09..62ae19b 100644
--- a/flink-runtime/src/test/java/org/apache/flink/runtime/taskmanager/TaskCancelAsyncProducerConsumerITCase.java
+++ b/flink-runtime/src/test/java/org/apache/flink/runtime/taskmanager/TaskCancelAsyncProducerConsumerITCase.java
@@ -37,7 +37,7 @@ import org.apache.flink.runtime.jobmanager.scheduler.SlotSharingGroup;
import org.apache.flink.runtime.minicluster.MiniCluster;
import org.apache.flink.runtime.minicluster.MiniClusterConfiguration;
import org.apache.flink.runtime.testingUtils.TestingUtils;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.types.LongValue;
import org.apache.flink.util.TestLogger;
@@ -53,7 +53,7 @@ import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertTrue;
-@Category(Flip6.class)
+@Category(New.class)
public class TaskCancelAsyncProducerConsumerITCase extends TestLogger {
// The Exceptions thrown by the producer/consumer Threads
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-scala-shell/src/test/scala/org/apache/flink/api/scala/ScalaShellITCase.scala
----------------------------------------------------------------------
diff --git a/flink-scala-shell/src/test/scala/org/apache/flink/api/scala/ScalaShellITCase.scala b/flink-scala-shell/src/test/scala/org/apache/flink/api/scala/ScalaShellITCase.scala
index 95a2999..cb6231a 100644
--- a/flink-scala-shell/src/test/scala/org/apache/flink/api/scala/ScalaShellITCase.scala
+++ b/flink-scala-shell/src/test/scala/org/apache/flink/api/scala/ScalaShellITCase.scala
@@ -322,7 +322,7 @@ object ScalaShellITCase {
@BeforeClass
def beforeAll(): Unit = {
configuration.setInteger(ConfigConstants.TASK_MANAGER_NUM_TASK_SLOTS, parallelism)
- configuration.setString(CoreOptions.MODE, CoreOptions.OLD_MODE)
+ configuration.setString(CoreOptions.MODE, CoreOptions.LEGACY_MODE)
cluster = Option(new StandaloneMiniCluster(configuration))
}
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-scala-shell/src/test/scala/org/apache/flink/api/scala/ScalaShellLocalStartupITCase.scala
----------------------------------------------------------------------
diff --git a/flink-scala-shell/src/test/scala/org/apache/flink/api/scala/ScalaShellLocalStartupITCase.scala b/flink-scala-shell/src/test/scala/org/apache/flink/api/scala/ScalaShellLocalStartupITCase.scala
index 34832e7..9365948 100644
--- a/flink-scala-shell/src/test/scala/org/apache/flink/api/scala/ScalaShellLocalStartupITCase.scala
+++ b/flink-scala-shell/src/test/scala/org/apache/flink/api/scala/ScalaShellLocalStartupITCase.scala
@@ -85,7 +85,7 @@ class ScalaShellLocalStartupITCase extends TestLogger {
System.setOut(new PrintStream(baos))
val configuration = new Configuration()
- configuration.setString(CoreOptions.MODE, CoreOptions.OLD_MODE)
+ configuration.setString(CoreOptions.MODE, CoreOptions.LEGACY_MODE)
val dir = temporaryFolder.newFolder()
BootstrapTools.writeConfiguration(configuration, new File(dir, "flink-conf.yaml"))
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/Flip6LocalStreamEnvironment.java
----------------------------------------------------------------------
diff --git a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/Flip6LocalStreamEnvironment.java b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/Flip6LocalStreamEnvironment.java
deleted file mode 100644
index 4cc23fc..0000000
--- a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/Flip6LocalStreamEnvironment.java
+++ /dev/null
@@ -1,111 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one or more
- * contributor license agreements. See the NOTICE file distributed with
- * this work for additional information regarding copyright ownership.
- * The ASF licenses this file to You under the Apache License, Version 2.0
- * (the "License"); you may not use this file except in compliance with
- * the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.apache.flink.streaming.api.environment;
-
-import org.apache.flink.annotation.Internal;
-import org.apache.flink.api.common.JobExecutionResult;
-import org.apache.flink.configuration.Configuration;
-import org.apache.flink.configuration.RestOptions;
-import org.apache.flink.configuration.TaskManagerOptions;
-import org.apache.flink.runtime.jobgraph.JobGraph;
-import org.apache.flink.runtime.minicluster.MiniCluster;
-import org.apache.flink.runtime.minicluster.MiniClusterConfiguration;
-import org.apache.flink.streaming.api.graph.StreamGraph;
-
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-/**
- * The Flip6LocalStreamEnvironment is a StreamExecutionEnvironment that runs the program locally,
- * multi-threaded, in the JVM where the environment is instantiated. It spawns an embedded
- * Flink cluster in the background and executes the program on that cluster.
- *
- * <p>When this environment is instantiated, it uses a default parallelism of {@code 1}. The default
- * parallelism can be set via {@link #setParallelism(int)}.
- */
-@Internal
-public class Flip6LocalStreamEnvironment extends LocalStreamEnvironment {
-
- private static final Logger LOG = LoggerFactory.getLogger(Flip6LocalStreamEnvironment.class);
-
- /**
- * Creates a new mini cluster stream environment that uses the default configuration.
- */
- public Flip6LocalStreamEnvironment() {
- this(null);
- }
-
- /**
- * Creates a new mini cluster stream environment that configures its local executor with the given configuration.
- *
- * @param config The configuration used to configure the local executor.
- */
- public Flip6LocalStreamEnvironment(Configuration config) {
- super(config);
- setParallelism(1);
- }
-
- /**
- * Executes the JobGraph of the on a mini cluster of CLusterUtil with a user
- * specified name.
- *
- * @param jobName
- * name of the job
- * @return The result of the job execution, containing elapsed time and accumulators.
- */
- @Override
- public JobExecutionResult execute(String jobName) throws Exception {
- // transform the streaming program into a JobGraph
- StreamGraph streamGraph = getStreamGraph();
- streamGraph.setJobName(jobName);
-
- JobGraph jobGraph = streamGraph.getJobGraph();
- jobGraph.setAllowQueuedScheduling(true);
-
- Configuration configuration = new Configuration();
- configuration.addAll(jobGraph.getJobConfiguration());
- configuration.setLong(TaskManagerOptions.MANAGED_MEMORY_SIZE, -1L);
-
- // add (and override) the settings with what the user defined
- configuration.addAll(this.conf);
-
- configuration.setInteger(RestOptions.REST_PORT, 0);
-
- MiniClusterConfiguration cfg = new MiniClusterConfiguration.Builder()
- .setConfiguration(configuration)
- .setNumSlotsPerTaskManager(jobGraph.getMaximumParallelism())
- .build();
-
- if (LOG.isInfoEnabled()) {
- LOG.info("Running job on local embedded Flink mini cluster");
- }
-
- MiniCluster miniCluster = new MiniCluster(cfg);
-
- try {
- miniCluster.start();
- configuration.setInteger(RestOptions.REST_PORT, miniCluster.getRestAddress().getPort());
-
- return miniCluster.executeJobBlocking(jobGraph);
- }
- finally {
- transformations.clear();
- miniCluster.close();
- }
- }
-}
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/LegacyLocalStreamEnvironment.java
----------------------------------------------------------------------
diff --git a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/LegacyLocalStreamEnvironment.java b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/LegacyLocalStreamEnvironment.java
new file mode 100644
index 0000000..8341ec4
--- /dev/null
+++ b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/LegacyLocalStreamEnvironment.java
@@ -0,0 +1,101 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.flink.streaming.api.environment;
+
+import org.apache.flink.annotation.Public;
+import org.apache.flink.api.common.JobExecutionResult;
+import org.apache.flink.configuration.ConfigConstants;
+import org.apache.flink.configuration.Configuration;
+import org.apache.flink.configuration.TaskManagerOptions;
+import org.apache.flink.runtime.jobgraph.JobGraph;
+import org.apache.flink.runtime.minicluster.LocalFlinkMiniCluster;
+import org.apache.flink.streaming.api.graph.StreamGraph;
+
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * The LocalStreamEnvironment is a StreamExecutionEnvironment that runs the program locally,
+ * multi-threaded, in the JVM where the environment is instantiated. It spawns an embedded
+ * Flink cluster in the background and executes the program on that cluster.
+ *
+ * <p>When this environment is instantiated, it uses a default parallelism of {@code 1}. The default
+ * parallelism can be set via {@link #setParallelism(int)}.
+ *
+ * <p>Local environments can also be instantiated through {@link StreamExecutionEnvironment#createLocalEnvironment()}
+ * and {@link StreamExecutionEnvironment#createLocalEnvironment(int)}. The former version will pick a
+ * default parallelism equal to the number of hardware contexts in the local machine.
+ */
+@Public
+public class LegacyLocalStreamEnvironment extends LocalStreamEnvironment {
+
+ private static final Logger LOG = LoggerFactory.getLogger(LocalStreamEnvironment.class);
+
+ public LegacyLocalStreamEnvironment() {
+ this(new Configuration());
+ }
+
+ /**
+ * Creates a new local stream environment that configures its local executor with the given configuration.
+ *
+ * @param config The configuration used to configure the local executor.
+ */
+ public LegacyLocalStreamEnvironment(Configuration config) {
+ super(config);
+ }
+
+ /**
+ * Executes the JobGraph of the on a mini cluster of CLusterUtil with a user
+ * specified name.
+ *
+ * @param jobName
+ * name of the job
+ * @return The result of the job execution, containing elapsed time and accumulators.
+ */
+ @Override
+ public JobExecutionResult execute(String jobName) throws Exception {
+ // transform the streaming program into a JobGraph
+ StreamGraph streamGraph = getStreamGraph();
+ streamGraph.setJobName(jobName);
+
+ JobGraph jobGraph = streamGraph.getJobGraph();
+
+ Configuration configuration = new Configuration();
+ configuration.addAll(jobGraph.getJobConfiguration());
+
+ configuration.setLong(TaskManagerOptions.MANAGED_MEMORY_SIZE, -1L);
+ configuration.setInteger(ConfigConstants.TASK_MANAGER_NUM_TASK_SLOTS, jobGraph.getMaximumParallelism());
+
+ // add (and override) the settings with what the user defined
+ configuration.addAll(getConfiguration());
+
+ if (LOG.isInfoEnabled()) {
+ LOG.info("Running job on local embedded Flink mini cluster");
+ }
+
+ LocalFlinkMiniCluster exec = new LocalFlinkMiniCluster(configuration, true);
+ try {
+ exec.start();
+ return exec.submitJobAndWait(jobGraph, getConfig().isSysoutLoggingEnabled());
+ }
+ finally {
+ transformations.clear();
+ exec.stop();
+ }
+ }
+}
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/LocalStreamEnvironment.java
----------------------------------------------------------------------
diff --git a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/LocalStreamEnvironment.java b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/LocalStreamEnvironment.java
index a53e2a3..935c78e 100644
--- a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/LocalStreamEnvironment.java
+++ b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/LocalStreamEnvironment.java
@@ -21,16 +21,19 @@ import org.apache.flink.annotation.Public;
import org.apache.flink.api.common.InvalidProgramException;
import org.apache.flink.api.common.JobExecutionResult;
import org.apache.flink.api.java.ExecutionEnvironment;
-import org.apache.flink.configuration.ConfigConstants;
import org.apache.flink.configuration.Configuration;
+import org.apache.flink.configuration.RestOptions;
import org.apache.flink.configuration.TaskManagerOptions;
import org.apache.flink.runtime.jobgraph.JobGraph;
-import org.apache.flink.runtime.minicluster.LocalFlinkMiniCluster;
+import org.apache.flink.runtime.minicluster.MiniCluster;
+import org.apache.flink.runtime.minicluster.MiniClusterConfiguration;
import org.apache.flink.streaming.api.graph.StreamGraph;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import javax.annotation.Nonnull;
+
/**
* The LocalStreamEnvironment is a StreamExecutionEnvironment that runs the program locally,
* multi-threaded, in the JVM where the environment is instantiated. It spawns an embedded
@@ -38,39 +41,38 @@ import org.slf4j.LoggerFactory;
*
* <p>When this environment is instantiated, it uses a default parallelism of {@code 1}. The default
* parallelism can be set via {@link #setParallelism(int)}.
- *
- * <p>Local environments can also be instantiated through {@link StreamExecutionEnvironment#createLocalEnvironment()}
- * and {@link StreamExecutionEnvironment#createLocalEnvironment(int)}. The former version will pick a
- * default parallelism equal to the number of hardware contexts in the local machine.
*/
@Public
public class LocalStreamEnvironment extends StreamExecutionEnvironment {
private static final Logger LOG = LoggerFactory.getLogger(LocalStreamEnvironment.class);
- /** The configuration to use for the local cluster. */
- protected final Configuration conf;
+ private final Configuration configuration;
/**
- * Creates a new local stream environment that uses the default configuration.
+ * Creates a new mini cluster stream environment that uses the default configuration.
*/
public LocalStreamEnvironment() {
- this(null);
+ this(new Configuration());
}
/**
- * Creates a new local stream environment that configures its local executor with the given configuration.
+ * Creates a new mini cluster stream environment that configures its local executor with the given configuration.
*
- * @param config The configuration used to configure the local executor.
+ * @param configuration The configuration used to configure the local executor.
*/
- public LocalStreamEnvironment(Configuration config) {
+ public LocalStreamEnvironment(@Nonnull Configuration configuration) {
if (!ExecutionEnvironment.areExplicitEnvironmentsAllowed()) {
throw new InvalidProgramException(
- "The LocalStreamEnvironment cannot be used when submitting a program through a client, " +
- "or running in a TestEnvironment context.");
+ "The LocalStreamEnvironment cannot be used when submitting a program through a client, " +
+ "or running in a TestEnvironment context.");
}
+ this.configuration = configuration;
+ setParallelism(1);
+ }
- this.conf = config == null ? new Configuration() : config;
+ protected Configuration getConfiguration() {
+ return configuration;
}
/**
@@ -88,28 +90,37 @@ public class LocalStreamEnvironment extends StreamExecutionEnvironment {
streamGraph.setJobName(jobName);
JobGraph jobGraph = streamGraph.getJobGraph();
+ jobGraph.setAllowQueuedScheduling(true);
Configuration configuration = new Configuration();
configuration.addAll(jobGraph.getJobConfiguration());
-
configuration.setLong(TaskManagerOptions.MANAGED_MEMORY_SIZE, -1L);
- configuration.setInteger(ConfigConstants.TASK_MANAGER_NUM_TASK_SLOTS, jobGraph.getMaximumParallelism());
// add (and override) the settings with what the user defined
- configuration.addAll(this.conf);
+ configuration.addAll(this.configuration);
+
+ configuration.setInteger(RestOptions.REST_PORT, 0);
+
+ MiniClusterConfiguration cfg = new MiniClusterConfiguration.Builder()
+ .setConfiguration(configuration)
+ .setNumSlotsPerTaskManager(jobGraph.getMaximumParallelism())
+ .build();
if (LOG.isInfoEnabled()) {
LOG.info("Running job on local embedded Flink mini cluster");
}
- LocalFlinkMiniCluster exec = new LocalFlinkMiniCluster(configuration, true);
+ MiniCluster miniCluster = new MiniCluster(cfg);
+
try {
- exec.start();
- return exec.submitJobAndWait(jobGraph, getConfig().isSysoutLoggingEnabled());
+ miniCluster.start();
+ configuration.setInteger(RestOptions.REST_PORT, miniCluster.getRestAddress().getPort());
+
+ return miniCluster.executeJobBlocking(jobGraph);
}
finally {
transformations.clear();
- exec.stop();
+ miniCluster.close();
}
}
}
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/RemoteStreamEnvironment.java
----------------------------------------------------------------------
diff --git a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/RemoteStreamEnvironment.java b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/RemoteStreamEnvironment.java
index 036cf4d..075a3cd 100644
--- a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/RemoteStreamEnvironment.java
+++ b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/RemoteStreamEnvironment.java
@@ -203,7 +203,7 @@ public class RemoteStreamEnvironment extends StreamExecutionEnvironment {
final ClusterClient<?> client;
try {
- if (CoreOptions.OLD_MODE.equals(configuration.getString(CoreOptions.MODE))) {
+ if (CoreOptions.LEGACY_MODE.equals(configuration.getString(CoreOptions.MODE))) {
client = new StandaloneClusterClient(configuration);
} else {
client = new RestClusterClient<>(configuration, "RemoteStreamEnvironment");
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/StreamExecutionEnvironment.java
----------------------------------------------------------------------
diff --git a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/StreamExecutionEnvironment.java b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/StreamExecutionEnvironment.java
index c39201c..fa81c27 100644
--- a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/StreamExecutionEnvironment.java
+++ b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/environment/StreamExecutionEnvironment.java
@@ -1652,10 +1652,10 @@ public abstract class StreamExecutionEnvironment {
public static LocalStreamEnvironment createLocalEnvironment(int parallelism, Configuration configuration) {
final LocalStreamEnvironment currentEnvironment;
- if (CoreOptions.FLIP6_MODE.equals(configuration.getString(CoreOptions.MODE))) {
- currentEnvironment = new Flip6LocalStreamEnvironment(configuration);
- } else {
+ if (CoreOptions.NEW_MODE.equals(configuration.getString(CoreOptions.MODE))) {
currentEnvironment = new LocalStreamEnvironment(configuration);
+ } else {
+ currentEnvironment = new LegacyLocalStreamEnvironment(configuration);
}
currentEnvironment.setParallelism(parallelism);
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/environment/LocalStreamEnvironmentITCase.java
----------------------------------------------------------------------
diff --git a/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/environment/LocalStreamEnvironmentITCase.java b/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/environment/LocalStreamEnvironmentITCase.java
index f302eda..84a1ce9 100644
--- a/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/environment/LocalStreamEnvironmentITCase.java
+++ b/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/environment/LocalStreamEnvironmentITCase.java
@@ -26,7 +26,7 @@ import org.junit.Test;
import static org.junit.Assert.assertEquals;
/**
- * Tests for {@link Flip6LocalStreamEnvironment}.
+ * Tests for {@link LocalStreamEnvironment}.
*/
@SuppressWarnings("serial")
public class LocalStreamEnvironmentITCase extends TestLogger {
@@ -37,7 +37,7 @@ public class LocalStreamEnvironmentITCase extends TestLogger {
*/
@Test
public void testRunIsolatedJob() throws Exception {
- Flip6LocalStreamEnvironment env = new Flip6LocalStreamEnvironment();
+ LocalStreamEnvironment env = new LocalStreamEnvironment();
assertEquals(1, env.getParallelism());
addSmallBoundedJob(env, 3);
@@ -50,7 +50,7 @@ public class LocalStreamEnvironmentITCase extends TestLogger {
*/
@Test
public void testMultipleJobsAfterAnother() throws Exception {
- Flip6LocalStreamEnvironment env = new Flip6LocalStreamEnvironment();
+ LocalStreamEnvironment env = new LocalStreamEnvironment();
addSmallBoundedJob(env, 3);
env.execute();
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-test-utils-parent/flink-test-utils-junit/src/main/java/org/apache/flink/testutils/category/Flip6.java
----------------------------------------------------------------------
diff --git a/flink-test-utils-parent/flink-test-utils-junit/src/main/java/org/apache/flink/testutils/category/Flip6.java b/flink-test-utils-parent/flink-test-utils-junit/src/main/java/org/apache/flink/testutils/category/Flip6.java
deleted file mode 100644
index f23e5b5..0000000
--- a/flink-test-utils-parent/flink-test-utils-junit/src/main/java/org/apache/flink/testutils/category/Flip6.java
+++ /dev/null
@@ -1,25 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements. See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership. The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * "License"); you may not use this file except in compliance
- * with the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.apache.flink.testutils.category;
-
-/**
- * Category marker interface for Flip-6 tests.
- */
-public interface Flip6 {
-}
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-test-utils-parent/flink-test-utils-junit/src/main/java/org/apache/flink/testutils/category/LegacyAndNew.java
----------------------------------------------------------------------
diff --git a/flink-test-utils-parent/flink-test-utils-junit/src/main/java/org/apache/flink/testutils/category/LegacyAndNew.java b/flink-test-utils-parent/flink-test-utils-junit/src/main/java/org/apache/flink/testutils/category/LegacyAndNew.java
new file mode 100644
index 0000000..180f87e
--- /dev/null
+++ b/flink-test-utils-parent/flink-test-utils-junit/src/main/java/org/apache/flink/testutils/category/LegacyAndNew.java
@@ -0,0 +1,26 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.flink.testutils.category;
+
+/**
+ * Category marker interface for tests relevant for the legacy and
+ * new architecture.
+ */
+public interface LegacyAndNew {
+}
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-test-utils-parent/flink-test-utils-junit/src/main/java/org/apache/flink/testutils/category/New.java
----------------------------------------------------------------------
diff --git a/flink-test-utils-parent/flink-test-utils-junit/src/main/java/org/apache/flink/testutils/category/New.java b/flink-test-utils-parent/flink-test-utils-junit/src/main/java/org/apache/flink/testutils/category/New.java
new file mode 100644
index 0000000..2f8fd44
--- /dev/null
+++ b/flink-test-utils-parent/flink-test-utils-junit/src/main/java/org/apache/flink/testutils/category/New.java
@@ -0,0 +1,25 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.flink.testutils.category;
+
+/**
+ * Category marker interface for tests based on the new architecture.
+ */
+public interface New {
+}
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-test-utils-parent/flink-test-utils-junit/src/main/java/org/apache/flink/testutils/category/OldAndFlip6.java
----------------------------------------------------------------------
diff --git a/flink-test-utils-parent/flink-test-utils-junit/src/main/java/org/apache/flink/testutils/category/OldAndFlip6.java b/flink-test-utils-parent/flink-test-utils-junit/src/main/java/org/apache/flink/testutils/category/OldAndFlip6.java
deleted file mode 100644
index fd24534..0000000
--- a/flink-test-utils-parent/flink-test-utils-junit/src/main/java/org/apache/flink/testutils/category/OldAndFlip6.java
+++ /dev/null
@@ -1,25 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements. See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership. The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * "License"); you may not use this file except in compliance
- * with the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing, software
- * distributed under the License is distributed on an "AS IS" BASIS,
- * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
- * See the License for the specific language governing permissions and
- * limitations under the License.
- */
-
-package org.apache.flink.testutils.category;
-
-/**
- * Category marker interface for old and Flip-6 tests.
- */
-public interface OldAndFlip6 {
-}
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-test-utils-parent/flink-test-utils/src/main/java/org/apache/flink/test/util/AbstractTestBase.java
----------------------------------------------------------------------
diff --git a/flink-test-utils-parent/flink-test-utils/src/main/java/org/apache/flink/test/util/AbstractTestBase.java b/flink-test-utils-parent/flink-test-utils/src/main/java/org/apache/flink/test/util/AbstractTestBase.java
index d73f624..c53ed8b 100644
--- a/flink-test-utils-parent/flink-test-utils/src/main/java/org/apache/flink/test/util/AbstractTestBase.java
+++ b/flink-test-utils-parent/flink-test-utils/src/main/java/org/apache/flink/test/util/AbstractTestBase.java
@@ -19,7 +19,7 @@
package org.apache.flink.test.util;
import org.apache.flink.configuration.Configuration;
-import org.apache.flink.testutils.category.OldAndFlip6;
+import org.apache.flink.testutils.category.LegacyAndNew;
import org.apache.flink.util.FileUtils;
import org.junit.ClassRule;
@@ -56,7 +56,7 @@ import java.io.IOException;
*
* </pre>
*/
-@Category(OldAndFlip6.class)
+@Category(LegacyAndNew.class)
public abstract class AbstractTestBase extends TestBaseUtils {
private static final int DEFAULT_PARALLELISM = 4;
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-test-utils-parent/flink-test-utils/src/main/java/org/apache/flink/test/util/MiniClusterResource.java
----------------------------------------------------------------------
diff --git a/flink-test-utils-parent/flink-test-utils/src/main/java/org/apache/flink/test/util/MiniClusterResource.java b/flink-test-utils-parent/flink-test-utils/src/main/java/org/apache/flink/test/util/MiniClusterResource.java
index 545be86..9b0ac77 100644
--- a/flink-test-utils-parent/flink-test-utils/src/main/java/org/apache/flink/test/util/MiniClusterResource.java
+++ b/flink-test-utils-parent/flink-test-utils/src/main/java/org/apache/flink/test/util/MiniClusterResource.java
@@ -55,7 +55,7 @@ public class MiniClusterResource extends ExternalResource {
private static final String CODEBASE_KEY = "codebase";
- private static final String FLIP6_CODEBASE = "flip6";
+ private static final String NEW_CODEBASE = "new";
private final MiniClusterResourceConfiguration miniClusterResourceConfiguration;
@@ -80,7 +80,7 @@ public class MiniClusterResource extends ExternalResource {
final boolean enableClusterClient) {
this(
miniClusterResourceConfiguration,
- Objects.equals(FLIP6_CODEBASE, System.getProperty(CODEBASE_KEY)) ? MiniClusterType.FLIP6 : MiniClusterType.OLD,
+ Objects.equals(NEW_CODEBASE, System.getProperty(CODEBASE_KEY)) ? MiniClusterType.NEW : MiniClusterType.LEGACY,
enableClusterClient);
}
@@ -104,7 +104,7 @@ public class MiniClusterResource extends ExternalResource {
public ClusterClient<?> getClusterClient() {
if (!enableClusterClient) {
// this check is technically only necessary for legacy clusters
- // we still fail here for flip6 to keep the behaviors in sync
+ // we still fail here to keep the behaviors in sync
throw new IllegalStateException("To use the client you must enable it with the constructor.");
}
@@ -164,18 +164,18 @@ public class MiniClusterResource extends ExternalResource {
private void startJobExecutorService(MiniClusterType miniClusterType) throws Exception {
switch (miniClusterType) {
- case OLD:
- startOldMiniCluster();
+ case LEGACY:
+ startLegacyMiniCluster();
break;
- case FLIP6:
- startFlip6MiniCluster();
+ case NEW:
+ startMiniCluster();
break;
default:
throw new FlinkRuntimeException("Unknown MiniClusterType " + miniClusterType + '.');
}
}
- private void startOldMiniCluster() throws Exception {
+ private void startLegacyMiniCluster() throws Exception {
final Configuration configuration = new Configuration(miniClusterResourceConfiguration.getConfiguration());
configuration.setInteger(ConfigConstants.LOCAL_NUMBER_TASK_MANAGER, miniClusterResourceConfiguration.getNumberTaskManagers());
configuration.setInteger(ConfigConstants.TASK_MANAGER_NUM_TASK_SLOTS, miniClusterResourceConfiguration.getNumberSlotsPerTaskManager());
@@ -190,7 +190,7 @@ public class MiniClusterResource extends ExternalResource {
}
}
- private void startFlip6MiniCluster() throws Exception {
+ private void startMiniCluster() throws Exception {
final Configuration configuration = miniClusterResourceConfiguration.getConfiguration();
// we need to set this since a lot of test expect this because TestBaseUtils.startCluster()
@@ -284,7 +284,7 @@ public class MiniClusterResource extends ExternalResource {
* Type of the mini cluster to start.
*/
public enum MiniClusterType {
- OLD,
- FLIP6
+ LEGACY,
+ NEW
}
}
http://git-wip-us.apache.org/repos/asf/flink/blob/af5279e9/flink-tests/src/test/java/org/apache/flink/runtime/jobmaster/JobMasterTriggerSavepointIT.java
----------------------------------------------------------------------
diff --git a/flink-tests/src/test/java/org/apache/flink/runtime/jobmaster/JobMasterTriggerSavepointIT.java b/flink-tests/src/test/java/org/apache/flink/runtime/jobmaster/JobMasterTriggerSavepointIT.java
index f9edfa6..db3313c 100644
--- a/flink-tests/src/test/java/org/apache/flink/runtime/jobmaster/JobMasterTriggerSavepointIT.java
+++ b/flink-tests/src/test/java/org/apache/flink/runtime/jobmaster/JobMasterTriggerSavepointIT.java
@@ -36,7 +36,7 @@ import org.apache.flink.runtime.jobgraph.tasks.AbstractInvokable;
import org.apache.flink.runtime.jobgraph.tasks.CheckpointCoordinatorConfiguration;
import org.apache.flink.runtime.jobgraph.tasks.JobCheckpointingSettings;
import org.apache.flink.test.util.AbstractTestBase;
-import org.apache.flink.testutils.category.Flip6;
+import org.apache.flink.testutils.category.New;
import org.apache.flink.util.ExceptionUtils;
import org.junit.Assume;
@@ -67,7 +67,7 @@ import static org.hamcrest.Matchers.isOneOf;
*
* @see org.apache.flink.runtime.jobmaster.JobMaster
*/
-@Category(Flip6.class)
+@Category(New.class)
public class JobMasterTriggerSavepointIT extends AbstractTestBase {
private static CountDownLatch invokeLatch;