Flink 异步IO访问外部数据(mysql篇)

  接上篇:【翻译】Flink 异步I / O访问外部数据

  最近看了大佬的博客,突然想起Async I/O方式是Blink 推给社区的一大重要功能,可以使用异步的方式获取外部数据,想着自己实现以下,项目上用的时候,可以不用现去找了。

  最开始想用scala 实现一个读取 hbase数据的demo,参照官网demo:  

1/** 2 * An implementation of the 'AsyncFunction' that sends requests and sets the callback. 3 */ 4class AsyncDatabaseRequest extends AsyncFunction[String, (String, String)] { 5 6 /** The database specific client that can issue concurrent requests with callbacks */ 7 lazy val client: DatabaseClient = new DatabaseClient(host, post, credentials) 8 9 /** The context used for the future callbacks */ 10 implicit lazy val executor: ExecutionContext = ExecutionContext.fromExecutor(Executors.directExecutor()) 11 12 13 override def asyncInvoke(str: String, resultFuture: ResultFuture[(String, String)]): Unit = { 14 15 // issue the asynchronous request, receive a future for the result 16 val resultFutureRequested: Future[String] = client.query(str) 17 18 // set the callback to be executed once the request by the client is complete 19 // the callback simply forwards the result to the result future 20 resultFutureRequested.onSuccess { 21 case result: String => resultFuture.complete(Iterable((str, result))) 22 } 23 } 24} 25 26// create the original stream 27val stream: DataStream[String] = ... 28 29// apply the async I/O transformation 30val resultStream: DataStream[(String, String)] = 31 AsyncDataStream.unorderedWait(stream, new AsyncDatabaseRequest(), 1000, TimeUnit.MILLISECONDS, 100)

失败了,上图标红的部分实现不了

1、Future 找不到可以用的实现类

2、unorderedWait 一直报错

源码example 里面也有Scala 的案例

1def main(args: Array[String]) { 2 val timeout = 10000L 3 4 val env = StreamExecutionEnvironment.getExecutionEnvironment 5 6 val input = env.addSource(new SimpleSource()) 7 8 val asyncMapped = AsyncDataStream.orderedWait(input, timeout, TimeUnit.MILLISECONDS, 10) { 9 (input, collector: ResultFuture[Int]) => 10 Future { 11 collector.complete(Seq(input)) 12 } (ExecutionContext.global) 13 } 14 15 asyncMapped.print() 16 17 env.execute("Async I/O job") 18 }

主要部分是这样的,菜鸡表示无力,想继承RichAsyncFunction,可以使用open 方法初始化链接。

网上博客翻了不少,大部分是翻译官网的原理,案例也没有可以执行的,苦恼。

失败了。

转为java版本的,昨天在群里问,有个大佬给我个Java版本的: https://github.com/perkinls/flink-local-train/blob/c8b4efe33620352aea0100adef4fae2a068a3b65/src/main/scala/com/lp/test/asyncio/AsyncIoSideTableJoinMysqlJava.java 还没看过,因为Java版的官网的案例能看懂。

下面开始上mysql 版本 的 源码(hbase 的还没测试过,本机的hbase 挂了):

业务如下:

接收kafka数据,转为user对象,调用async,使用user.id 查询对应的phone,放回user对象,输出

 主类:

1import com.alibaba.fastjson.JSON; 2import com.venn.common.Common; 3import org.apache.flink.formats.json.JsonNodeDeserializationSchema; 4import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.node.ObjectNode; 5import org.apache.flink.streaming.api.datastream.AsyncDataStream; 6import org.apache.flink.streaming.api.datastream.DataStream; 7import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; 8import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer; 9import java.util.concurrent.TimeUnit; 10 11 12public class AsyncMysqlRequest { 13 14 public static void main(String[] args) throws Exception { 15 16 final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); 17 FlinkKafkaConsumer<ObjectNode> source = new FlinkKafkaConsumer<>("async", new JsonNodeDeserializationSchema(), Common.getProp()); 18 19 // 接收kafka数据,转为User 对象 20 DataStream<User> input = env.addSource(source).map(value -> { 21 String id = value.get("id").asText(); 22 String username = value.get("username").asText(); 23 String password = value.get("password").asText(); 24 25 return new User(id, username, password); 26 }); 27 // 异步IO 获取mysql数据, timeout 时间 1s,容量 10(超过10个请求,会反压上游节点) 28 DataStream async = AsyncDataStream.unorderedWait(input, new AsyncFunctionForMysqlJava(), 1000, TimeUnit.MICROSECONDS, 10); 29 30 async.map(user -> { 31 32 return JSON.toJSON(user).toString(); 33 }) 34 .print(); 35 36 env.execute("asyncForMysql"); 37 38 } 39}

函数类:

1import org.apache.flink.configuration.Configuration; 2import org.apache.flink.streaming.api.functions.async.ResultFuture; 3import org.apache.flink.streaming.api.functions.async.RichAsyncFunction; 4import org.slf4j.Logger; 5import org.slf4j.LoggerFactory; 6import java.util.ArrayList; 7import java.util.Collections; 8import java.util.List; 9import java.util.concurrent.*; 10 11public class AsyncFunctionForMysqlJava extends RichAsyncFunction<AsyncUser, AsyncUser> { 12 13 14 Logger logger = LoggerFactory.getLogger(AsyncFunctionForMysqlJava.class); 15 private transient MysqlClient client; 16 private transient ExecutorService executorService; 17 18 /** 19 * open 方法中初始化链接 20 * 21 * @param parameters 22 * @throws Exception 23 */ 24 @Override 25 public void open(Configuration parameters) throws Exception { 26 logger.info("async function for mysql java open ..."); 27 super.open(parameters); 28 29 client = new MysqlClient(); 30 executorService = Executors.newFixedThreadPool(30); 31 } 32 33 /** 34 * use asyncUser.getId async get asyncUser phone 35 * 36 * @param asyncUser 37 * @param resultFuture 38 * @throws Exception 39 */ 40 @Override 41 public void asyncInvoke(AsyncUser asyncUser, ResultFuture<AsyncUser> resultFuture) throws Exception { 42 43 executorService.submit(() -> { 44 // submit query 45 System.out.println("submit query : " + asyncUser.getId() + "-1-" + System.currentTimeMillis()); 46 AsyncUser tmp = client.query1(asyncUser); 47 // 一定要记得放回 resultFuture,不然数据全部是timeout 的 48 resultFuture.complete(Collections.singletonList(tmp)); 49 }); 50 } 51 52 @Override 53 public void timeout(AsyncUser input, ResultFuture<AsyncUser> resultFuture) throws Exception { 54 logger.warn("Async function for hbase timeout"); 55 List<AsyncUser> list = new ArrayList(); 56 input.setPhone("timeout"); 57 list.add(input); 58 resultFuture.complete(list); 59 } 60 61 /** 62 * close function 63 * 64 * @throws Exception 65 */ 66 @Override 67 public void close() throws Exception { 68 logger.info("async function for mysql java close ..."); 69 super.close(); 70 } 71}

MysqlClient:

1import com.venn.flink.util.MathUtil; 2import org.apache.flink.shaded.netty4.io.netty.channel.DefaultEventLoop; 3import org.apache.flink.shaded.netty4.io.netty.util.concurrent.Future; 4import org.apache.flink.shaded.netty4.io.netty.util.concurrent.SucceededFuture; 5 6import java.sql.DriverManager; 7import java.sql.PreparedStatement; 8import java.sql.ResultSet; 9import java.sql.SQLException; 10 11public class MysqlClient { 12 13 private static String jdbcUrl = "jdbc:mysql://192.168.229.128:3306?useSSL=false&allowPublicKeyRetrieval=true"; 14 private static String username = "root"; 15 private static String password = "123456"; 16 private static String driverName = "com.mysql.jdbc.Driver"; 17 private static java.sql.Connection conn; 18 private static PreparedStatement ps; 19 20 static { 21 try { 22 Class.forName(driverName); 23 conn = DriverManager.getConnection(jdbcUrl, username, password); 24 ps = conn.prepareStatement("select phone from async.async_test where id = ?"); 25 } catch (ClassNotFoundException | SQLException e) { 26 e.printStackTrace(); 27 } 28 } 29 30 /** 31 * execute query 32 * 33 * @param user 34 * @return 35 */ 36 public AsyncUser query1(AsyncUser user) { 37 38 try { 39 Thread.sleep(10); 40 } catch (InterruptedException e) { 41 e.printStackTrace(); 42 } 43 44 String phone = "0000"; 45 try { 46 ps.setString(1, user.getId()); 47 ResultSet rs = ps.executeQuery(); 48 if (!rs.isClosed() && rs.next()) { 49 phone = rs.getString(1); 50 } 51 System.out.println("execute query : " + user.getId() + "-2-" + "phone : " + phone + "-" + System.currentTimeMillis()); 52 } catch (SQLException e) { 53 e.printStackTrace(); 54 } 55 user.setPhone(phone); 56 return user; 57 58 } 59 60 // 测试代码 61 public static void main(String[] args) { 62 MysqlClient mysqlClient = new MysqlClient(); 63 64 AsyncUser asyncUser = new AsyncUser(); 65 asyncUser.setId("526"); 66 long start = System.currentTimeMillis(); 67 asyncUser = mysqlClient.query1(asyncUser); 68 69 System.out.println("end : " + (System.currentTimeMillis() - start)); 70 System.out.println(asyncUser.toString()); 71 } 72}

函数类(错误示范:asyncInvoke 方法中阻塞查询数据库,是同步的):

1import org.apache.flink.configuration.Configuration; 2import org.apache.flink.streaming.api.functions.async.ResultFuture; 3import org.apache.flink.streaming.api.functions.async.RichAsyncFunction; 4import org.slf4j.Logger; 5import org.slf4j.LoggerFactory; 6import java.sql.DriverManager; 7import java.sql.PreparedStatement; 8import java.sql.ResultSet; 9import java.util.ArrayList; 10import java.util.List; 11 12public class AsyncFunctionForMysqlJava extends RichAsyncFunction<User, User> { 13 14 // 链接 15 private static String jdbcUrl = "jdbc:mysql://192.168.229.128:3306?useSSL=false"; 16 private static String username = "root"; 17 private static String password = "123456"; 18 private static String driverName = "com.mysql.jdbc.Driver"; 19 20 21 java.sql.Connection conn; 22 PreparedStatement ps; 23 Logger logger = LoggerFactory.getLogger(AsyncFunctionForMysqlJava.class); 24 25 /** 26 * open 方法中初始化链接 27 * @param parameters 28 * @throws Exception 29 */ 30 @Override 31 public void open(Configuration parameters) throws Exception { 32 logger.info("async function for hbase java open ..."); 33 super.open(parameters); 34 35 Class.forName(driverName); 36 conn = DriverManager.getConnection(jdbcUrl, username, password); 37 ps = conn.prepareStatement("select phone from async.async_test where id = ?"); 38 } 39 40 /** 41 * use user.getId async get user phone 42 * 43 * @param user 44 * @param resultFuture 45 * @throws Exception 46 */ 47 @Override 48 public void asyncInvoke(User user, ResultFuture<User> resultFuture) throws Exception { 49 // 使用 user id 查询 50 ps.setString(1, user.getId()); 51 ResultSet rs = ps.executeQuery(); 52 String phone = null; 53 if (rs.next()) { 54 phone = rs.getString(1); 55 } 56 user.setPhone(phone); 57 List<User> list = new ArrayList(); 58 list.add(user); 59 // 放回 result 队列 60 resultFuture.complete(list); 61 } 62 63 @Override 64 public void timeout(User input, ResultFuture<User> resultFuture) throws Exception { 65 logger.info("Async function for hbase timeout"); 66 List<User> list = new ArrayList(); 67 list.add(input); 68 resultFuture.complete(list); 69 } 70 71 /** 72 * close function 73 * 74 * @throws Exception 75 */ 76 @Override 77 public void close() throws Exception { 78 logger.info("async function for hbase java close ..."); 79 super.close(); 80 conn.close(); 81 } 82}

测试数据如下:

1{"id" : 1, "username" : "venn", "password" : 1561709530935} 2{"id" : 2, "username" : "venn", "password" : 1561709536029} 3{"id" : 3, "username" : "venn", "password" : 1561709541033} 4{"id" : 4, "username" : "venn", "password" : 1561709546037} 5{"id" : 5, "username" : "venn", "password" : 1561709551040} 6{"id" : 6, "username" : "venn", "password" : 1561709556044} 7{"id" : 7, "username" : "venn", "password" : 1561709561048}

执行结果如下:

1submit query : 1-1-1562763486845 2submit query : 2-1-1562763486846 3submit query : 3-1-1562763486846 4submit query : 4-1-1562763486849 5submit query : 5-1-1562763486849 6submit query : 6-1-1562763486859 7submit query : 7-1-1562763486913 8submit query : 8-1-1562763486967 9submit query : 9-1-1562763487021 10execute query : 1-2-phone : 12345678910-1562763487316 111> {"password":"1562763486506","phone":"12345678910","id":"1","username":"venn"} 12submit query : 10-1-1562763487408 13submit query : 11-1-1562763487408 14execute query : 9-2-phone : 1562661110630-1562763487633 151> {"password":"1562763487017","phone":"1562661110630","id":"9","username":"venn"} # 这里可以看到异步,提交查询的到 11 了,执行查询 的只有 1/9,返回了 1/9(unorderedWait 调用) 16submit query : 12-1-1562763487634 17execute query : 8-2-phone : 1562661110627-1562763487932 181> {"password":"1562763486963","phone":"1562661110627","id":"8","username":"venn"} 19submit query : 13-1-1562763487933 20execute query : 7-2-phone : 1562661110624-1562763488228 211> {"password":"1562763486909","phone":"1562661110624","id":"7","username":"venn"} 22submit query : 14-1-1562763488230 23execute query : 6-2-phone : 1562661110622-1562763488526 241> {"password":"1562763486855","phone":"1562661110622","id":"6","username":"venn"} 25submit query : 15-1-1562763488527 26execute query : 4-2-phone : 12345678913-1562763488832 271> {"password":"1562763486748","phone":"12345678913","id":"4","username":"venn"}

hbase、redis或其他实现类似

欢迎关注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(

MySQL部分从库上面因为大量的临时表tmp_table造成慢查询

背景描述Time:20190124T00:08:14.70572408:00User@Host:@Id:Schema:sentrymetaLast_errno:0Killed:0Query_time:0.315758Lock_

皕杰报表之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 )

Flink 异步IO访问外部数据(mysql篇) - HelloWorld