栏目分类:
子分类:
返回
名师互学网用户登录
快速导航关闭
当前搜索
当前分类
子分类
实用工具
热门搜索
名师互学网 > IT > 前沿技术 > 大数据 > 大数据系统

flink入门(wordcount)

flink入门(wordcount)

Flink快速上手

1.在IDEA创建maven工程FlinkTutorial

2.在pom.xml中添加依赖和maven插件

 
        
            org.apache.flink
            flink-scala_2.11
            1.10.2
        
        
        
            org.apache.flink
            flink-streaming-scala_2.11
            1.10.2
        
     	
            org.apache.flink
            flink-connector-kafka-0.11_2.12
            1.10.1
        
        
            org.apache.bahir
            flink-connector-redis_2.11
            1.0
            
                
                    org.apache.flink
                    flink-streaming-java_2.11
                
            
        
        
            org.apache.flink
            flink-connector-elasticsearch6_2.12
            1.10.1
        
        
            mysql
            mysql-connector-java
            5.1.44
        
        
            org.apache.flink
            flink-statebackend-rocksdb_2.12
            1.10.1
        
        
            org.apache.flink
            flink-table-planner_2.12
            1.10.1
        
        
            org.apache.flink
            flink-table-planner-blink_2.12
            1.10.1
        
        
            org.apache.flink
            flink-csv
            1.10.1
        
        
            org.apache.flink
            flink-json
            1.10.1
        
 

    
    
    
    net.alchim31.maven
    scala-maven-plugin
    3.4.6
    
        
        
        
            compile
        
         
    
    
    
    org.apache.maven.plugins
    maven-assembly-plugin
    3.0.0
    
        
            jar-with-dependencies
        
    
    
        
            make-assembly
            package
            
                single
            
        
    
    
    

在项目添加scala支持

批处理wordcount
import org.apache.flink.api.scala.ExecutionEnvironment
import org.apache.flink.api.scala._

object WordCount01 {
  def main(args: Array[String]): Unit = {
    // 创建一个批处理环境
    val env = ExecutionEnvironment.getExecutionEnvironment
    // 从文件中读取数据
 val inputData = env.readTextFile("FlinkBasis/src/resources/hello.txt")
    val data:DataSet[String] = env.readTextFile(inputpath)
    val result = data.flatMap(_.split(" ")).map((_,1)).
      groupBy(0) //以第一个字段分组
      .sum(1)  //对第二个子端求和
    result.print()
  }
}

流处理wordcount
import org.apache.flink.api.java.utils.ParameterTool
import org.apache.flink.streaming.api.scala._

object StreamWordCount {
  def main(args: Array[String]): Unit = {
    val env = StreamExecutionEnvironment.getExecutionEnvironment
    // 设置并行线程,开发环境中设置并行度是自己电脑的核心数
    //    env.setParallelism(2)
    // 从外部命令中提取参数,作为socket主机名和端口号
//    val paramTool: ParameterTool = ParameterTool.fromArgs(args)
//    val host: String = paramTool.get("host")
//    val port: Int = paramTool.getInt("port")
    //接受socket文本流
    val data = env.socketTextStream("localhost", 7777)
    //进行wordcount
    val result = data.flatMap(_.split(" ")).filter(_.nonEmpty).map((_, 1)).keyBy(0).sum(1)
//    result.print()
    //如果不希望出现线程index就在print语句中设置并行度为1
    result.print().setParallelism(1)
    env.execute("stream word count")
  }
}

在本地windows安装nc,cmd执行nc -l -p 7777或者在本地linux(微软商店的ubuntu)使用nc -lk 7777

发送数据

转载请注明:文章转载自 www.mshxw.com
本文地址:https://www.mshxw.com/it/663842.html
我们一直用心在做
关于我们 文章归档 网站地图 联系我们

版权所有 (c)2021-2022 MSHXW.COM

ICP备案号:晋ICP备2021003244-6号