Flink 中定时加载外部数据

社区中有好几个同学问过这样的场景:  

flink 任务中,source 进来的数据,需要连接数据库里面的字段,再做后面的处理

这里假设一个 ETL 的场景,输入数据包含两个字段 “type, userid....” ,需要根据 type,连接一张 mysql 的配置表,关联 type 对应的具体内容。相对于输入数据的数量,type 的值是很少的(这里默认只有10种), 所以对应配置表就只有10条数据,配置是会定时修改的(比如跑批补充数据),配置的修改必须在一定时间内生效。

实时 ETL,需要用里面的一个字段去关联数据库,补充其他数据,进来的数据中关联字段是很单一的(就10个),对应数据库的数据也很少,如果用 异步 IO,感觉会比较傻(浪费资源、性能还不好)。同时数据库的数据是会不定时修改的,所以不能在启动的时候一次性加载。

Flink 现在对应这种场景可以使用  Boradcase state 做,如:基于Broadcast 状态的Flink Etl Demo

这里想说的是另一种更简单的方法: 使用定时器,定时加载数据库的数据 (就是简单的Java定时器)

先说一下代码流程:

1、自定义的 source,输入逗号分隔的两个字段

2、使用 RichMapFunction  转换数据,在 open 中定义定时器,定时触发查询 mysql 的任务,并将结果放到一个 map 中

3、输入数据关联 map 的数据,然后输出

先看下数据库中的数据:

1mysql> select * from timer; 2+------+------+ 3| id | name | 4+------+------+ 5| 0 | 0zOq | 6| 1 | 1hKC | 7| 2 | 2ibM | 8| 3 | 3fCe | 9| 4 | 4TaM | 10| 5 | 5URU | 11| 6 | 6WhP | 12| 7 | 7zjn | 13| 8 | 8Szl | 14| 9 | 9blS | 15+------+------+ 1610 rows in set (0.01 sec)

总共10条数据,id 就是对应的关联字段,需要填充的数据是 name

下面是主要的代码:// 自定义的source,输出 x,xxx 格式随机字符

1val input = env.addSource(new TwoStringSource) 2 val stream = input.map(new RichMapFunction[String, String] { 3 4 val jdbcUrl = "jdbc:mysql://venn:3306?useSSL=false&allowPublicKeyRetrieval=true" 5 val username = "root" 6 val password = "123456" 7 val driverName = "com.mysql.jdbc.Driver" 8 var conn: Connection = null 9 var ps: PreparedStatement = null 10 val map = new util.HashMap[String, String]() 11 12 override def open(parameters: Configuration): Unit = { 13 logger.info("init....") 14 query() 15 // new Timer 16 val timer = new Timer(true) 17 // schedule is 10 second 定义了一个10秒的定时器,定时执行查询数据库的方法 18 timer.schedule(new TimerTask { 19 override def run(): Unit = { 20 query() 21 } 22 }, 10000, 10000) 23 24 } 25 26 override def map(value: String): String = { 27 // concat input and mysql data,简单关联输出 28 value + "-" + map.get(value.split(",")(0)) 29 } 30 31 /** 32 * query mysql for get new config data 33 */ 34 def query() = { 35 logger.info("query mysql") 36 try { 37 Class.forName(driverName) 38 conn = DriverManager.getConnection(jdbcUrl, username, password) 39 ps = conn.prepareStatement("select id,name from venn.timer") 40 val rs = ps.executeQuery 41 42 while (!rs.isClosed && rs.next) { 43 val id = rs.getString(1) 44 val name = rs.getString(2) // 将结果放到 map 中 45 map.put(id, name) 46 } 47 logger.info("get config from db size : {}", map.size()) 48 49 } catch { 50 case e@(_: ClassNotFoundException | _: SQLException) => 51 e.printStackTrace() 52 } finally { 53 ps.close() 54 conn.close() 55 } 56 } 57 }) 58// .print() 59 60 61 val sink = new FlinkKafkaProducer[String]("timer_out" 62 , new MyKafkaSerializationSchema[String]() 63 , Common.getProp 64 , FlinkKafkaProducer.Semantic.EXACTLY_ONCE) 65 stream.addSink(sink)

简单的Java定时器:

1val timer = new Timer(true) 2// schedule is 10 second, 5 second between successive task executions 3timer.schedule(new TimerTask { 4 override def run(): Unit = { 5 query() 6 } 7}, 10000, 10000)

------------------20200327 改---------------------

之前 博客写的有问题,public void schedule(TimerTask task, long delay, long period) 的第三个参数才是重复执行的时间间隔,0 是不执行,我之前写的时候放上去的案例,调用的 Timer 的构造方法是: public void schedule(TimerTask task, long delay) 只会在 delay 时间后调用一次,并不会重复执行,不需要 调用 : public void schedule(TimerTask task, long delay, long period)   这样的构造方法,才能真正的定时执行。

使用之前的方法执行的,会看到query 方法执行了两次,是 open 中主动调用了一次和 之后调度了一次,定时器就结束了。

感谢社区大佬指出

同时社区还有大佬指出 : ScheduledExecutorService 会比 timer 更好;理由: Timer里边的逻辑失败的话不会抛出任何异常,直接结束,建议用ScheduledExecutorService替换Timer并且捕获下异常看看

------------------------------------

看下输出的数据:

17,N-7zjn 27,C-7zjn 37,U-7zjn 44,T-4TaM 57,J-7zjn 69,R-9blS 74,C-4TaM 89,T-9blS 94,A-4TaM 106,I-6WhP 119,U-9blS

注:“-” 之前是原始数据,后面是关联后的数据

部署到服务器上定时器的调度:

12019-09-28 18:28:13,476 INFO com.venn.stream.api.timer.CustomerTimerDemo$ - query mysql 22019-09-28 18:28:13,480 INFO com.venn.stream.api.timer.CustomerTimerDemo$ - get config from db size : 10 32019-09-28 18:28:18,553 INFO org.apache.flink.streaming.api.functions.sink.TwoPhaseCommitSinkFunction - FlinkKafkaProducer 0/1 - checkpoint 17 complete, committing transaction TransactionHolder{handle=KafkaTransactionState [transactionalId=null, producerId=-1, epoch=-1], transactionStartTime=1569666488499} from checkpoint 17 42019-09-28 18:28:23,476 INFO com.venn.stream.api.timer.CustomerTimerDemo$ - query mysql 52019-09-28 18:28:23,481 INFO com.venn.stream.api.timer.CustomerTimerDemo$ - get config from db size : 10 62019-09-28 18:28:28,549 INFO org.apache.flink.streaming.api.functions.sink.TwoPhaseCommitSinkFunction - FlinkKafkaProducer 0/1 - checkpoint 18 complete, committing transaction TransactionHolder{handle=KafkaTransactionState [transactionalId=null, producerId=-1, epoch=-1], transactionStartTime=1569666498505} from checkpoint 18 72019-09-28 18:28:33,477 INFO com.venn.stream.api.timer.CustomerTimerDemo$ - query mysql 82019-09-28 18:28:33,484 INFO com.venn.stream.api.timer.CustomerTimerDemo$ - get config from db size : 10

十秒调度一次

 欢迎关注Flink菜鸟公众号,会不定期更新Flink(开发技术)相关的推文

 

点赞
收藏

评论区

加载中...

相关推荐

MySQL:[Err] 1292 - Incorrect datetime value: ‘0000-00-00 00:00:00‘ for column ‘CREATE_TIME‘ at row 1

文章目录问题用navicat导入数据时,报错:原因这是因为当前的MySQL不支持datetime为0的情况。解决修改sql\mode:sql\mode:SQLMode定义了MySQL应支持的SQL语法、数据校验等,这样可以更容易地在不同的环境中使用MySQL。全局s

Oracle 分组与拼接字符串同时使用

SELECTT.,ROWNUMIDFROM(SELECTT.EMPLID,T.NAME,T.BU,T.REALDEPART,T.FORMATDATE,SUM(T.S0)S0,MAX(UPDATETIME)CREATETIME,LISTAGG(TOCHAR(

皕杰报表之UUID

​在我们用皕杰报表工具设计填报报表时,如何在新增行里自动增加id呢?能新增整数排序id吗?目前可以在新增行里自动增加id,但只能用uuid函数增加UUID编码,不能新增整数排序id。uuid函数说明:获取一个UUID,可以在填报表中用来创建数据ID语法:uuid()或uuid(sep)参数说明:sep布尔值,生成的uuid中是否包含分隔符'',缺省为

手写Java HashMap源码

HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程HashMap的使用教程22

2020年前端实用代码段,为你的工作保驾护航

有空的时候,自己总结了几个代码段,在开发中也经常使用,谢谢。1、使用解构获取json数据let jsonData  id: 1,status: "OK",data: 'a', 'b';let  id, status, data: number   jsonData;console.log(id, status, number )

Twitter的分布式自增ID算法snowflake (Java版)

概述分布式系统中,有一些需要使用全局唯一ID的场景,这种时候为了防止ID冲突可以使用36位的UUID,但是UUID有一些缺点,首先他相对比较长,另外UUID一般是无序的。有些时候我们希望能使用一种简单一些的ID,并且希望ID能够按照时间有序生成。而twitter的snowflake解决了这种需求,最初Twitter把存储系统从MySQL迁移

Flink 中定时加载外部数据 - HelloWorld