暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

Flink任务编写流程总结

IT那活儿 2025-10-13
79

点击上方“IT那活儿”公众号--专注于企业全栈运维技术分享,不管IT什么活儿,干就完了!!!



概述

本文档讲述如何使用flink的java接口来编写一个流式处理任务。
环境准备:
  • IntelliJIEDA2018+
  • Java1.8+


编写过程

2.1 创建maven工程
File,选择new project,选择Maven,ProjectSDK选择java1.8,然后点击next
设置合适的GroupId和ArtifactId,点击确定后工程自动生成。
2.2 配置工程的pom.xml
加入java接口和scala接口的依赖:
加入kafka的客户端依赖:
2.3 编写flink消费kafka流式处理代码
右键点击src,main,java文件夹,New一个javaclass,取名FlinkTestTask。
添加如下代码:
  • 第12行是flink获取流式处理的执行环境;
  • 第15到17行设置kafka的属性,必须包含bootstrap.server和group.id两个必须属性,其他属性可以参考kafka官网客户端的默认属性,必要时候可以在程序中修改。
  • 第18行通过flink的kafka消费connector新建一个消费kafka的test_topic的对象;
  • 第20行设置kafka通过最新的消息开始消费;
  • 第23行将新建的kafka消费添加的flink程序的执行环境的数据源source中,这样就可以开始消费kafka中的数据了。
2.4 处理kafka中的数据
假设kafka的数据是默认的json格式,我们选择合适的json转换工具包,把数据转换为map类,然后进行处理。
在pom.xml中添加alibaba的fastjson依赖:
在代码中把kafka获取到的String数据转换为map格式:
用到了flink的flatmap算子,把string格式的流转换为map格式的流。
2.5 提取需要的信息存入目标kafka
提取map中的price字段,把价格存入目标kafka中:
2.6 开始运行
注意这段代码一定不能少,否则任务运行报错。
2.7 打包
选用自带插件打包:
如图,选中install,点击右上角播放按钮即可打包完成,打包完成后,在目录下target目录获取jar包。
2.8 发布到yarn环境运行
flinkrun-myarn-cluster-p 8 -yjm 1024 -ytm 4096 -ynmflink_test-cFlinkTest  Taskxxx.jar
  • 8表示并行度。
  • 1024是jobmanager占用的内存大小,4096是taskmanager占用的内存大小。

END


本文作者:郭 宪(上海新炬中北团队)

本文来源:“IT那活儿”公众号

文章转载自IT那活儿,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论