You are viewing a plain text version of this content. The canonical link for it is here.
Posted to commits@ignite.apache.org by il...@apache.org on 2019/12/31 09:42:37 UTC

[ignite] branch master updated: IGNITE-12356 Migrate Flink module to ignite-extensions - Fixes #7222.

This is an automated email from the ASF dual-hosted git repository.

ilyak pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ignite.git


The following commit(s) were added to refs/heads/master by this push:
     new f550b95  IGNITE-12356 Migrate Flink module to ignite-extensions - Fixes #7222.
f550b95 is described below

commit f550b95547f672409663d7b0052939e62026dbbf
Author: samaitra <sa...@gmail.com>
AuthorDate: Tue Dec 31 12:38:01 2019 +0300

    IGNITE-12356 Migrate Flink module to ignite-extensions - Fixes #7222.
    
    Signed-off-by: Ilya Kasnacheev <il...@gmail.com>
---
 modules/flink/README.txt                           |  33 ---
 modules/flink/licenses/apache-2.0.txt              | 202 -------------------
 modules/flink/pom.xml                              | 194 ------------------
 .../org/apache/ignite/sink/flink/IgniteSink.java   | 197 ------------------
 .../org/apache/ignite/sink/flink/package-info.java |  22 --
 .../apache/ignite/source/flink/IgniteSource.java   | 223 ---------------------
 .../ignite/source/flink/TaskRemoteFilter.java      |  60 ------
 .../apache/ignite/source/flink/package-info.java   |  21 --
 .../ignite/sink/flink/FlinkIgniteSinkSelfTest.java |  84 --------
 .../sink/flink/FlinkIgniteSinkSelfTestSuite.java   |  29 ---
 .../source/flink/FlinkIgniteSourceSelfTest.java    | 154 --------------
 .../flink/FlinkIgniteSourceSelfTestSuite.java      |  30 ---
 .../flink/src/test/resources/example-ignite.xml    |  73 -------
 parent/pom.xml                                     |   8 -
 pom.xml                                            |   1 -
 15 files changed, 1331 deletions(-)

diff --git a/modules/flink/README.txt b/modules/flink/README.txt
deleted file mode 100644
index a198b1b..0000000
--- a/modules/flink/README.txt
+++ /dev/null
@@ -1,33 +0,0 @@
-Apache Ignite Flink Sink Module
------------------------------------
-
-Apache Ignite Flink Sink module is a streaming connector to inject Flink data into Ignite cache.
-
-Starting data transfer to Ignite can be done with the following steps.
-
-1. Import Ignite Flink Sink Module in Maven Project
-
-If you are using Maven to manage dependencies of your project, you can add Flink module
-dependency like this (replace '${ignite.version}' with actual Ignite version you are
-interested in):
-
-<project xmlns="http://maven.apache.org/POM/4.0.0"
-    xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
-    xsi:schemaLocation="http://maven.apache.org/POM/4.0.0
-                        http://maven.apache.org/xsd/maven-4.0.0.xsd">
-    ...
-    <dependencies>
-        ...
-        <dependency>
-            <groupId>org.apache.ignite</groupId>
-            <artifactId>ignite-flink</artifactId>
-            <version>${ignite.version}</version>
-        </dependency>
-        ...
-    </dependencies>
-    ...
-</project>
-
-2. Create an Ignite configuration file (see example-ignite.xml) and make sure it is accessible from the sink.
-
-3. Make sure your data input to the sink is specified. For example `input.addSink(igniteSinkObject)`
diff --git a/modules/flink/licenses/apache-2.0.txt b/modules/flink/licenses/apache-2.0.txt
deleted file mode 100644
index d645695..0000000
--- a/modules/flink/licenses/apache-2.0.txt
+++ /dev/null
@@ -1,202 +0,0 @@
-
-                                 Apache License
-                           Version 2.0, January 2004
-                        http://www.apache.org/licenses/
-
-   TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION
-
-   1. Definitions.
-
-      "License" shall mean the terms and conditions for use, reproduction,
-      and distribution as defined by Sections 1 through 9 of this document.
-
-      "Licensor" shall mean the copyright owner or entity authorized by
-      the copyright owner that is granting the License.
-
-      "Legal Entity" shall mean the union of the acting entity and all
-      other entities that control, are controlled by, or are under common
-      control with that entity. For the purposes of this definition,
-      "control" means (i) the power, direct or indirect, to cause the
-      direction or management of such entity, whether by contract or
-      otherwise, or (ii) ownership of fifty percent (50%) or more of the
-      outstanding shares, or (iii) beneficial ownership of such entity.
-
-      "You" (or "Your") shall mean an individual or Legal Entity
-      exercising permissions granted by this License.
-
-      "Source" form shall mean the preferred form for making modifications,
-      including but not limited to software source code, documentation
-      source, and configuration files.
-
-      "Object" form shall mean any form resulting from mechanical
-      transformation or translation of a Source form, including but
-      not limited to compiled object code, generated documentation,
-      and conversions to other media types.
-
-      "Work" shall mean the work of authorship, whether in Source or
-      Object form, made available under the License, as indicated by a
-      copyright notice that is included in or attached to the work
-      (an example is provided in the Appendix below).
-
-      "Derivative Works" shall mean any work, whether in Source or Object
-      form, that is based on (or derived from) the Work and for which the
-      editorial revisions, annotations, elaborations, or other modifications
-      represent, as a whole, an original work of authorship. For the purposes
-      of this License, Derivative Works shall not include works that remain
-      separable from, or merely link (or bind by name) to the interfaces of,
-      the Work and Derivative Works thereof.
-
-      "Contribution" shall mean any work of authorship, including
-      the original version of the Work and any modifications or additions
-      to that Work or Derivative Works thereof, that is intentionally
-      submitted to Licensor for inclusion in the Work by the copyright owner
-      or by an individual or Legal Entity authorized to submit on behalf of
-      the copyright owner. For the purposes of this definition, "submitted"
-      means any form of electronic, verbal, or written communication sent
-      to the Licensor or its representatives, including but not limited to
-      communication on electronic mailing lists, source code control systems,
-      and issue tracking systems that are managed by, or on behalf of, the
-      Licensor for the purpose of discussing and improving the Work, but
-      excluding communication that is conspicuously marked or otherwise
-      designated in writing by the copyright owner as "Not a Contribution."
-
-      "Contributor" shall mean Licensor and any individual or Legal Entity
-      on behalf of whom a Contribution has been received by Licensor and
-      subsequently incorporated within the Work.
-
-   2. Grant of Copyright License. Subject to the terms and conditions of
-      this License, each Contributor hereby grants to You a perpetual,
-      worldwide, non-exclusive, no-charge, royalty-free, irrevocable
-      copyright license to reproduce, prepare Derivative Works of,
-      publicly display, publicly perform, sublicense, and distribute the
-      Work and such Derivative Works in Source or Object form.
-
-   3. Grant of Patent License. Subject to the terms and conditions of
-      this License, each Contributor hereby grants to You a perpetual,
-      worldwide, non-exclusive, no-charge, royalty-free, irrevocable
-      (except as stated in this section) patent license to make, have made,
-      use, offer to sell, sell, import, and otherwise transfer the Work,
-      where such license applies only to those patent claims licensable
-      by such Contributor that are necessarily infringed by their
-      Contribution(s) alone or by combination of their Contribution(s)
-      with the Work to which such Contribution(s) was submitted. If You
-      institute patent litigation against any entity (including a
-      cross-claim or counterclaim in a lawsuit) alleging that the Work
-      or a Contribution incorporated within the Work constitutes direct
-      or contributory patent infringement, then any patent licenses
-      granted to You under this License for that Work shall terminate
-      as of the date such litigation is filed.
-
-   4. Redistribution. You may reproduce and distribute copies of the
-      Work or Derivative Works thereof in any medium, with or without
-      modifications, and in Source or Object form, provided that You
-      meet the following conditions:
-
-      (a) You must give any other recipients of the Work or
-          Derivative Works a copy of this License; and
-
-      (b) You must cause any modified files to carry prominent notices
-          stating that You changed the files; and
-
-      (c) You must retain, in the Source form of any Derivative Works
-          that You distribute, all copyright, patent, trademark, and
-          attribution notices from the Source form of the Work,
-          excluding those notices that do not pertain to any part of
-          the Derivative Works; and
-
-      (d) If the Work includes a "NOTICE" text file as part of its
-          distribution, then any Derivative Works that You distribute must
-          include a readable copy of the attribution notices contained
-          within such NOTICE file, excluding those notices that do not
-          pertain to any part of the Derivative Works, in at least one
-          of the following places: within a NOTICE text file distributed
-          as part of the Derivative Works; within the Source form or
-          documentation, if provided along with the Derivative Works; or,
-          within a display generated by the Derivative Works, if and
-          wherever such third-party notices normally appear. The contents
-          of the NOTICE file are for informational purposes only and
-          do not modify the License. You may add Your own attribution
-          notices within Derivative Works that You distribute, alongside
-          or as an addendum to the NOTICE text from the Work, provided
-          that such additional attribution notices cannot be construed
-          as modifying the License.
-
-      You may add Your own copyright statement to Your modifications and
-      may provide additional or different license terms and conditions
-      for use, reproduction, or distribution of Your modifications, or
-      for any such Derivative Works as a whole, provided Your use,
-      reproduction, and distribution of the Work otherwise complies with
-      the conditions stated in this License.
-
-   5. Submission of Contributions. Unless You explicitly state otherwise,
-      any Contribution intentionally submitted for inclusion in the Work
-      by You to the Licensor shall be under the terms and conditions of
-      this License, without any additional terms or conditions.
-      Notwithstanding the above, nothing herein shall supersede or modify
-      the terms of any separate license agreement you may have executed
-      with Licensor regarding such Contributions.
-
-   6. Trademarks. This License does not grant permission to use the trade
-      names, trademarks, service marks, or product names of the Licensor,
-      except as required for reasonable and customary use in describing the
-      origin of the Work and reproducing the content of the NOTICE file.
-
-   7. Disclaimer of Warranty. Unless required by applicable law or
-      agreed to in writing, Licensor provides the Work (and each
-      Contributor provides its Contributions) on an "AS IS" BASIS,
-      WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or
-      implied, including, without limitation, any warranties or conditions
-      of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A
-      PARTICULAR PURPOSE. You are solely responsible for determining the
-      appropriateness of using or redistributing the Work and assume any
-      risks associated with Your exercise of permissions under this License.
-
-   8. Limitation of Liability. In no event and under no legal theory,
-      whether in tort (including negligence), contract, or otherwise,
-      unless required by applicable law (such as deliberate and grossly
-      negligent acts) or agreed to in writing, shall any Contributor be
-      liable to You for damages, including any direct, indirect, special,
-      incidental, or consequential damages of any character arising as a
-      result of this License or out of the use or inability to use the
-      Work (including but not limited to damages for loss of goodwill,
-      work stoppage, computer failure or malfunction, or any and all
-      other commercial damages or losses), even if such Contributor
-      has been advised of the possibility of such damages.
-
-   9. Accepting Warranty or Additional Liability. While redistributing
-      the Work or Derivative Works thereof, You may choose to offer,
-      and charge a fee for, acceptance of support, warranty, indemnity,
-      or other liability obligations and/or rights consistent with this
-      License. However, in accepting such obligations, You may act only
-      on Your own behalf and on Your sole responsibility, not on behalf
-      of any other Contributor, and only if You agree to indemnify,
-      defend, and hold each Contributor harmless for any liability
-      incurred by, or claims asserted against, such Contributor by reason
-      of your accepting any such warranty or additional liability.
-
-   END OF TERMS AND CONDITIONS
-
-   APPENDIX: How to apply the Apache License to your work.
-
-      To apply the Apache License to your work, attach the following
-      boilerplate notice, with the fields enclosed by brackets "[]"
-      replaced with your own identifying information. (Don't include
-      the brackets!)  The text should be enclosed in the appropriate
-      comment syntax for the file format. We also recommend that a
-      file or class name and description of purpose be included on the
-      same "printed page" as the copyright notice for easier
-      identification within third-party archives.
-
-   Copyright [yyyy] [name of copyright owner]
-
-   Licensed 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.
diff --git a/modules/flink/pom.xml b/modules/flink/pom.xml
deleted file mode 100644
index 5a65a92..0000000
--- a/modules/flink/pom.xml
+++ /dev/null
@@ -1,194 +0,0 @@
-<?xml version="1.0" encoding="UTF-8"?>
-
-<!--
-  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.
--->
-
-<!--
-    POM file.
--->
-<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
-    <modelVersion>4.0.0</modelVersion>
-
-    <parent>
-        <groupId>org.apache.ignite</groupId>
-        <artifactId>ignite-parent</artifactId>
-        <version>1</version>
-        <relativePath>../../parent</relativePath>
-    </parent>
-
-    <artifactId>ignite-flink</artifactId>
-    <version>2.9.0-SNAPSHOT</version>
-    <url>http://ignite.apache.org</url>
-
-    <properties>
-        <flink.version>1.5.0</flink.version>
-        <kryo-serializers.version>0.42</kryo-serializers.version>
-    </properties>
-
-    <dependencies>
-        <dependency>
-            <groupId>org.apache.ignite</groupId>
-            <artifactId>ignite-core</artifactId>
-            <version>${project.version}</version>
-        </dependency>
-
-        <dependency>
-            <groupId>org.apache.flink</groupId>
-            <artifactId>flink-java</artifactId>
-            <version>${flink.version}</version>
-            <exclusions>
-                <exclusion>
-                    <groupId>log4j</groupId>
-                    <artifactId>log4j</artifactId>
-                </exclusion>
-                <exclusion>
-                    <groupId>org.slf4j</groupId>
-                    <artifactId>slf4j-log4j12</artifactId>
-                </exclusion>
-                <exclusion>
-                    <groupId>commons-logging</groupId>
-                    <artifactId>commons-logging</artifactId>
-                </exclusion>
-                <exclusion>
-                    <groupId>org.slf4j</groupId>
-                    <artifactId>slf4j-simple</artifactId>
-                </exclusion>
-                <exclusion>
-                    <groupId>org.slf4j</groupId>
-                    <artifactId>log4j-over-slf4j</artifactId>
-                </exclusion>
-                <exclusion>
-                    <groupId>org.apache.zookeeper</groupId>
-                    <artifactId>zookeeper</artifactId>
-                </exclusion>
-                <exclusion>
-                    <groupId>commons-beanutils</groupId>
-                    <artifactId>commons-beanutils</artifactId>
-                </exclusion>
-                <exclusion>
-                    <groupId>commons-beanutils</groupId>
-                    <artifactId>commons-beanutils-bean-collections</artifactId>
-                </exclusion>
-                <exclusion>
-                    <groupId>commons-codec</groupId>
-                    <artifactId>commons-codec</artifactId>
-                </exclusion>
-            </exclusions>
-        </dependency>
-
-        <dependency>
-            <groupId>org.apache.flink</groupId>
-            <artifactId>flink-streaming-java_2.11</artifactId>
-            <version>${flink.version}</version>
-            <exclusions>
-                <exclusion>
-                    <groupId>log4j</groupId>
-                    <artifactId>log4j</artifactId>
-                </exclusion>
-                <exclusion>
-                    <groupId>org.slf4j</groupId>
-                    <artifactId>slf4j-log4j12</artifactId>
-                </exclusion>
-                <exclusion>
-                    <groupId>commons-logging</groupId>
-                    <artifactId>commons-logging</artifactId>
-                </exclusion>
-                <exclusion>
-                    <groupId>org.slf4j</groupId>
-                    <artifactId>slf4j-simple</artifactId>
-                </exclusion>
-                <exclusion>
-                    <groupId>org.slf4j</groupId>
-                    <artifactId>log4j-over-slf4j</artifactId>
-                </exclusion>
-                <exclusion>
-                    <groupId>org.apache.zookeeper</groupId>
-                    <artifactId>zookeeper</artifactId>
-                </exclusion>
-            </exclusions>
-        </dependency>
-
-        <dependency>
-            <groupId>org.apache.flink</groupId>
-            <artifactId>flink-clients_2.11</artifactId>
-            <version>${flink.version}</version>
-            <exclusions>
-                <exclusion>
-                    <groupId>log4j</groupId>
-                    <artifactId>log4j</artifactId>
-                </exclusion>
-                <exclusion>
-                    <groupId>org.slf4j</groupId>
-                    <artifactId>slf4j-log4j12</artifactId>
-                </exclusion>
-                <exclusion>
-                    <groupId>commons-logging</groupId>
-                    <artifactId>commons-logging</artifactId>
-                </exclusion>
-                <exclusion>
-                    <groupId>org.slf4j</groupId>
-                    <artifactId>slf4j-simple</artifactId>
-                </exclusion>
-                <exclusion>
-                    <groupId>org.slf4j</groupId>
-                    <artifactId>log4j-over-slf4j</artifactId>
-                </exclusion>
-                <exclusion>
-                    <groupId>org.apache.zookeeper</groupId>
-                    <artifactId>zookeeper</artifactId>
-                </exclusion>
-            </exclusions>
-        </dependency>
-
-        <dependency>
-            <groupId>org.apache.flink</groupId>
-            <artifactId>flink-test-utils_2.11</artifactId>
-            <version>${flink.version}</version>
-            <scope>test</scope>
-        </dependency>
-
-        <dependency>
-            <groupId>org.apache.ignite</groupId>
-            <artifactId>ignite-core</artifactId>
-            <version>${project.version}</version>
-            <type>test-jar</type>
-            <scope>test</scope>
-        </dependency>
-
-        <dependency>
-            <groupId>org.apache.ignite</groupId>
-            <artifactId>ignite-log4j</artifactId>
-            <version>${project.version}</version>
-            <scope>test</scope>
-        </dependency>
-
-        <dependency>
-            <groupId>org.apache.ignite</groupId>
-            <artifactId>ignite-spring</artifactId>
-            <version>${project.version}</version>
-            <scope>test</scope>
-        </dependency>
-
-        <dependency>
-            <groupId>org.mockito</groupId>
-            <artifactId>mockito-all</artifactId>
-            <version>${mockito.version}</version>
-            <scope>test</scope>
-        </dependency>
-    </dependencies>
-
-</project>
diff --git a/modules/flink/src/main/java/org/apache/ignite/sink/flink/IgniteSink.java b/modules/flink/src/main/java/org/apache/ignite/sink/flink/IgniteSink.java
deleted file mode 100644
index 8deb0d7..0000000
--- a/modules/flink/src/main/java/org/apache/ignite/sink/flink/IgniteSink.java
+++ /dev/null
@@ -1,197 +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.ignite.sink.flink;
-
-import java.util.Map;
-import org.apache.flink.configuration.Configuration;
-import org.apache.flink.streaming.api.functions.sink.RichSinkFunction;
-import org.apache.ignite.Ignite;
-import org.apache.ignite.IgniteDataStreamer;
-import org.apache.ignite.IgniteException;
-import org.apache.ignite.IgniteIllegalStateException;
-import org.apache.ignite.IgniteLogger;
-import org.apache.ignite.Ignition;
-import org.apache.ignite.internal.util.typedef.internal.A;
-
-/**
- * Apache Flink Ignite sink implemented as a RichSinkFunction.
- */
-public class IgniteSink<IN> extends RichSinkFunction<IN> {
-    /** Default flush frequency. */
-    private static final long DFLT_FLUSH_FREQ = 10000L;
-
-    /** Logger. */
-    private transient IgniteLogger log;
-
-    /** Automatic flush frequency. */
-    private long autoFlushFrequency = DFLT_FLUSH_FREQ;
-
-    /** Enables overwriting existing values in cache. */
-    private boolean allowOverwrite = false;
-
-    /** Flag for stopped state. */
-    private volatile boolean stopped = true;
-
-    /** Ignite instance. */
-    protected transient Ignite ignite;
-
-    /** Ignite Data streamer instance. */
-    protected transient IgniteDataStreamer streamer;
-
-    /** Ignite grid configuration file. */
-    protected final String igniteCfgFile;
-
-    /** Cache name. */
-    protected final String cacheName;
-
-    /**
-     * Gets the cache name.
-     *
-     * @return Cache name.
-     */
-    public String getCacheName() {
-        return cacheName;
-    }
-
-    /**
-     * Gets Ignite configuration file.
-     *
-     * @return Configuration file.
-     */
-    public String getIgniteConfigFile() {
-        return igniteCfgFile;
-    }
-
-    /**
-     * Gets the Ignite instance.
-     *
-     * @return Ignite instance.
-     */
-    public Ignite getIgnite() {
-        return ignite;
-    }
-
-    /**
-     * Obtains data flush frequency.
-     *
-     * @return Flush frequency.
-     */
-    public long getAutoFlushFrequency() {
-        return autoFlushFrequency;
-    }
-
-    /**
-     * Specifies data flush frequency into the grid.
-     *
-     * @param autoFlushFrequency Flush frequency.
-     */
-    public void setAutoFlushFrequency(long autoFlushFrequency) {
-        this.autoFlushFrequency = autoFlushFrequency;
-    }
-
-    /**
-     * Obtains flag for enabling overwriting existing values in cache.
-     *
-     * @return True if overwriting is allowed, false otherwise.
-     */
-    public boolean getAllowOverwrite() {
-        return allowOverwrite;
-    }
-
-    /**
-     * Enables overwriting existing values in cache.
-     *
-     * @param allowOverwrite Flag value.
-     */
-    public void setAllowOverwrite(boolean allowOverwrite) {
-        this.allowOverwrite = allowOverwrite;
-    }
-
-    /**
-     * Default IgniteSink constructor.
-     *
-     * @param cacheName Cache name.
-     */
-    public IgniteSink(String cacheName, String igniteCfgFile) {
-        this.cacheName = cacheName;
-        this.igniteCfgFile = igniteCfgFile;
-    }
-
-    /**
-     * Starts streamer.
-     *
-     * @throws IgniteException If failed.
-     */
-    @Override
-    public void open(Configuration parameter) {
-        A.notNull(igniteCfgFile, "Ignite config file");
-        A.notNull(cacheName, "Cache name");
-
-        try {
-            // if an ignite instance is already started in same JVM then use it.
-            this.ignite = Ignition.ignite();
-        } catch (IgniteIllegalStateException e) {
-            this.ignite = Ignition.start(igniteCfgFile);
-        }
-
-        this.ignite.getOrCreateCache(cacheName);
-
-        this.log = this.ignite.log();
-
-        this.streamer = this.ignite.dataStreamer(cacheName);
-        this.streamer.autoFlushFrequency(autoFlushFrequency);
-        this.streamer.allowOverwrite(allowOverwrite);
-
-        stopped = false;
-    }
-
-    /**
-     * Stops streamer.
-     *
-     * @throws IgniteException If failed.
-     */
-    @Override
-    public void close() {
-        if (stopped)
-            return;
-
-        stopped = true;
-
-        this.streamer.close();
-    }
-
-    /**
-     * Transfers data into grid. It is called when new data
-     * arrives to the sink, and forwards it to {@link IgniteDataStreamer}.
-     *
-     * @param in IN.
-     */
-    @SuppressWarnings("unchecked")
-    @Override
-    public void invoke(IN in) {
-        try {
-            if (!(in instanceof Map))
-                throw new IgniteException("Map as a streamer input is expected!");
-
-            this.streamer.addData((Map)in);
-        }
-        catch (Exception e) {
-            log.error("Error while processing IN of " + cacheName, e);
-        }
-    }
-}
diff --git a/modules/flink/src/main/java/org/apache/ignite/sink/flink/package-info.java b/modules/flink/src/main/java/org/apache/ignite/sink/flink/package-info.java
deleted file mode 100644
index 7b1437c..0000000
--- a/modules/flink/src/main/java/org/apache/ignite/sink/flink/package-info.java
+++ /dev/null
@@ -1,22 +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 description. -->
- * IgniteSink -- streaming connector integration with Apache Flink.
- */
-package org.apache.ignite.sink.flink;
diff --git a/modules/flink/src/main/java/org/apache/ignite/source/flink/IgniteSource.java b/modules/flink/src/main/java/org/apache/ignite/source/flink/IgniteSource.java
deleted file mode 100644
index 2dd670a..0000000
--- a/modules/flink/src/main/java/org/apache/ignite/source/flink/IgniteSource.java
+++ /dev/null
@@ -1,223 +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.ignite.source.flink;
-
-import java.util.ArrayList;
-import java.util.List;
-import java.util.UUID;
-import java.util.concurrent.BlockingQueue;
-import java.util.concurrent.LinkedBlockingQueue;
-import java.util.concurrent.TimeUnit;
-import org.apache.flink.streaming.api.functions.source.RichParallelSourceFunction;
-import org.apache.ignite.Ignite;
-import org.apache.ignite.IgniteException;
-import org.apache.ignite.events.CacheEvent;
-import org.apache.ignite.internal.util.typedef.X;
-import org.apache.ignite.internal.util.typedef.internal.A;
-import org.apache.ignite.lang.IgniteBiPredicate;
-import org.apache.ignite.lang.IgnitePredicate;
-import org.apache.ignite.resources.IgniteInstanceResource;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-/**
- * Apache Flink Ignite source implemented as a RichParallelSourceFunction.
- */
-public class IgniteSource extends RichParallelSourceFunction<CacheEvent> {
-    /** Serial version uid. */
-    private static final long serialVersionUID = 1L;
-
-    /** Logger. */
-    private static final Logger log = LoggerFactory.getLogger(IgniteSource.class);
-
-    /** Default max number of events taken from the buffer at once. */
-    private static final int DFLT_EVT_BATCH_SIZE = 1;
-
-    /** Default number of milliseconds timeout for event buffer queue operation. */
-    private static final int DFLT_EVT_BUFFER_TIMEOUT = 10;
-
-    /** Event buffer. */
-    private BlockingQueue<CacheEvent> evtBuf = new LinkedBlockingQueue<>();
-
-    /** Remote Listener id. */
-    private UUID rmtLsnrId;
-
-    /** Flag for isRunning state. */
-    private volatile boolean isRunning;
-
-    /** Max number of events taken from the buffer at once. */
-    private int evtBatchSize = DFLT_EVT_BATCH_SIZE;
-
-    /** Number of milliseconds timeout for event buffer queue operation. */
-    private int evtBufTimeout = DFLT_EVT_BUFFER_TIMEOUT;
-
-    /** Local listener. */
-    private final TaskLocalListener locLsnr = new TaskLocalListener();
-
-    /** Ignite instance. */
-    @IgniteInstanceResource
-    private transient Ignite ignite;
-
-    /** Cache name. */
-    private final String cacheName;
-
-    /**
-     * Sets Ignite instance.
-     *
-     * @param ignite Ignite instance.
-     */
-    public void setIgnite(Ignite ignite) {
-        this.ignite = ignite;
-    }
-
-    /**
-     * Sets Event Batch Size.
-     *
-     * @param evtBatchSize Event Batch Size.
-     */
-    public void setEvtBatchSize(int evtBatchSize) {
-        this.evtBatchSize = evtBatchSize;
-    }
-
-    /**
-     * Sets Event Buffer timeout.
-     *
-     * @param evtBufTimeout Event Buffer timeout.
-     */
-    public void setEvtBufTimeout(int evtBufTimeout) {
-        this.evtBufTimeout = evtBufTimeout;
-    }
-
-    /**
-     * @return Local Task Listener
-     */
-    TaskLocalListener getLocLsnr() {
-        return locLsnr;
-    }
-
-    /**
-     * Default IgniteSource constructor.
-     *
-     * @param cacheName Cache name.
-     */
-    public IgniteSource(String cacheName) {
-        this.cacheName = cacheName;
-    }
-
-    /**
-     * Starts Ignite source.
-     *
-     * @param filter User defined filter.
-     * @param cacheEvts Converts comma-delimited cache events strings to Ignite internal representation.
-     */
-    @SuppressWarnings("unchecked")
-    public void start(IgnitePredicate<CacheEvent> filter, int... cacheEvts) {
-        A.notNull(cacheName, "Cache name");
-
-        TaskRemoteFilter rmtLsnr = new TaskRemoteFilter(cacheName, filter);
-
-        try {
-            synchronized (this) {
-                if (isRunning)
-                    return;
-
-                isRunning = true;
-
-                rmtLsnrId = ignite.events(ignite.cluster().forCacheNodes(cacheName))
-                    .remoteListen(locLsnr, rmtLsnr, cacheEvts);
-            }
-        }
-        catch (IgniteException e) {
-            log.error("Failed to register event listener!", e);
-
-            throw e;
-        }
-    }
-
-    /**
-     * Transfers data from grid.
-     *
-     * @param ctx SourceContext.
-     */
-    @Override public void run(SourceContext<CacheEvent> ctx) {
-        List<CacheEvent> evts = new ArrayList<>(evtBatchSize);
-
-        try {
-            while (isRunning) {
-                // block here for some time if there is no events from source
-                CacheEvent firstEvt = evtBuf.poll(1, TimeUnit.SECONDS);
-
-                if (firstEvt != null)
-                    evts.add(firstEvt);
-
-                if (evtBuf.drainTo(evts, evtBatchSize) > 0) {
-                    synchronized (ctx.getCheckpointLock()) {
-                        for (CacheEvent evt : evts)
-                            ctx.collect(evt);
-
-                        evts.clear();
-                    }
-                }
-            }
-        }
-        catch (Exception e) {
-            if (X.hasCause(e, InterruptedException.class))
-                return; // Executing thread can be interrupted see cancel() javadoc.
-
-            log.error("Error while processing cache event of " + cacheName, e);
-        }
-    }
-
-    /** {@inheritDoc} */
-    @Override public void cancel() {
-        synchronized (this) {
-            if (!isRunning)
-                return;
-
-            isRunning = false;
-
-            if (rmtLsnrId != null && ignite != null) {
-                ignite.events(ignite.cluster().forCacheNodes(cacheName))
-                    .stopRemoteListen(rmtLsnrId);
-
-                rmtLsnrId = null;
-            }
-        }
-    }
-
-    /**
-     * Local listener buffering cache events to be further sent to Flink.
-     */
-    private class TaskLocalListener implements IgniteBiPredicate<UUID, CacheEvent> {
-        /** {@inheritDoc} */
-        @Override public boolean apply(UUID id, CacheEvent evt) {
-            try {
-                if (!evtBuf.offer(evt, evtBufTimeout, TimeUnit.MILLISECONDS))
-                    log.error("Failed to buffer event {}", evt.name());
-            }
-            catch (InterruptedException ignored) {
-                log.error("Failed to buffer event using local task listener {}", evt.name());
-
-                Thread.currentThread().interrupt(); // Restore interrupt flag.
-            }
-
-            return true;
-        }
-    }
-}
-
diff --git a/modules/flink/src/main/java/org/apache/ignite/source/flink/TaskRemoteFilter.java b/modules/flink/src/main/java/org/apache/ignite/source/flink/TaskRemoteFilter.java
deleted file mode 100644
index 4c89d25..0000000
--- a/modules/flink/src/main/java/org/apache/ignite/source/flink/TaskRemoteFilter.java
+++ /dev/null
@@ -1,60 +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.ignite.source.flink;
-
-import org.apache.ignite.Ignite;
-import org.apache.ignite.cache.affinity.Affinity;
-import org.apache.ignite.events.CacheEvent;
-import org.apache.ignite.lang.IgnitePredicate;
-import org.apache.ignite.resources.IgniteInstanceResource;
-
-/**
- * Remote filter.
- */
-public class TaskRemoteFilter implements IgnitePredicate<CacheEvent> {
-    /** Serial version Id. */
-    private static final long serialVersionUID = 1L;
-
-    /** Ignite Instance Resource. */
-    @IgniteInstanceResource
-    private Ignite ignite;
-
-    /** Cache name. */
-    private final String cacheName;
-
-    /** User-defined filter. */
-    private final IgnitePredicate<CacheEvent> filter;
-
-    /**
-     * @param cacheName Cache name.
-     * @param filter IgnitePredicate.
-     */
-    TaskRemoteFilter(String cacheName, IgnitePredicate<CacheEvent> filter) {
-        this.cacheName = cacheName;
-        this.filter = filter;
-    }
-
-    /** {@inheritDoc} */
-    @Override public boolean apply(CacheEvent evt) {
-        Affinity<Object> affinity = ignite.affinity(cacheName);
-
-        // Process this event. Ignored on backups.
-        return affinity.isPrimary(ignite.cluster().localNode(), evt.key()) &&
-                (filter == null || filter.apply(evt));
-    }
-}
diff --git a/modules/flink/src/main/java/org/apache/ignite/source/flink/package-info.java b/modules/flink/src/main/java/org/apache/ignite/source/flink/package-info.java
deleted file mode 100644
index adc33fc..0000000
--- a/modules/flink/src/main/java/org/apache/ignite/source/flink/package-info.java
+++ /dev/null
@@ -1,21 +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.
- */
-
-/**
- * IgniteSource -- source connector integration with Apache Flink.
- */
-package org.apache.ignite.source.flink;
diff --git a/modules/flink/src/test/java/org/apache/ignite/sink/flink/FlinkIgniteSinkSelfTest.java b/modules/flink/src/test/java/org/apache/ignite/sink/flink/FlinkIgniteSinkSelfTest.java
deleted file mode 100644
index 25c8950..0000000
--- a/modules/flink/src/test/java/org/apache/ignite/sink/flink/FlinkIgniteSinkSelfTest.java
+++ /dev/null
@@ -1,84 +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.ignite.sink.flink;
-
-import java.util.HashMap;
-import java.util.Map;
-import org.apache.flink.configuration.Configuration;
-import org.apache.flink.streaming.api.datastream.DataStream;
-import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
-import org.apache.ignite.testframework.junits.common.GridCommonAbstractTest;
-import org.junit.Test;
-
-/**
- * Tests for {@link IgniteSink}.
- */
-public class FlinkIgniteSinkSelfTest extends GridCommonAbstractTest {
-    /** Cache name. */
-    private static final String TEST_CACHE = "testCache";
-
-    /** Ignite test configuration file. */
-    private static final String GRID_CONF_FILE = "modules/flink/src/test/resources/example-ignite.xml";
-
-    @Test
-    public void testIgniteSink() throws Exception {
-        Configuration configuration = new Configuration();
-
-        IgniteSink igniteSink = new IgniteSink(TEST_CACHE, GRID_CONF_FILE);
-
-        igniteSink.setAllowOverwrite(true);
-
-        igniteSink.setAutoFlushFrequency(1L);
-
-        igniteSink.open(configuration);
-
-        Map<String, String> myData = new HashMap<>();
-        myData.put("testData", "testValue");
-
-        igniteSink.invoke(myData);
-
-        /** waiting for a small duration for the cache flush to complete */
-        Thread.sleep(2000);
-
-        assertEquals("testValue", igniteSink.getIgnite().getOrCreateCache(TEST_CACHE).get("testData"));
-    }
-
-    @Test
-    public void testIgniteSinkStreamExecution() throws Exception {
-        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
-
-        IgniteSink igniteSink = new IgniteSink(TEST_CACHE, GRID_CONF_FILE);
-
-        igniteSink.setAllowOverwrite(true);
-
-        igniteSink.setAutoFlushFrequency(1);
-
-        Map<String, String> myData = new HashMap<>();
-        myData.put("testdata", "testValue");
-        DataStream<Map> stream = env.fromElements(myData);
-
-        stream.addSink(igniteSink);
-        try {
-            env.execute();
-        }
-        catch (Exception e) {
-            e.printStackTrace();
-            fail("Stream execution process failed.");
-        }
-    }
-}
diff --git a/modules/flink/src/test/java/org/apache/ignite/sink/flink/FlinkIgniteSinkSelfTestSuite.java b/modules/flink/src/test/java/org/apache/ignite/sink/flink/FlinkIgniteSinkSelfTestSuite.java
deleted file mode 100644
index 7890e4c..0000000
--- a/modules/flink/src/test/java/org/apache/ignite/sink/flink/FlinkIgniteSinkSelfTestSuite.java
+++ /dev/null
@@ -1,29 +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.ignite.sink.flink;
-
-import org.junit.runner.RunWith;
-import org.junit.runners.Suite;
-
-/**
- * Apache Flink sink tests.
- */
-@RunWith(Suite.class)
-@Suite.SuiteClasses({FlinkIgniteSinkSelfTest.class})
-public class FlinkIgniteSinkSelfTestSuite {
-}
diff --git a/modules/flink/src/test/java/org/apache/ignite/source/flink/FlinkIgniteSourceSelfTest.java b/modules/flink/src/test/java/org/apache/ignite/source/flink/FlinkIgniteSourceSelfTest.java
deleted file mode 100644
index f59007e..0000000
--- a/modules/flink/src/test/java/org/apache/ignite/source/flink/FlinkIgniteSourceSelfTest.java
+++ /dev/null
@@ -1,154 +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.ignite.source.flink;
-
-import java.util.UUID;
-import org.apache.flink.api.common.functions.RuntimeContext;
-import org.apache.flink.streaming.api.functions.source.SourceFunction;
-import org.apache.flink.streaming.api.operators.StreamingRuntimeContext;
-import org.apache.ignite.Ignite;
-import org.apache.ignite.IgniteCluster;
-import org.apache.ignite.IgniteEvents;
-import org.apache.ignite.cluster.ClusterGroup;
-import org.apache.ignite.events.CacheEvent;
-import org.apache.ignite.events.EventType;
-import org.apache.ignite.internal.IgniteInternalFuture;
-import org.apache.ignite.internal.util.lang.GridAbsPredicate;
-import org.apache.ignite.lang.IgniteBiPredicate;
-import org.apache.ignite.testframework.GridTestUtils;
-import org.apache.ignite.testframework.junits.common.GridCommonAbstractTest;
-import org.junit.After;
-import org.junit.Before;
-import org.junit.Test;
-
-import static org.mockito.Mockito.mock;
-import static org.mockito.Mockito.times;
-import static org.mockito.Mockito.verify;
-import static org.mockito.Mockito.when;
-
-/**
- * Tests for {@link IgniteSource}.
- */
-public class FlinkIgniteSourceSelfTest extends GridCommonAbstractTest {
-    /** Cache name. */
-    private static final String TEST_CACHE = "testCache";
-
-    /** Flink source context. */
-    private SourceFunction.SourceContext<CacheEvent> ctx;
-
-    /** Ignite instance. */
-    private Ignite ignite;
-
-    /** Cluster Group */
-    private ClusterGroup clsGrp;
-
-    /** Ignite Source instance */
-    private IgniteSource igniteSrc;
-
-    /** */
-    @SuppressWarnings("unchecked")
-    @Before
-    public void setUpTest() throws Exception {
-        ctx = mock(SourceFunction.SourceContext.class);
-        ignite = mock(Ignite.class);
-        clsGrp = mock(ClusterGroup.class);
-
-        IgniteEvents igniteEvts = mock(IgniteEvents.class);
-        IgniteCluster igniteCluster = mock(IgniteCluster.class);
-        TaskRemoteFilter taskRemoteFilter = mock(TaskRemoteFilter.class);
-
-        when(ctx.getCheckpointLock()).thenReturn(new Object());
-        when(ignite.events(clsGrp)).thenReturn(igniteEvts);
-        when(ignite.cluster()).thenReturn(igniteCluster);
-
-        igniteSrc = new IgniteSource(TEST_CACHE);
-        igniteSrc.setIgnite(ignite);
-        igniteSrc.setEvtBatchSize(1);
-        igniteSrc.setEvtBufTimeout(1);
-        igniteSrc.setRuntimeContext(createRuntimeContext());
-
-        IgniteBiPredicate locLsnr = igniteSrc.getLocLsnr();
-
-        when(igniteEvts.remoteListen(locLsnr, taskRemoteFilter, EventType.EVT_CACHE_OBJECT_PUT ))
-            .thenReturn(UUID.randomUUID());
-
-        when(igniteCluster.forCacheNodes(TEST_CACHE)).thenReturn(clsGrp);
-    }
-
-    /**  */
-    @After
-    public void tearDownTest() {
-        igniteSrc.cancel();
-    }
-
-    /** Creates streaming runtime context */
-    private RuntimeContext createRuntimeContext() {
-        StreamingRuntimeContext runtimeCtx = mock(StreamingRuntimeContext.class);
-
-        when(runtimeCtx.isCheckpointingEnabled()).thenReturn(true);
-
-        return runtimeCtx;
-    }
-
-    /**
-     * Tests Ignite source start operation.
-     *
-     * @throws Exception If failed.
-     */
-    @Test
-    public void testIgniteSourceStart() throws Exception {
-        igniteSrc.start(null, EventType.EVT_CACHE_OBJECT_PUT);
-
-        verify(ignite.events(clsGrp), times(1));
-    }
-
-    /**
-     * Tests Ignite source run operation.
-     *
-     * @throws Exception If failed.
-     */
-    @Test
-    public void testIgniteSourceRun() throws Exception {
-        IgniteInternalFuture f = GridTestUtils.runAsync(new Runnable() {
-            @Override public void run() {
-                try {
-                    igniteSrc.start(null, EventType.EVT_CACHE_OBJECT_PUT);
-
-                    igniteSrc.run(ctx);
-                }
-                catch (Throwable e) {
-                    igniteSrc.cancel();
-
-                   throw new AssertionError("Unexpected failure.", e);
-                }
-            }
-        });
-
-        long endTime = System.currentTimeMillis() + 2000;
-
-        GridTestUtils.waitForCondition(new GridAbsPredicate() {
-            @Override public boolean apply() {
-                return f.isDone() || System.currentTimeMillis() > endTime;
-            }
-        }, 3000);
-
-        igniteSrc.cancel();
-
-        f.get(3000);
-    }
-}
diff --git a/modules/flink/src/test/java/org/apache/ignite/source/flink/FlinkIgniteSourceSelfTestSuite.java b/modules/flink/src/test/java/org/apache/ignite/source/flink/FlinkIgniteSourceSelfTestSuite.java
deleted file mode 100644
index 7070402..0000000
--- a/modules/flink/src/test/java/org/apache/ignite/source/flink/FlinkIgniteSourceSelfTestSuite.java
+++ /dev/null
@@ -1,30 +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.ignite.source.flink;
-
-import org.junit.runner.RunWith;
-import org.junit.runners.Suite;
-
-/**
- * Apache Flink source tests.
- */
-@RunWith(Suite.class)
-@Suite.SuiteClasses({FlinkIgniteSourceSelfTest.class})
-public class FlinkIgniteSourceSelfTestSuite {
-}
-
diff --git a/modules/flink/src/test/resources/example-ignite.xml b/modules/flink/src/test/resources/example-ignite.xml
deleted file mode 100644
index d4f4dc1..0000000
--- a/modules/flink/src/test/resources/example-ignite.xml
+++ /dev/null
@@ -1,73 +0,0 @@
-<?xml version="1.0" encoding="UTF-8"?>
-
-<!--
-  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.
--->
-
-<!--
-    Ignite configuration with all defaults and enabled events.
-    Used for testing IgniteSink running Ignite in a client mode.
--->
-<beans xmlns="http://www.springframework.org/schema/beans"
-       xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
-       xmlns:util="http://www.springframework.org/schema/util"
-       xsi:schemaLocation="
-        http://www.springframework.org/schema/beans
-        http://www.springframework.org/schema/beans/spring-beans.xsd
-        http://www.springframework.org/schema/util
-        http://www.springframework.org/schema/util/spring-util.xsd">
-    <bean id="ignite.cfg" class="org.apache.ignite.configuration.IgniteConfiguration">
-        <!-- Enable client mode. -->
-        <property name="clientMode" value="false"/>
-
-        <!-- Cache accessed from IgniteSink. -->
-        <property name="cacheConfiguration">
-            <list>
-                <!-- Partitioned cache example configuration with configurations adjusted to server nodes'. -->
-                <bean class="org.apache.ignite.configuration.CacheConfiguration">
-                    <property name="atomicityMode" value="ATOMIC"/>
-                    <property name="name" value="testCache"/>
-                </bean>
-            </list>
-        </property>
-
-        <!-- Enable cache events. -->
-        <property name="includeEventTypes">
-            <list>
-                <!-- Cache events (only EVT_CACHE_OBJECT_PUT for tests). -->
-                <util:constant static-field="org.apache.ignite.events.EventType.EVT_CACHE_OBJECT_PUT"/>
-                <util:constant static-field="org.apache.ignite.events.EventType.EVT_CACHE_OBJECT_READ"/>
-                <util:constant static-field="org.apache.ignite.events.EventType.EVT_CACHE_OBJECT_REMOVED"/>
-
-            </list>
-        </property>
-
-        <!-- Explicitly configure TCP discovery SPI to provide list of initial nodes. -->
-        <property name="discoverySpi">
-            <bean class="org.apache.ignite.spi.discovery.tcp.TcpDiscoverySpi">
-                <property name="ipFinder">
-                    <bean class="org.apache.ignite.spi.discovery.tcp.ipfinder.vm.TcpDiscoveryVmIpFinder">
-                        <property name="addresses">
-                            <list>
-                                <value>127.0.0.1:47500..47509</value>
-                            </list>
-                        </property>
-                    </bean>
-                </property>
-            </bean>
-        </property>
-    </bean>
-</beans>
diff --git a/parent/pom.xml b/parent/pom.xml
index 27e3042..fe963f9 100644
--- a/parent/pom.xml
+++ b/parent/pom.xml
@@ -495,14 +495,6 @@
                                 <packages>org.apache.ignite.osgi*</packages>
                             </group>
                             <group>
-                                <title>Flink Sink Integration</title>
-                                <packages>org.apache.ignite.sink.flink*</packages>
-                            </group>
-                            <group>
-                                <title>Flink Source Integration</title>
-                                <packages>org.apache.ignite.source.flink*</packages>
-                            </group>
-                            <group>
                                 <title>SpringData integration</title>
                                 <packages>org.apache.ignite.springdata.repository*</packages>
                             </group>
diff --git a/pom.xml b/pom.xml
index 33f6d21..15ce173 100644
--- a/pom.xml
+++ b/pom.xml
@@ -89,7 +89,6 @@
         <module>modules/web/ignite-appserver-test</module>
         <module>modules/web/ignite-websphere-test</module>
         <module>modules/cassandra</module>
-        <module>modules/flink</module>
         <module>modules/kubernetes</module>
         <module>modules/zeromq</module>
         <module>modules/rocketmq</module>