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>