【Flink】双流处理:实时对账实现

2024-08-29 10:32

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

Flink双流处理:实时对账实现

  • 一、基础概念
  • 二、双流处理的方法
    • Connect
    • Union
    • Join
  • 三、实战:实时对账实现
    • 需求描述
    • 需求分析
    • 代码实现
  • 相关阅读

更多内容详见:https://github.com/pierre94/flink-notes

一、基础概念

主要是两种处理模式:

  • Connect/Join
  • Union

二、双流处理的方法

Connect

DataStream,DataStream → ConnectedStreams

连接两个保持他们类型的数据流,两个数据流被Connect之后,只是被放在了一个同一个流中,内部依然保持各自的数据和形式不发生任何变化,两个流相互独立。

Connect后使用CoProcessFunction、CoMap、CoFlatMap、KeyedCoProcessFunction等API 对两个流分别处理。如CoMap:

val warning = high.map( sensorData => (sensorData.id, sensorData.temperature) )
val connected = warning.connect(low)val coMap = connected.map(
warningData => (warningData._1, warningData._2, "warning"),
lowData => (lowData.id, "healthy")
)

(ConnectedStreams → DataStream 功能与 map 一样,对 ConnectedStreams 中的每一个流分别进行 map 和 flatMap 处理。)

疑问,既然两个流内部独立,那Connect 后有什么意义呢?

Connect后的两条流可以共享状态,在对账等场景具有重大意义!

Union


DataStream → DataStream:对两个或者两个以上的 DataStream 进行 union 操作,产生一个包含所有 DataStream 元素的新 DataStream。

val unionStream: DataStream[StartUpLog] = appStoreStream.union(otherStream) unionStream.print("union:::")

注意:Union 可以操作多个流,而Connect只能对两个流操作

Join

Join是基于Connect更高层的一个实现,结合Window实现。

相关知识点比较多,详细文档见: https://ci.apache.org/projects/flink/flink-docs-release-1.10/dev/stream/operators/joining.html

三、实战:实时对账实现

需求描述

有两个时间Event1、Event2,第一个字段是时间id,第二个字段是时间戳,需要对两者进行实时对账。当其中一个事件缺失、延迟时要告警出来。

需求分析

类似之前的订单超时告警需求。之前数据源是一个流,我们在function里面进行一些改写。这里我们分别使用Event1和Event2两个流进行Connect处理。

// 事件1
case class Event1(id: Long, eventTime: Long)
// 事件2
case class Event2(id: Long, eventTime: Long)
// 输出结果
case class Result(id: Long, warnings: String)

代码实现

scala实现

涉及知识点:

  • 双流Connect
  • 使用OutputTag侧输出
  • KeyedCoProcessFunction(processElement1、processElement2)使用
  • ValueState使用
  • 定时器onTimer使用

启动两个TCP服务:

nc -lh 9999
nc -lk 9998

注意:nc启动的是服务端、flink启动的是客户端

import java.text.SimpleDateFormatimport org.apache.flink.api.common.state.{ValueState, ValueStateDescriptor}
import org.apache.flink.streaming.api.TimeCharacteristic
import org.apache.flink.streaming.api.functions.co.KeyedCoProcessFunction
import org.apache.flink.streaming.api.scala.{StreamExecutionEnvironment, _}
import org.apache.flink.util.Collectorobject CoTest {val simpleDateFormat = new SimpleDateFormat("dd/MM/yyyy:HH:mm:ss")val txErrorOutputTag = new OutputTag[Result]("txErrorOutputTag")def main(args: Array[String]): Unit = {val env = StreamExecutionEnvironment.getExecutionEnvironmentenv.setStreamTimeCharacteristic(TimeCharacteristic.EventTime)env.setParallelism(1)val event1Stream = env.socketTextStream("127.0.0.1", 9999).map(data => {val dataArray = data.split(",")Event1(dataArray(0).trim.toLong, simpleDateFormat.parse(dataArray(1).trim).getTime)}).assignAscendingTimestamps(_.eventTime * 1000L).keyBy(_.id)val event2Stream = env.socketTextStream("127.0.0.1", 9998).map(data => {val dataArray = data.split(",")Event2(dataArray(0).trim.toLong, simpleDateFormat.parse(dataArray(1).trim).getTime)}).assignAscendingTimestamps(_.eventTime * 1000L).keyBy(_.id)val coStream = event1Stream.connect(event2Stream).process(new CoTestProcess())//    union 必须是同一条类型的流//    val unionStream = event1Stream.union(event2Stream)//    unionStream.print()coStream.print("ok")coStream.getSideOutput(txErrorOutputTag).print("txError")env.execute("union test")}//共享状态class CoTestProcess() extends KeyedCoProcessFunction[Long,Event1, Event2, Result] {lazy val event1State: ValueState[Boolean]= getRuntimeContext.getState(new ValueStateDescriptor[Boolean]("event1-state", classOf[Boolean]))lazy val event2State: ValueState[Boolean]= getRuntimeContext.getState(new ValueStateDescriptor[Boolean]("event2-state", classOf[Boolean]))override def processElement1(value: Event1, ctx: KeyedCoProcessFunction[Long, Event1, Event2, Result]#Context, out: Collector[Result]): Unit = {if (event2State.value()) {event2State.clear()out.collect(Result(value.id, "ok"))} else {event1State.update(true)//等待一分钟ctx.timerService().registerEventTimeTimer(value.eventTime + 1000L * 60)}}override def processElement2(value: Event2, ctx: KeyedCoProcessFunction[Long, Event1, Event2, Result]#Context, out: Collector[Result]): Unit = {if (event1State.value()) {event1State.clear()out.collect(Result(value.id, "ok"))} else {event2State.update(true)ctx.timerService().registerEventTimeTimer(value.eventTime + 1000L * 60)}}override def onTimer(timestamp: Long, ctx: KeyedCoProcessFunction[Long, Event1, Event2, Result]#OnTimerContext, out: Collector[Result]): Unit = {if(event1State.value()){ctx.output(txErrorOutputTag,Result(ctx.getCurrentKey,s"no event2,timestamp:$timestamp"))event1State.clear()}else if(event2State.value()){ctx.output(txErrorOutputTag,Result(ctx.getCurrentKey,s"no event1,timestamp:$timestamp"))event2State.clear()}}}}

相关阅读

《Flink状态编程: 订单超时告警》:
https://blog.csdn.net/u013128262/article/details/104648592

《github:Flink学习笔记》:
https://github.com/pierre94/flink-notes


原创声明,本文系作者授权云+社区发表,未经许可,不得转载。

https://cloud.tencent.com/developer/article/1596145

这篇关于【Flink】双流处理:实时对账实现的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



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

相关文章

Vue中动态权限到按钮的完整实现方案详解

《Vue中动态权限到按钮的完整实现方案详解》这篇文章主要为大家详细介绍了Vue如何在现有方案的基础上加入对路由的增、删、改、查权限控制,感兴趣的小伙伴可以跟随小编一起学习一下... 目录一、数据库设计扩展1.1 修改路由表(routes)1.2 修改角色与路由权限表(role_routes)二、后端接口设计

C#集成DeepSeek模型实现AI私有化的流程步骤(本地部署与API调用教程)

《C#集成DeepSeek模型实现AI私有化的流程步骤(本地部署与API调用教程)》本文主要介绍了C#集成DeepSeek模型实现AI私有化的方法,包括搭建基础环境,如安装Ollama和下载DeepS... 目录前言搭建基础环境1、安装 Ollama2、下载 DeepSeek R1 模型客户端 ChatBo

Qt实现发送HTTP请求的示例详解

《Qt实现发送HTTP请求的示例详解》这篇文章主要为大家详细介绍了如何通过Qt实现发送HTTP请求,文中的示例代码讲解详细,具有一定的借鉴价值,感兴趣的小伙伴可以跟随小编一起学习一下... 目录1、添加network模块2、包含改头文件3、创建网络访问管理器4、创建接口5、创建网络请求对象6、创建一个回复对

C++实现回文串判断的两种高效方法

《C++实现回文串判断的两种高效方法》文章介绍了两种判断回文串的方法:解法一通过创建新字符串来处理,解法二在原字符串上直接筛选判断,两种方法都使用了双指针法,文中通过代码示例讲解的非常详细,需要的朋友... 目录一、问题描述示例二、解法一:将字母数字连接到新的 string思路代码实现代码解释复杂度分析三、

grom设置全局日志实现执行并打印sql语句

《grom设置全局日志实现执行并打印sql语句》本文主要介绍了grom设置全局日志实现执行并打印sql语句,包括设置日志级别、实现自定义Logger接口以及如何使用GORM的默认logger,通过这些... 目录gorm中的自定义日志gorm中日志的其他操作日志级别Debug自定义 Loggergorm中的

Spring Boot整合消息队列RabbitMQ的实现示例

《SpringBoot整合消息队列RabbitMQ的实现示例》本文主要介绍了SpringBoot整合消息队列RabbitMQ的实现示例,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的... 目录RabbitMQ 简介与安装1. RabbitMQ 简介2. RabbitMQ 安装Spring

Gin框架中的GET和POST表单处理的实现

《Gin框架中的GET和POST表单处理的实现》Gin框架提供了简单而强大的机制来处理GET和POST表单提交的数据,通过c.Query、c.PostForm、c.Bind和c.Request.For... 目录一、GET表单处理二、POST表单处理1. 使用c.PostForm获取表单字段:2. 绑定到结

mysql8.0无备份通过idb文件恢复数据的方法、idb文件修复和tablespace id不一致处理

《mysql8.0无备份通过idb文件恢复数据的方法、idb文件修复和tablespaceid不一致处理》文章描述了公司服务器断电后数据库故障的过程,作者通过查看错误日志、重新初始化数据目录、恢复备... 周末突然接到一位一年多没联系的妹妹打来电话,“刘哥,快来救救我”,我脑海瞬间冒出妙瓦底,电信火苲马扁.

springMVC返回Http响应的实现

《springMVC返回Http响应的实现》本文主要介绍了在SpringBoot中使用@Controller、@ResponseBody和@RestController注解进行HTTP响应返回的方法,... 目录一、返回页面二、@Controller和@ResponseBody与RestController

nginx中重定向的实现

《nginx中重定向的实现》本文主要介绍了Nginx中location匹配和rewrite重定向的规则与应用,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有一定的参考学习价值,需要的朋友们下... 目录一、location1、 location匹配2、 location匹配的分类2.1 精确匹配2