点击上方“IT那活儿”公众号--专注于企业全栈运维技术分享,不管IT什么活儿,干就完了!!!
概述
IntelliJIEDA2018+ Java1.8+
编写过程



第12行是flink获取流式处理的执行环境; 第15到17行设置kafka的属性,必须包含bootstrap.server和group.id两个必须属性,其他属性可以参考kafka官网客户端的默认属性,必要时候可以在程序中修改。 第18行通过flink的kafka消费connector新建一个消费kafka的test_topic的对象; 第20行设置kafka通过最新的消息开始消费; 第23行将新建的kafka消费添加的flink程序的执行环境的数据源source中,这样就可以开始消费kafka中的数据了。





flinkrun-myarn-cluster-p 8 -yjm 1024 -ytm 4096 -ynmflink_test-cFlinkTest Taskxxx.jar
8表示并行度。 1024是jobmanager占用的内存大小,4096是taskmanager占用的内存大小。

本文作者:郭 宪(上海新炬中北团队)
本文来源:“IT那活儿”公众号

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




