这是我的代码
import org.apache.flink.streaming.api.scala._
import org.apache.flink.api.java.utils.ParameterTool
object Demo03WCStream {
def main(args: Array[String]): Unit = {
// 从外部命令中获取参数
val params: ParameterTool = ParameterTool.fromArgs(args)
val host: String = params.get("host")
val port: Int = params.getInt("port")
//创建流式处理API
val env = StreamExecutionEnvironment.getExecutionEnvironment
//设置全局并行度为1
env.setParallelism(1)
//禁用job chain
env.disableOperatorChaining()
val ds: DataStream[String] = env.socketTextStream(host, port)
val ds01: DataStream[(String, Int)] =
ds.flatMap(_.split(" ")).setParallelism(2)
.map((_, 1)).setParallelism(2).disableChaining()
.keyBy(0)
.sum(1).setParallelism(2)
ds01.print("wc")
env.execute("wc job")
}
}
这是报错
Failed to execute goal org.apache.maven.plugins:maven-clean-plugin:2.5:clean (default-clean) on project flink1.12: Failed to clean project: Failed to delete E:\IDEA\flink1.12\target