美文网首页
初学flink-02-流式wordCount

初学flink-02-流式wordCount

作者: 默言少语 | 来源:发表于2019-12-09 15:34 被阅读0次

    pom

    <?xml version="1.0" encoding="UTF-8"?>
    <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>
    
        <groupId>com.stt.flink</groupId>
        <artifactId>flink-demo</artifactId>
        <version>1.0-SNAPSHOT</version>
    
    
        <dependencies>
            <dependency>
                <groupId>org.apache.flink</groupId>
                <artifactId>flink-scala_2.11</artifactId>
                <version>1.7.2</version>
            </dependency>
            <!-- https://mvnrepository.com/artifact/org.apache.flink/flink-streaming-scala -->
            <dependency>
                <groupId>org.apache.flink</groupId>
                <artifactId>flink-streaming-scala_2.11</artifactId>
                <version>1.7.2</version>
            </dependency>
        </dependencies>
    
    
        <build>
            <plugins>
                <!-- 该插件用于将Scala代码编译成class文件 -->
                <plugin>
                    <groupId>net.alchim31.maven</groupId>
                    <artifactId>scala-maven-plugin</artifactId>
                    <version>3.4.6</version>
                    <executions>
                        <execution>
                            <!-- 声明绑定到maven的compile阶段 -->
                            <goals>
                                <goal>testCompile</goal>
                            </goals>
                        </execution>
                    </executions>
                </plugin>
                <plugin>
                    <groupId>org.apache.maven.plugins</groupId>
                    <artifactId>maven-assembly-plugin</artifactId>
                    <version>3.0.0</version>
                    <configuration>
                        <descriptorRefs>
                            <descriptorRef>jar-with-dependencies</descriptorRef>
                        </descriptorRefs>
                    </configuration>
                    <executions>
                        <execution>
                            <id>make-assembly</id>
                            <phase>package</phase>
                            <goals>
                                <goal>single</goal>
                            </goals>
                        </execution>
                    </executions>
                </plugin>
            </plugins>
        </build>
    </project>
    

    scala

    package com.stt.flink
    
    import org.apache.flink.streaming.api.scala.{DataStream, StreamExecutionEnvironment}
    import org.apache.flink.api.scala._ // 隐式转换需要
    
    // 流处理wordCount程序
    object Ch_02_streamWordCount {
      def main(args: Array[String]): Unit = {
        // 创建流式处理的执行环境
        val env: StreamExecutionEnvironment = StreamExecutionEnvironment.getExecutionEnvironment
    
        // 接收一个socket文本流
        val socketDataStream: DataStream[String] = env.socketTextStream("node01",8888)
    
        // 对每条数据进行处理
        val wordCountDataStream: DataStream[(String, Int)] = socketDataStream
          .flatMap(_.split("\\s+"))
          .filter(_.nonEmpty)
          .map((_, 1))
          .keyBy(0) // DataSet 有groupBy操作,在DataStream中有keyBy操作相似的效果
          .sum(1)
    
        wordCountDataStream.print()
    
        // 启动executor
        env.execute()
      }
    }
    

    启动nc服务

    [ttshe@node01~]$ nc -lk 8888
    hello flink
    flink
    hello
    
    • 结果
      • 前面的标号表示slot的索引
    SLF4J: Failed to load class "org.slf4j.impl.StaticLoggerBinder".
    SLF4J: Defaulting to no-operation (NOP) logger implementation
    SLF4J: See http://www.slf4j.org/codes.html#StaticLoggerBinder for further details.
    4> (hello,1)
    10> (flink,1)
    10> (flink,2)
    4> (hello,2)
    

    相关文章

      网友评论

          本文标题:初学flink-02-流式wordCount

          本文链接:https://www.haomeiwen.com/subject/tegbgctx.html