文章详情

短信预约-IT技能 免费直播动态提醒

请输入下面的图形验证码

提交验证

短信预约提醒成功

flink连接消费kafka实例

2023-06-02 16:44

关注

这篇文章主要讲解了“flink连接消费kafka实例”,文中的讲解内容简单清晰,易于学习与理解,下面请大家跟着小编的思路慢慢深入,一起来研究和学习“flink连接消费kafka实例”吧!

package flink.streamingimport java.util.Propertiesimport org.apache.flink.streaming.api.scala._import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumerimport org.apache.flink.api.common.serialization.SimpleStringSchemaimport org.apache.flink.streaming.api.windowing.time.Timeobject StreamingTest {  def main(args: Array[String]): Unit = {    val kafkaProps = new Properties()    //kafka的一些属性    kafkaProps.setProperty("bootstrap.servers", "bigdata01:9092")    //所在的消费组    kafkaProps.setProperty("group.id", "group1")    //获取当前的执行环境    val evn = StreamExecutionEnvironment.getExecutionEnvironment    //kafka的consumer,test1是要消费的topic    val kafkaSource = new FlinkKafkaConsumer[String]("test1",new SimpleStringSchema,kafkaProps)    //设置从最新的offset开始消费    kafkaSource.setStartFromLatest()    //自动提交offset    kafkaSource.setCommitOffsetsOnCheckpoints(true)        //flink的checkpoint的时间间隔    evn.enableCheckpointing(5000)    //添加consumer    val stream = evn.addSource(kafkaSource)    stream.setParallelism(3)    val text = stream.flatMap{ _.toLowerCase().split("\\W+")filter { _.nonEmpty} }          .map{(_,1)}          .keyBy(0)          .timeWindow(Time.seconds(5))          .sum(1)         text.print()     //启动执行        evn.execute("kafkawd")                      }}

//

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.0</modelVersion>  <groupId>hgs</groupId>  <artifactId>flink_lesson</artifactId>  <version>1.0.0</version>  <packaging>jar</packaging>  <name>flink_lesson</name>  <url>http://maven.apache.org</url>  <properties>    <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>  </properties>  <dependencies>    <dependency>      <groupId>junit</groupId>      <artifactId>junit</artifactId>      <version>4.1</version>      <scope>test</scope>    </dependency>        <!-- https://mvnrepository.com/artifact/org.apache.flink/flink-core --><dependency>    <groupId>org.apache.flink</groupId>    <artifactId>flink-core</artifactId>    <version>1.7.1</version></dependency><!-- https://mvnrepository.com/artifact/org.apache.flink/flink-clients --><dependency>    <groupId>org.apache.flink</groupId>    <artifactId>flink-clients_2.12</artifactId>    <version>1.7.1</version></dependency><!-- https://mvnrepository.com/artifact/org.apache.flink/flink-streaming-scala --><dependency>    <groupId>org.apache.flink</groupId>    <artifactId>flink-streaming-scala_2.12</artifactId>    <version>1.7.1</version></dependency><!-- https://mvnrepository.com/artifact/org.apache.flink/flink-connector-kafka-0.11 --><dependency>  <groupId>org.apache.flink</groupId>  <artifactId>flink-connector-kafka_2.12</artifactId>  <version>1.7.1</version></dependency><!-- https://mvnrepository.com/artifact/io.netty/netty-all --><dependency>    <groupId>io.netty</groupId>    <artifactId>netty-all</artifactId>    <version>4.1.32.Final</version></dependency>  </dependencies>       <build>        <plugins>            <plugin>                <artifactId>maven-assembly-plugin</artifactId>                <version>2.6</version>                <configuration>                             <archive>                        <manifest>                            <!-- 我运行这个jar所运行的主类 -->                            <mainClass>hgs.flink_lesson.WordCount</mainClass>                        </manifest>                    </archive>                                         <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>                          <plugin>                <groupId>org.apache.maven.plugins</groupId>                <artifactId>maven-compiler-plugin</artifactId>                <configuration>                    <source>1.8</source>                    <target>1.8</target>                </configuration>            </plugin>              <plugin><groupId>net.alchim31.maven</groupId><artifactId>scala-maven-plugin</artifactId><version>3.2.0</version><executions><execution><goals><goal>compile</goal><goal>testCompile</goal>    </goals><configuration><args><!-- <arg>-make:transitive</arg> -->                <arg>-dependencyfile</arg>                <arg>${project.build.directory}/.scala_dependencies</arg>              </args></configuration></execution></executions></plugin><plugin><groupId>org.apache.maven.plugins</groupId><artifactId>maven-surefire-plugin</artifactId><version>2.18.1</version><configuration><useFile>false</useFile><disableXmlReport>true</disableXmlReport><!-- If you have classpath issue like NoDefClassError,... --><!-- useManifestOnlyJar>false</useManifestOnlyJar --><includes><include>***Suite.*</include></includes></configuration></plugin>                 </plugins>    </build></project>

感谢各位的阅读,以上就是“flink连接消费kafka实例”的内容了,经过本文的学习后,相信大家对flink连接消费kafka实例这一问题有了更深刻的体会,具体使用情况还需要大家实践验证。这里是编程网,小编将为大家推送更多相关知识点的文章,欢迎关注!

阅读原文内容投诉

免责声明:

① 本站未注明“稿件来源”的信息均来自网络整理。其文字、图片和音视频稿件的所属权归原作者所有。本站收集整理出于非商业性的教育和科研之目的,并不意味着本站赞同其观点或证实其内容的真实性。仅作为临时的测试数据,供内部测试之用。本站并未授权任何人以任何方式主动获取本站任何信息。

② 本站未注明“稿件来源”的临时测试数据将在测试完成后最终做删除处理。有问题或投稿请发送至: 邮箱/279061341@qq.com QQ/279061341

软考中级精品资料免费领

  • 历年真题答案解析
  • 备考技巧名师总结
  • 高频考点精准押题
  • 2024年上半年信息系统项目管理师第二批次真题及答案解析(完整版)

    难度     807人已做
    查看
  • 【考后总结】2024年5月26日信息系统项目管理师第2批次考情分析

    难度     351人已做
    查看
  • 【考后总结】2024年5月25日信息系统项目管理师第1批次考情分析

    难度     314人已做
    查看
  • 2024年上半年软考高项第一、二批次真题考点汇总(完整版)

    难度     433人已做
    查看
  • 2024年上半年系统架构设计师考试综合知识真题

    难度     221人已做
    查看

相关文章

发现更多好内容

猜你喜欢

AI推送时光机
位置:首页-资讯-后端开发
咦!没有更多了?去看看其它编程学习网 内容吧
首页课程
资料下载
问答资讯