码农知识堂 - 1000bd
  •   Python
  •   PHP
  •   JS/TS
  •   JAVA
  •   C/C++
  •   C#
  •   GO
  •   Kotlin
  •   Swift
  • 【Flink实战】Flink自定义的Source 数据源案例-并行度调整结合WebUI


    🚀 作者 :“大数据小禅”

    🚀 文章简介 :【Flink实战】玩转Flink里面核心的Source Operator实战

    🚀 欢迎小伙伴们 点赞👍、收藏⭐、留言💬


    目录导航

        • 什么是Flink的并行度
        • Flink自定义的Source 数据源案例-并行度调整结合WebUI

    什么是Flink的并行度

    • Flink的并行度是指在Flink应用程序中并行执行任务的级别或程度。它决定了任务在Flink集群中的并发执行程度,即任务被划分成多少个并行的子任务。

    • 在Flink中,可以通过设置并行度来控制任务的并行执行。并行度是根据数据或计算的特性来确定的,可以根据任务的特点和所需的处理能力进行调优。

    • 将一个任务的并行度设置为N意味着将该任务分成N个并行的子任务,这些子任务可以在Flink集群的不同节点上同时执行。Flink会根据配置的并行度自动对任务进行数据切分和任务调度,以实现高效的并行处理。

    • 选择合适的并行度需要在平衡性、吞吐量和可伸缩性之间权衡。较高的并行度可以提高任务的处理能力和吞吐量,但也会增加系统的资源需求和管理成本。较低的并行度可能导致资源浪费和性能瓶颈。

    • 在设计Flink应用程序时,可以根据任务之间的依赖关系、数据流量、数据分布以及可用的资源来选择合适的并行度。可以通过调整并行度来优化任务的性能,平衡任务的负载,提高整体的处理能力。-

    Flink自定义的Source 数据源案例-并行度调整结合WebUI

    • 开启webui
      取消掉默认并行度为1,因为默认的并行度是8,也就是8个线程 默认的并行度就是系统的核数在这里插入图片描述
      在这里插入图片描述
    StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironmentWithWebUI(new Configuration());
    
    • 1
    • 设置不同的并行度
      Solt的数量就是设置的最大并行度的数量
      在这里插入图片描述
      在这里插入图片描述
    public static void main(String[] args) throws Exception {
    
            //构建执行任务环境以及任务的启动的入口, 存储全局相关的参数
            //StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
            StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironmentWithWebUI(new Configuration());
            env.setRuntimeMode(RuntimeExecutionMode.AUTOMATIC);
            env.setParallelism(2);
    
            DataStream<VideoOrder> videoOrderDS =  env.addSource(new VideoOrderSource());
    
            DataStream<VideoOrder> filterDS = videoOrderDS.filter(new FilterFunction<VideoOrder>() {
                @Override
                public boolean filter(VideoOrder videoOrder) throws Exception {
                    return videoOrder.getMoney()>5;
                }
            }).setParallelism(3);
    
            filterDS.print().setParallelism(4);
    
            //DataStream需要调用execute,可以取个名称
            env.execute("source job");
        }
    
    
    
    • 1
    • 2
    • 3
    • 4
    • 5
    • 6
    • 7
    • 8
    • 9
    • 10
    • 11
    • 12
    • 13
    • 14
    • 15
    • 16
    • 17
    • 18
    • 19
    • 20
    • 21
    • 22
    • 23
    • 24

    数据流中最大的并行度,就是算子链中最大算子的数量,比如source 2个并行度,filter 4个,sink 4个,最大就是4
    在这里插入图片描述
    在这里插入图片描述

  • 相关阅读:
    基础复习——图形定制——图形Drawable——形状图形——九宫格图片——状态列表图形...
    xxx.ko 驱动模块加载报错 “unknown symbol in module or invalid parameter”
    一加手机线刷2024版,param预载失败/MSM刷机工具报错
    计算机学院改考后,网络空间安全学院也改考了!南京理工大学计算机考研
    万字长文保姆级教你制作自己的多功能QQ机器人
    问题引入:多个线程读写同一共享变量是否存在并发问题?
    390. 消除游戏
    好心情:长期服用抗抑郁药,怎么把肝损伤降到最低?
    数据分析 | Pandas 200道练习题,每日10道题,学完必成大神(4)
    在Kubernetes中实现gRPC流量负载均衡
  • 原文地址:https://blog.csdn.net/weixin_45574790/article/details/132857034
  • 最新文章
  • 沪漂五周年了:我越来越迷茫了
    Agentic Skill Routing 实战:别再把所有 Skill 塞进 AI Agent 上下文
    MySQL-Seconds_behind_master的精度误差
    [MAF预定义ChatClient中间件-03]CachingChatClient——利用缓存省钱省时间
    AI的至暗历史:从万众期待到被政府撤资,AI的两次死亡徘徊
    Agent OS :五种驯服不确定性的范式
    PortSwigger SQL注入LAB11
    数据库即时编译JIT
    [Begin]AI Learn Data Day 0
    深度学习进阶(二十七)现代 LLM 的核心架构设计其二:SwiGLU
  • 热门文章
  • 十款代码表白小特效 一个比一个浪漫 赶紧收藏起来吧!!!
    奉劝各位学弟学妹们,该打造你的技术影响力了!
    五年了,我在 CSDN 的两个一百万。
    Java俄罗斯方块,老程序员花了一个周末,连接中学年代!
    面试官都震惊,你这网络基础可以啊!
    你真的会用百度吗?我不信 — 那些不为人知的搜索引擎语法
    心情不好的时候,用 Python 画棵樱花树送给自己吧
    通宵一晚做出来的一款类似CS的第一人称射击游戏Demo!原来做游戏也不是很难,连憨憨学妹都学会了!
    13 万字 C 语言从入门到精通保姆级教程2021 年版
    10行代码集2000张美女图,Python爬虫120例,再上征途
小工具 小游戏
Copyright © 2022 侵权请联系2656653265@qq.com    京ICP备2022015340号-1

京公网安备 11010502049817号