FlinkX流控实现

2024-08-28 05:32
文章标签 实现 流控 flinkx

本文主要是介绍FlinkX流控实现,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

FlinkX流控实现

流量控制防止并发性能过高对源数据库造成影响。

在FlinkX中,流量控制是采用Byte为单位/s进行流量控制的。

配置参数实例:

{“speed”: {"bytes": 0}
}

当 bytes > 0 时,才会开启流量控制。

整个计算的速率是根据整个系统中的指标,按照每秒的窗口,实时计算出限流的速率的。通过对读取记录的限流,但是没有背压。


image-20200524100353781

限流原理

实现逻辑

首先我们看一下读取器的初始化过程,

InputFormat初始化

public void open(InputSplit inputSplit) throws IOException {checkIfCreateSplitFailed(inputSplit);if(!inited){// 初始化累加器收集器,每秒从FlinkAPI读取对应累加器的值,并设置为global值initAccumulatorCollector();// 初始化指标累加器,每次调用nextRecord时提交指标initStatisticsAccumulator();// 开启限流器openByteRateLimiter();initRestoreInfo();if(restoreConfig.isRestore()){formatState.setNumOfSubTask(indexOfSubTask);}inited = true;}openInternal(inputSplit);}

本次只看前三个。

  1. 初始化累加器收集器,每秒从FlinkAPI读取对应累加器的值,并设置为global值(下文中计算速率是有用到)

  2. 初始化指标累加器,每次调用nextRecord时提交指标

  3. 开启限流器

  4. 初始化Restore配置(本章不讲,后续章节有用到)

我们重点详解一下前三个步骤:

在详解每一个步骤之前,首先了解下在数据同步过程中具体的指标

指标详情

分类指标名称含义
读取指标numRead累计读取数据条数
byteRead累计读取数据字节数
readDuration读取数据的总时间
写入指标numWrite累计写入数据条数
byteWrite累计写入数据字节数
writeDuration写入数据的总时间
错误指标nErrors累计错误记录数
nullErrors累计空指针错误记录数
duplicateErrors累计主键冲突错误记录数
conversionErrors累计类型转换错误记录数
otherErrors累计其它错误记录数

全局指标实现

如何控制全局限流,很重要的一环就是收集到全局系统的关键状况,无论是微服务调用还是读取限流本质都是同一个道理。首先需要找一个全局存储提供这些指标的存储和更新,FlinkX在这里使用的Flinx的累加器

image-20200524101615078

指标初始化
private void initStatisticsAccumulator(){numReadCounter = getRuntimeContext

这篇关于FlinkX流控实现的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/1113862

相关文章

Spring Boot 实现 IP 限流的原理、实践与利弊解析

《SpringBoot实现IP限流的原理、实践与利弊解析》在SpringBoot中实现IP限流是一种简单而有效的方式来保障系统的稳定性和可用性,本文给大家介绍SpringBoot实现IP限... 目录一、引言二、IP 限流原理2.1 令牌桶算法2.2 漏桶算法三、使用场景3.1 防止恶意攻击3.2 控制资源

springboot下载接口限速功能实现

《springboot下载接口限速功能实现》通过Redis统计并发数动态调整每个用户带宽,核心逻辑为每秒读取并发送限定数据量,防止单用户占用过多资源,确保整体下载均衡且高效,本文给大家介绍spring... 目录 一、整体目标 二、涉及的主要类/方法✅ 三、核心流程图解(简化) 四、关键代码详解1️⃣ 设置

Nginx 配置跨域的实现及常见问题解决

《Nginx配置跨域的实现及常见问题解决》本文主要介绍了Nginx配置跨域的实现及常见问题解决,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下面随着小编来... 目录1. 跨域1.1 同源策略1.2 跨域资源共享(CORS)2. Nginx 配置跨域的场景2.1

Python中提取文件名扩展名的多种方法实现

《Python中提取文件名扩展名的多种方法实现》在Python编程中,经常会遇到需要从文件名中提取扩展名的场景,Python提供了多种方法来实现这一功能,不同方法适用于不同的场景和需求,包括os.pa... 目录技术背景实现步骤方法一:使用os.path.splitext方法二:使用pathlib模块方法三

CSS实现元素撑满剩余空间的五种方法

《CSS实现元素撑满剩余空间的五种方法》在日常开发中,我们经常需要让某个元素占据容器的剩余空间,本文将介绍5种不同的方法来实现这个需求,并分析各种方法的优缺点,感兴趣的朋友一起看看吧... css实现元素撑满剩余空间的5种方法 在日常开发中,我们经常需要让某个元素占据容器的剩余空间。这是一个常见的布局需求

HTML5 getUserMedia API网页录音实现指南示例小结

《HTML5getUserMediaAPI网页录音实现指南示例小结》本教程将指导你如何利用这一API,结合WebAudioAPI,实现网页录音功能,从获取音频流到处理和保存录音,整个过程将逐步... 目录1. html5 getUserMedia API简介1.1 API概念与历史1.2 功能与优势1.3

Java实现删除文件中的指定内容

《Java实现删除文件中的指定内容》在日常开发中,经常需要对文本文件进行批量处理,其中,删除文件中指定内容是最常见的需求之一,下面我们就来看看如何使用java实现删除文件中的指定内容吧... 目录1. 项目背景详细介绍2. 项目需求详细介绍2.1 功能需求2.2 非功能需求3. 相关技术详细介绍3.1 Ja

使用Python和OpenCV库实现实时颜色识别系统

《使用Python和OpenCV库实现实时颜色识别系统》:本文主要介绍使用Python和OpenCV库实现的实时颜色识别系统,这个系统能够通过摄像头捕捉视频流,并在视频中指定区域内识别主要颜色(红... 目录一、引言二、系统概述三、代码解析1. 导入库2. 颜色识别函数3. 主程序循环四、HSV色彩空间详解

PostgreSQL中MVCC 机制的实现

《PostgreSQL中MVCC机制的实现》本文主要介绍了PostgreSQL中MVCC机制的实现,通过多版本数据存储、快照隔离和事务ID管理实现高并发读写,具有一定的参考价值,感兴趣的可以了解一下... 目录一 MVCC 基本原理python1.1 MVCC 核心概念1.2 与传统锁机制对比二 Postg

SpringBoot整合Flowable实现工作流的详细流程

《SpringBoot整合Flowable实现工作流的详细流程》Flowable是一个使用Java编写的轻量级业务流程引擎,Flowable流程引擎可用于部署BPMN2.0流程定义,创建这些流程定义的... 目录1、流程引擎介绍2、创建项目3、画流程图4、开发接口4.1 Java 类梳理4.2 查看流程图4