加入收藏 | 设为首页 | 会员中心 | 我要投稿 鹤壁站长网 (https://www.0392zz.cn/)- 分布式云、存储数据、视频终端、媒体处理、内容创作!
当前位置: 首页 > 运营中心 > 网站设计 > 教程 > 正文

Apache Flink 漫谈系列 - SQL概览

发布时间:2018-11-14 23:44:22 所属栏目:教程 来源:孙金城
导读:一、SQL简述 SQL是Structured Query Language的缩写,最初是由美国计算机科学家Donald D. Chamberlin和Raymond F. Boyce在20世纪70年代早期从 Early History of SQL 中了解关系模型后在IBM开发的。该版本最初称为[SEQUEL: A Structured English Query Lang

我们创建一个SqlOverviewITCase.scala 用于接下来介绍Flink SQL算子的功能体验。代码如下:

  1. import org.apache.flink.api.scala._ 
  2. import org.apache.flink.runtime.state.StateBackend 
  3. import org.apache.flink.runtime.state.memory.MemoryStateBackend 
  4. import org.apache.flink.streaming.api.TimeCharacteristic 
  5. import org.apache.flink.streaming.api.functions.sink.RichSinkFunction 
  6. import org.apache.flink.streaming.api.functions.source.SourceFunction 
  7. import org.apache.flink.streaming.api.functions.source.SourceFunction.SourceContext 
  8. import org.apache.flink.streaming.api.scala.StreamExecutionEnvironment 
  9. import org.apache.flink.streaming.api.watermark.Watermark 
  10. import org.apache.flink.table.api.TableEnvironment 
  11. import org.apache.flink.table.api.scala._ 
  12. import org.apache.flink.types.Row 
  13. import org.junit.rules.TemporaryFolder 
  14. import org.junit.{Rule, Test} 
  15.  
  16. import scala.collection.mutable 
  17. import scala.collection.mutable.ArrayBuffer 
  18.  
  19. class SqlOverviewITCase { 
  20. val _tempFolder = new TemporaryFolder 
  21.  
  22. @Rule 
  23. def tempFolder: TemporaryFolder = _tempFolder 
  24.  
  25. def getStateBackend: StateBackend = { 
  26. new MemoryStateBackend() 
  27. } 
  28.  
  29. // 客户表数据 
  30. val customer_data = new mutable.MutableList[(String, String, String)] 
  31. customer_data.+=(("c_001", "Kevin", "from JinLin")) 
  32. customer_data.+=(("c_002", "Sunny", "from JinLin")) 
  33. customer_data.+=(("c_003", "JinCheng", "from HeBei")) 
  34.  
  35.  
  36. // 订单表数据 
  37. val order_data = new mutable.MutableList[(String, String, String, String)] 
  38. order_data.+=(("o_001", "c_002", "2018-11-05 10:01:01", "iphone")) 
  39. order_data.+=(("o_002", "c_001", "2018-11-05 10:01:55", "ipad")) 
  40. order_data.+=(("o_003", "c_001", "2018-11-05 10:03:44", "flink book")) 
  41.  
  42. // 商品销售表数据 
  43. val item_data = Seq( 
  44. Left((1510365660000L, (1510365660000L, 20, "ITEM001", "Electronic"))), 
  45. Right((1510365660000L)), 
  46. Left((1510365720000L, (1510365720000L, 50, "ITEM002", "Electronic"))), 
  47. Right((1510365720000L)), 
  48. Left((1510365780000L, (1510365780000L, 30, "ITEM003", "Electronic"))), 
  49. Left((1510365780000L, (1510365780000L, 60, "ITEM004", "Electronic"))), 
  50. Right((1510365780000L)), 
  51. Left((1510365900000L, (1510365900000L, 40, "ITEM005", "Electronic"))), 
  52. Right((1510365900000L)), 
  53. Left((1510365960000L, (1510365960000L, 20, "ITEM006", "Electronic"))), 
  54. Right((1510365960000L)), 
  55. Left((1510366020000L, (1510366020000L, 70, "ITEM007", "Electronic"))), 
  56. Right((1510366020000L)), 
  57. Left((1510366080000L, (1510366080000L, 20, "ITEM008", "Clothes"))), 
  58. Right((151036608000L))) 
  59.  
  60. // 页面访问表数据 
  61. val pageAccess_data = Seq( 
  62. Left((1510365660000L, (1510365660000L, "ShangHai", "U0010"))), 
  63. Right((1510365660000L)), 
  64. Left((1510365660000L, (1510365660000L, "BeiJing", "U1001"))), 
  65. Right((1510365660000L)), 
  66. Left((1510366200000L, (1510366200000L, "BeiJing", "U2032"))), 
  67. Right((1510366200000L)), 
  68. Left((1510366260000L, (1510366260000L, "BeiJing", "U1100"))), 
  69. Right((1510366260000L)), 
  70. Left((1510373400000L, (1510373400000L, "ShangHai", "U0011"))), 
  71. Right((1510373400000L))) 
  72.  
  73. // 页面访问量表数据2 
  74. val pageAccessCount_data = Seq( 
  75. Left((1510365660000L, (1510365660000L, "ShangHai", 100))), 
  76. Right((1510365660000L)), 
  77. Left((1510365660000L, (1510365660000L, "BeiJing", 86))), 
  78. Right((1510365660000L)), 
  79. Left((1510365960000L, (1510365960000L, "BeiJing", 210))), 
  80. Right((1510366200000L)), 
  81. Left((1510366200000L, (1510366200000L, "BeiJing", 33))), 
  82. Right((1510366200000L)), 
  83. Left((1510373400000L, (1510373400000L, "ShangHai", 129))), 
  84. Right((1510373400000L))) 
  85.  
  86. // 页面访问表数据3 
  87. val pageAccessSession_data = Seq( 
  88. Left((1510365660000L, (1510365660000L, "ShangHai", "U0011"))), 
  89. Right((1510365660000L)), 
  90. Left((1510365720000L, (1510365720000L, "ShangHai", "U0012"))), 
  91. Right((1510365720000L)), 
  92. Left((1510365720000L, (1510365720000L, "ShangHai", "U0013"))), 
  93. Right((1510365720000L)), 
  94. Left((1510365900000L, (1510365900000L, "ShangHai", "U0015"))), 
  95. Right((1510365900000L)), 
  96. Left((1510366200000L, (1510366200000L, "ShangHai", "U0011"))), 
  97. Right((1510366200000L)), 
  98. Left((1510366200000L, (1510366200000L, "BeiJing", "U2010"))), 
  99. Right((1510366200000L)), 
  100. Left((1510366260000L, (1510366260000L, "ShangHai", "U0011"))), 
  101. Right((1510366260000L)), 
  102. Left((1510373760000L, (1510373760000L, "ShangHai", "U0410"))), 
  103. Right((1510373760000L))) 
  104.  
  105. def procTimePrint(sql: String): Unit = { 
  106. // Streaming 环境 
  107. val env = StreamExecutionEnvironment.getExecutionEnvironment 
  108. val tEnv = TableEnvironment.getTableEnvironment(env) 
  109.  
  110. // 将order_tab, customer_tab 注册到catalog 
  111. val customer = env.fromCollection(customer_data).toTable(tEnv).as('c_id, 'c_name, 'c_desc) 
  112. val order = env.fromCollection(order_data).toTable(tEnv).as('o_id, 'c_id, 'o_time, 'o_desc) 
  113.  
  114. tEnv.registerTable("order_tab", order) 
  115. tEnv.registerTable("customer_tab", customer) 
  116.  
  117. val result = tEnv.sqlQuery(sql).toRetractStream[Row] 
  118. val sink = new RetractingSink 
  119. result.addSink(sink) 
  120. env.execute() 
  121. } 
  122.  
  123. def rowTimePrint(sql: String): Unit = { 
  124. // Streaming 环境 
  125. val env = StreamExecutionEnvironment.getExecutionEnvironment 
  126. env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime) 
  127. env.setStateBackend(getStateBackend) 
  128. env.setParallelism(1) 
  129. val tEnv = TableEnvironment.getTableEnvironment(env) 
  130.  
  131. // 将item_tab, pageAccess_tab 注册到catalog 
  132. val item = 
  133. env.addSource(new EventTimeSourceFunction[(Long, Int, String, String)](item_data)) 
  134. .toTable(tEnv, 'onSellTime, 'price, 'itemID, 'itemType, 'rowtime.rowtime) 
  135.  
  136. val pageAccess = 
  137. env.addSource(new EventTimeSourceFunction[(Long, String, String)](pageAccess_data)) 
  138. .toTable(tEnv, 'accessTime, 'region, 'userId, 'rowtime.rowtime) 
  139.  
  140. val pageAccessCount = 
  141. env.addSource(new EventTimeSourceFunction[(Long, String, Int)](pageAccessCount_data)) 
  142. .toTable(tEnv, 'accessTime, 'region, 'accessCount, 'rowtime.rowtime) 
  143.  
  144. val pageAccessSession = 
  145. env.addSource(new EventTimeSourceFunction[(Long, String, String)](pageAccessSession_data)) 
  146. .toTable(tEnv, 'accessTime, 'region, 'userId, 'rowtime.rowtime) 
  147.  
  148. tEnv.registerTable("item_tab", item) 
  149. tEnv.registerTable("pageAccess_tab", pageAccess) 
  150. tEnv.registerTable("pageAccessCount_tab", pageAccessCount) 
  151. tEnv.registerTable("pageAccessSession_tab", pageAccessSession) 
  152.  
  153. val result = tEnv.sqlQuery(sql).toRetractStream[Row] 
  154. val sink = new RetractingSink 
  155. result.addSink(sink) 
  156. env.execute() 
  157.  
  158. } 
  159.  
  160. @Test 
  161. def testSelect(): Unit = { 
  162. val sql = "替换想要测试的SQL" 
  163. // 非window 相关用 procTimePrint(sql) 
  164. // Window 相关用 rowTimePrint(sql) 
  165. } 
  166.  
  167. } 
  168.  
  169. // 自定义Sink 
  170. final class RetractingSink extends RichSinkFunction[(Boolean, Row)] { 
  171. var retractedResults: ArrayBuffer[String] = mutable.ArrayBuffer.empty[String] 
  172.  
  173. def invoke(v: (Boolean, Row)) { 
  174. retractedResults.synchronized { 
  175. val vvalue = v._2.toString 
  176. if (v._1) { 
  177. retractedResults += value 
  178. } else { 
  179. val idx = retractedResults.indexOf(value) 
  180. if (idx >= 0) { 
  181. retractedResults.remove(idx) 
  182. } else { 
  183. throw new RuntimeException("Tried to retract a value that wasn't added first. " + 
  184. "This is probably an incorrectly implemented test. " + 
  185. "Try to set the parallelism of the sink to 1.") 
  186. } 
  187. } 
  188. } 
  189. retractedResults.sorted.foreach(println(_)) 
  190. } 
  191. } 
  192.  
  193. // Water mark 生成器 
  194. class EventTimeSourceFunction[T]( 
  195. dataWithTimestampList: Seq[Either[(Long, T), Long]]) extends SourceFunction[T] { 
  196. override def run(ctx: SourceContext[T]): Unit = { 
  197. dataWithTimestampList.foreach { 
  198. case Left(t) => ctx.collectWithTimestamp(t._2, t._1) 
  199. case Right(w) => ctx.emitWatermark(new Watermark(w)) 
  200. } 
  201. } 
  202.  
  203. override def cancel(): Unit = ??? 
  204. } 

五、Select

(编辑:鹤壁站长网)

【声明】本站内容均来自网络,其相关言论仅代表作者个人观点,不代表本站立场。若无意侵犯到您的权利,请及时与联系站长删除相关内容!

推荐文章
    热点阅读