5: win10 使用IDEA调试编译flink 1.14
1: 下载
https://dlcdn.apache.org/flink/flink-1.14.2/flink-1.14.2-bin-scala_2.12.tgz
2: JAVA, Maven, IDEA都安装好。
3:搞一个flink的官网的程序
创建一个maven 工程
$ mvn archetype:generate \
-DarchetypeGroupId=org.apache.flink \
-DarchetypeArtifactId=flink-walkthrough-datastream-java \
-DarchetypeVersion=1.14.2 \
-DgroupId=frauddetection \
-DartifactId=frauddetection \
-Dversion=0.1 \
-Dpackage=spendreport \
-DinteractiveMode=false
pom.xml
<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.0modelVersion> <groupId>frauddetectiongroupId> <artifactId>frauddetectionartifactId> <version>0.1version> <packaging>jarpackaging> <name>Flink Walkthrough DataStream Javaname> <url>https://flink.apache.orgurl> <properties> <project.build.sourceEncoding>UTF-8project.build.sourceEncoding> <flink.version>1.14.2flink.version> <target.java.version>1.8target.java.version> <scala.binary.version>2.11scala.binary.version> <maven.compiler.source>${target.java.version}maven.compiler.source> <maven.compiler.target>${target.java.version}maven.compiler.target> <log4j.version>2.16.0log4j.version> properties> <repositories> <repository> <id>apache.snapshotsid> <name>Apache Development Snapshot Repositoryname> <url>https://repository.apache.org/content/repositories/snapshots/url> <releases> <enabled>falseenabled> releases> <snapshots> <enabled>trueenabled> snapshots> repository> repositories> <dependencies> <dependency> <groupId>org.apache.flinkgroupId> <artifactId>flink-walkthrough-common_${scala.binary.version}artifactId> <version>${flink.version}version> dependency> <dependency> <groupId>org.apache.flinkgroupId> <artifactId>flink-streaming-java_${scala.binary.version}artifactId> <version>${flink.version}version> <scope>providedscope> dependency> <dependency> <groupId>org.apache.flinkgroupId> <artifactId>flink-clients_${scala.binary.version}artifactId> <version>${flink.version}version> <scope>providedscope> dependency> <dependency> <groupId>org.apache.logging.log4jgroupId> <artifactId>log4j-slf4j-implartifactId> <version>${log4j.version}version> <scope>runtimescope> dependency> <dependency> <groupId>org.apache.logging.log4jgroupId> <artifactId>log4j-apiartifactId> <version>${log4j.version}version> <scope>runtimescope> dependency> <dependency> <groupId>org.apache.logging.log4jgroupId> <artifactId>log4j-coreartifactId> <version>${log4j.version}version> <scope>runtimescope> dependency> dependencies> <build> <plugins> <plugin> <groupId>org.apache.maven.pluginsgroupId> <artifactId>maven-compiler-pluginartifactId> <version>3.1version> <configuration> <source>${target.java.version}source> <target>${target.java.version}target> configuration> plugin> <plugin> <groupId>org.apache.maven.pluginsgroupId> <artifactId>maven-shade-pluginartifactId> <version>3.0.0version> <executions> <execution> <phase>packagephase> <goals> <goal>shadegoal> goals> <configuration> <artifactSet> <excludes> <exclude>org.apache.flink:flink-shaded-force-shadingexclude> <exclude>com.google.code.findbugs:jsr305exclude> <exclude>org.slf4j:*exclude> <exclude>org.apache.logging.log4j:*exclude> excludes> artifactSet> <filters> <filter> <artifact>*:*artifact> <excludes> <exclude>META-INF/*.SFexclude> <exclude>META-INF/*.DSAexclude> <exclude>META-INF/*.RSAexclude> excludes> filter> filters> <transformers> <transformer implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer"> <mainClass>spendreport.FraudDetectionJobmainClass> transformer> transformers> configuration> execution> executions> plugin> plugins> <pluginManagement> <plugins> <plugin> <groupId>org.eclipse.m2egroupId> <artifactId>lifecycle-mappingartifactId> <version>1.0.0version> <configuration> <lifecycleMappingMetadata> <pluginExecutions> <pluginExecution> <pluginExecutionFilter> <groupId>org.apache.maven.pluginsgroupId> <artifactId>maven-shade-pluginartifactId> <versionRange>[3.0.0,)versionRange> <goals> <goal>shadegoal> goals> pluginExecutionFilter> <action> <ignore/> action> pluginExecution> <pluginExecution> <pluginExecutionFilter> <groupId>org.apache.maven.pluginsgroupId> <artifactId>maven-compiler-pluginartifactId> <versionRange>[3.1,)versionRange> <goals> <goal>testCompilegoal> <goal>compilegoal> goals> pluginExecutionFilter> <action> <ignore/> action> pluginExecution> pluginExecutions> lifecycleMappingMetadata> configuration> plugin> plugins> pluginManagement> build> project>
FraudDetectionJob.java
/* * 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 spendreport; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.walkthrough.common.sink.AlertSink; import org.apache.flink.walkthrough.common.entity.Alert; import org.apache.flink.walkthrough.common.entity.Transaction; import org.apache.flink.walkthrough.common.source.TransactionSource; /** * Skeleton code for the datastream walkthrough */ public class FraudDetectionJob { public static void main(String[] args) { StreamExecutionEnvironment env = null; try { env = StreamExecutionEnvironment.getExecutionEnvironment(); DataStreamtransactions = env .addSource(new TransactionSource()) .name("transactions"); DataStream alerts = transactions .keyBy(Transaction::getAccountId) .process(new FraudDetector()) .name("fraud-detector"); alerts .addSink(new AlertSink()) .name("send-alerts"); env.execute("Fraud Detection"); } catch (java.lang.Exception e) { e.printStackTrace(); } } }
FraudDetector.java
package spendreport; import org.apache.flink.api.common.state.ValueState; import org.apache.flink.api.common.state.ValueStateDescriptor; import org.apache.flink.api.common.typeinfo.Types; import org.apache.flink.configuration.Configuration; import org.apache.flink.streaming.api.functions.KeyedProcessFunction; import org.apache.flink.util.Collector; import org.apache.flink.walkthrough.common.entity.Alert; import org.apache.flink.walkthrough.common.entity.Transaction; public class FraudDetector extends KeyedProcessFunction{ private static final long serialVersionUID = 1L; private static final double SMALL_AMOUNT = 1.00; private static final double LARGE_AMOUNT = 500.00; private static final long ONE_MINUTE = 60 * 1000; private transient ValueState flagState; private transient ValueState timerState; @Override public void open(Configuration parameters) { ValueStateDescriptor flagDescriptor = new ValueStateDescriptor<>( "flag", Types.BOOLEAN); flagState = getRuntimeContext().getState(flagDescriptor); ValueStateDescriptor timerDescriptor = new ValueStateDescriptor<>( "timer-state", Types.LONG); timerState = getRuntimeContext().getState(timerDescriptor); } @Override public void processElement( Transaction transaction, Context context, Collector collector) throws Exception { // Get the current state for the current key Boolean lastTransactionWasSmall = flagState.value(); // Check if the flag is set if (lastTransactionWasSmall != null) { if (transaction.getAmount() > LARGE_AMOUNT) { //Output an alert downstream Alert alert = new Alert(); alert.setId(transaction.getAccountId()); collector.collect(alert); } // Clean up our state cleanUp(context); } if (transaction.getAmount() < SMALL_AMOUNT) { // set the flag to true flagState.update(true); long timer = context.timerService().currentProcessingTime() + ONE_MINUTE; context.timerService().registerProcessingTimeTimer(timer); timerState.update(timer); } } @Override public void onTimer(long timestamp, OnTimerContext ctx, Collector out) { // remove flag after 1 minute timerState.clear(); flagState.clear(); } private void cleanUp(Context ctx) throws Exception { // delete timer Long timer = timerState.value(); ctx.timerService().deleteProcessingTimeTimer(timer); // clean up all state timerState.clear(); flagState.clear(); } }
4: IDEA里导入flink libarary
导入flink的lib和op连个包。
然后就可以Idea 运行和调试