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

Apache Flink 漫谈系列 - 流表对偶(duality)性

发布时间:2018-11-11 10:12:26 所属栏目:教程 来源:孙金城
导读:实际问题 很多大数据计算产品,都对用户提供了SQL API,比如Hive, Spark, Flink等,那么SQL作为传统关系数据库的查询语言,是应用在批查询场景的。Hive和Spark本质上都是Batch的计算模式(在《Apache Flink 漫谈系列 - 概述》我们介绍过Spark是Micro Batchi

打开binlog.txt 内容如下:

  1. /*!50530 SET @@SESSION.PSEUDO_SLAVE_MODE=1*/; 
  2. /*!50003 SET @OLD_COMPLETION_TYPE=@@COMPLETION_TYPE,COMPLETION_TYPE=0*/; 
  3. DELIMITER /*!*/; 
  4. # at 4 
  5. #180430 22:29:33 server id 1  end_log_pos 124 CRC32 0xff61797c  Start: binlog v 4, server v 8.0.11 created 180430 22:29:33 at startup 
  6. # Warning: this binlog is either in use or was not closed properly. 
  7. ROLLBACK/*!*/; 
  8. # at 124 
  9. #180430 22:29:33 server id 1  end_log_pos 155 CRC32 0x629ae755  Previous-GTIDs 
  10. # [empty] 
  11. # at 155 
  12. #180430 22:32:11 server id 1  end_log_pos 228 CRC32 0xbde49fca  Anonymous_GTID  last_committed=0        sequence_number=1       rbr_only=no     original_committed_timestamp=1525098731207902   immediate_commit_timestamp=1525098731207902     transaction_length=213 
  13. # original_commit_timestamp=1525098731207902 (2018-04-30 22:32:11.207902 CST) 
  14. # immediate_commit_timestamp=1525098731207902 (2018-04-30 22:32:11.207902 CST) 
  15. /*!80001 SET @@session.original_commit_timestamp=1525098731207902*//*!*/; 
  16. SET @@SESSION.GTID_NEXT= 'ANONYMOUS'/*!*/; 
  17. # at 228 
  18. #180430 22:32:11 server id 1  end_log_pos 368 CRC32 0xe5f330e7  Query   thread_id=9     exec_time=0     error_code=0    Xid = 22 
  19. use `Apache Flinkdb`/*!*/; 
  20. SET TIMESTAMP=1525098731/*!*/; 
  21. SET @@session.pseudo_thread_id=9/*!*/; 
  22. SET @@session.foreign_key_checks=1, @@session.sql_auto_is_null=0, @@session.unique_checks=1, @@session.autocommit=1/*!*/; 
  23. SET @@session.sql_mode=1168113696/*!*/; 
  24. SET @@session.auto_increment_increment=1, @@session.auto_increment_offset=1/*!*/; 
  25. /*!C utf8mb4 *//*!*/; 
  26. SET @@session.character_set_client=255,@@session.collation_connection=255,@@session.collation_server=255/*!*/; 
  27. SET @@session.lc_time_names=0/*!*/; 
  28. SET @@session.collation_database=DEFAULT/*!*/; 
  29. /*!80005 SET @@session.default_collation_for_utf8mb4=255*//*!*/; 
  30. DROP TABLE `tab` /* generated by server */ 
  31. /*!*/; 
  32. # at 368 
  33. #180430 22:32:21 server id 1  end_log_pos 443 CRC32 0x50e5acb7  Anonymous_GTID  last_committed=1        sequence_number=2       rbr_only=no     original_committed_timestamp=1525098741628960   immediate_commit_timestamp=1525098741628960     transaction_length=302 
  34. # original_commit_timestamp=1525098741628960 (2018-04-30 22:32:21.628960 CST) 
  35. # immediate_commit_timestamp=1525098741628960 (2018-04-30 22:32:21.628960 CST) 
  36. /*!80001 SET @@session.original_commit_timestamp=1525098741628960*//*!*/; 
  37. SET @@SESSION.GTID_NEXT= 'ANONYMOUS'/*!*/; 
  38. # at 443 
  39. #180430 22:32:21 server id 1  end_log_pos 670 CRC32 0xe1353dd6  Query   thread_id=9     exec_time=0     error_code=0    Xid = 23 
  40. SET TIMESTAMP=1525098741/*!*/; 
  41. create table tab( 
  42.    id INT NOT NULL AUTO_INCREMENT, 
  43.    user VARCHAR(100) NOT NULL, 
  44.    clicks INT NOT NULL, 
  45.    PRIMARY KEY (id) 
  46. ) 
  47. /*!*/; 
  48. # at 670 
  49. #180430 22:36:53 server id 1  end_log_pos 745 CRC32 0xcf436fbb  Anonymous_GTID  last_committed=2        sequence_number=3       rbr_only=yes    original_committed_timestamp=1525099013988373   immediate_commit_timestamp=1525099013988373     transaction_length=301 
  50. /*!50718 SET TRANSACTION ISOLATION LEVEL READ COMMITTED*//*!*/; 
  51. # original_commit_timestamp=1525099013988373 (2018-04-30 22:36:53.988373 CST) 
  52. # immediate_commit_timestamp=1525099013988373 (2018-04-30 22:36:53.988373 CST) 
  53. /*!80001 SET @@session.original_commit_timestamp=1525099013988373*//*!*/; 
  54. SET @@SESSION.GTID_NEXT= 'ANONYMOUS'/*!*/; 
  55. # at 745 
  56. #180430 22:36:53 server id 1  end_log_pos 823 CRC32 0x71c64dd2  Query   thread_id=9     exec_time=0     error_code=0 
  57. SET TIMESTAMP=1525099013/*!*/; 
  58. BEGIN 
  59. /*!*/; 
  60. # at 823 
  61. #180430 22:36:53 server id 1  end_log_pos 890 CRC32 0x63792f6b  Table_map: `Apache Flinkdb`.`tab` mapped to number 96 
  62. # at 890 
  63. #180430 22:36:53 server id 1  end_log_pos 940 CRC32 0xf2dade22  Write_rows: table id 96 flags: STMT_END_F 
  64. ### INSERT INTO `Apache Flinkdb`.`tab` 
  65. ### SET 
  66. ###   @11=1 
  67. ###   @2='Mary' 
  68. ###   @3=1 
  69. # at 940 
  70. #180430 22:36:53 server id 1  end_log_pos 971 CRC32 0x7db3e61e  Xid = 25 
  71. COMMIT/*!*/; 
  72. # at 971 
  73. #180430 22:37:06 server id 1  end_log_pos 1046 CRC32 0xd05dd12c         Anonymous_GTID  last_committed=3        sequence_number=4       rbr_only=yes    original_committed_timestamp=1525099026328547   immediate_commit_timestamp=1525099026328547     transaction_length=300 
  74. /*!50718 SET TRANSACTION ISOLATION LEVEL READ COMMITTED*//*!*/; 
  75. # original_commit_timestamp=1525099026328547 (2018-04-30 22:37:06.328547 CST) 
  76. # immediate_commit_timestamp=1525099026328547 (2018-04-30 22:37:06.328547 CST) 
  77. /*!80001 SET @@session.original_commit_timestamp=1525099026328547*//*!*/; 
  78. SET @@SESSION.GTID_NEXT= 'ANONYMOUS'/*!*/; 
  79. # at 1046 
  80. #180430 22:37:06 server id 1  end_log_pos 1124 CRC32 0x80f259e0         Query   thread_id=9     exec_time=0     error_code=0 
  81. SET TIMESTAMP=1525099026/*!*/; 
  82. BEGIN 
  83. /*!*/; 
  84. # at 1124 
  85. #180430 22:37:06 server id 1  end_log_pos 1191 CRC32 0x255903ba         Table_map: `Apache Flinkdb`.`tab` mapped to number 96 
  86. # at 1191 
  87. #180430 22:37:06 server id 1  end_log_pos 1240 CRC32 0xe76bfc79         Write_rows: table id 96 flags: STMT_END_F 
  88. ### INSERT INTO `Apache Flinkdb`.`tab` 
  89. ### SET 
  90. ###   @1=2 
  91. ###   @2='Bob' 
  92. ###   @3=1 
  93. # at 1240 
  94. #180430 22:37:06 server id 1  end_log_pos 1271 CRC32 0x83cddfef         Xid = 26 
  95. COMMIT/*!*/; 
  96. # at 1271 
  97. #180430 22:37:15 server id 1  end_log_pos 1346 CRC32 0x7095baee         Anonymous_GTID  last_committed=4        sequence_number=5       rbr_only=yes    original_committed_timestamp=1525099035811597   immediate_commit_timestamp=1525099035811597     transaction_length=326 
  98. /*!50718 SET TRANSACTION ISOLATION LEVEL READ COMMITTED*//*!*/; 
  99. # original_commit_timestamp=1525099035811597 (2018-04-30 22:37:15.811597 CST) 
  100. # immediate_commit_timestamp=1525099035811597 (2018-04-30 22:37:15.811597 CST) 
  101. /*!80001 SET @@session.original_commit_timestamp=1525099035811597*//*!*/; 
  102. SET @@SESSION.GTID_NEXT= 'ANONYMOUS'/*!*/; 
  103. # at 1346 
  104. #180430 22:37:15 server id 1  end_log_pos 1433 CRC32 0x70ef97e2         Query   thread_id=9     exec_time=0     error_code=0 
  105. SET TIMESTAMP=1525099035/*!*/; 
  106. BEGIN 
  107. /*!*/; 
  108. # at 1433 
  109. #180430 22:37:15 server id 1  end_log_pos 1500 CRC32 0x75f1f399         Table_map: `Apache Flinkdb`.`tab` mapped to number 96 
  110. # at 1500 
  111. #180430 22:37:15 server id 1  end_log_pos 1566 CRC32 0x256bd4b8         Update_rows: table id 96 flags: STMT_END_F 
  112. ### UPDATE `Apache Flinkdb`.`tab` 
  113. ### WHERE 
  114. ###   @11=1 
  115. ###   @2='Mary' 
  116. ###   @3=1 
  117. ### SET 
  118. ###   @11=1 
  119. ###   @2='Mary' 
  120. ###   @3=2 
  121. # at 1566 
  122. #180430 22:37:15 server id 1  end_log_pos 1597 CRC32 0x93c86579         Xid = 27 
  123. COMMIT/*!*/; 
  124. # at 1597 
  125. #180430 22:37:27 server id 1  end_log_pos 1672 CRC32 0xe8bd63e7         Anonymous_GTID  last_committed=5        sequence_number=6       rbr_only=yes    original_committed_timestamp=1525099047219517   immediate_commit_timestamp=1525099047219517     transaction_length=300 
  126. /*!50718 SET TRANSACTION ISOLATION LEVEL READ COMMITTED*//*!*/; 
  127. # original_commit_timestamp=1525099047219517 (2018-04-30 22:37:27.219517 CST) 
  128. # immediate_commit_timestamp=1525099047219517 (2018-04-30 22:37:27.219517 CST) 
  129. /*!80001 SET @@session.original_commit_timestamp=1525099047219517*//*!*/; 
  130. SET @@SESSION.GTID_NEXT= 'ANONYMOUS'/*!*/; 
  131. # at 1672 
  132. #180430 22:37:27 server id 1  end_log_pos 1750 CRC32 0x5356c3c7         Query   thread_id=9     exec_time=0     error_code=0 
  133. SET TIMESTAMP=1525099047/*!*/; 
  134. BEGIN 
  135. /*!*/; 
  136. # at 1750 
  137. #180430 22:37:27 server id 1  end_log_pos 1817 CRC32 0x37e6b1ce         Table_map: `Apache Flinkdb`.`tab` mapped to number 96 
  138. # at 1817 
  139. #180430 22:37:27 server id 1  end_log_pos 1866 CRC32 0x6ab1bbe6         Write_rows: table id 96 flags: STMT_END_F 
  140. ### INSERT INTO `Apache Flinkdb`.`tab` 
  141. ### SET 
  142. ###   @1=3 
  143. ###   @2='Llz' 
  144. ###   @3=1 
  145. # at 1866 
  146. #180430 22:37:27 server id 1  end_log_pos 1897 CRC32 0x3b62b153         Xid = 28 
  147. COMMIT/*!*/; 
  148. # at 1897 
  149. #180430 22:37:36 server id 1  end_log_pos 1972 CRC32 0x603134c1         Anonymous_GTID  last_committed=6        sequence_number=7       rbr_only=yes    original_committed_timestamp=1525099056866022   immediate_commit_timestamp=1525099056866022     transaction_length=324 
  150. /*!50718 SET TRANSACTION ISOLATION LEVEL READ COMMITTED*//*!*/; 
  151. # original_commit_timestamp=1525099056866022 (2018-04-30 22:37:36.866022 CST) 
  152. # immediate_commit_timestamp=1525099056866022 (2018-04-30 22:37:36.866022 CST) 
  153. /*!80001 SET @@session.original_commit_timestamp=1525099056866022*//*!*/; 
  154. SET @@SESSION.GTID_NEXT= 'ANONYMOUS'/*!*/; 
  155. # at 1972 
  156. #180430 22:37:36 server id 1  end_log_pos 2059 CRC32 0xe17df4e4         Query   thread_id=9     exec_time=0     error_code=0 
  157. SET TIMESTAMP=1525099056/*!*/; 
  158. BEGIN 
  159. /*!*/; 
  160. # at 2059 
  161. #180430 22:37:36 server id 1  end_log_pos 2126 CRC32 0x53888b05         Table_map: `Apache Flinkdb`.`tab` mapped to number 96 
  162. # at 2126 
  163. #180430 22:37:36 server id 1  end_log_pos 2190 CRC32 0x85f34996         Update_rows: table id 96 flags: STMT_END_F 
  164. ### UPDATE `Apache Flinkdb`.`tab` 
  165. ### WHERE 
  166. ###   @1=2 
  167. ###   @2='Bob' 
  168. ###   @3=1 
  169. ### SET 
  170. ###   @1=2 
  171. ###   @2='Bob' 
  172. ###   @3=2 
  173. # at 2190 
  174. #180430 22:37:36 server id 1  end_log_pos 2221 CRC32 0x877f1e23         Xid = 29 
  175. COMMIT/*!*/; 
  176. # at 2221 
  177. #180430 22:37:45 server id 1  end_log_pos 2296 CRC32 0xfbc7e868         Anonymous_GTID  last_committed=7        sequence_number=8       rbr_only=yes    original_committed_timestamp=1525099065089940   immediate_commit_timestamp=1525099065089940     transaction_length=326 
  178. /*!50718 SET TRANSACTION ISOLATION LEVEL READ COMMITTED*//*!*/; 
  179. # original_commit_timestamp=1525099065089940 (2018-04-30 22:37:45.089940 CST) 
  180. # immediate_commit_timestamp=1525099065089940 (2018-04-30 22:37:45.089940 CST) 
  181. /*!80001 SET @@session.original_commit_timestamp=1525099065089940*//*!*/; 
  182. SET @@SESSION.GTID_NEXT= 'ANONYMOUS'/*!*/; 
  183. # at 2296 
  184. #180430 22:37:45 server id 1  end_log_pos 2383 CRC32 0x8a514364         Query   thread_id=9     exec_time=0     error_code=0 
  185. SET TIMESTAMP=1525099065/*!*/; 
  186. BEGIN 
  187. /*!*/; 
  188. # at 2383 
  189. #180430 22:37:45 server id 1  end_log_pos 2450 CRC32 0xdf18ca60         Table_map: `Apache Flinkdb`.`tab` mapped to number 96 
  190. # at 2450 
  191. #180430 22:37:45 server id 1  end_log_pos 2516 CRC32 0xd50de69f         Update_rows: table id 96 flags: STMT_END_F 
  192. ### UPDATE `Apache Flinkdb`.`tab` 
  193. ### WHERE 
  194. ###   @11=1 
  195. ###   @2='Mary' 
  196. ###   @3=2 
  197. ### SET 
  198. ###   @11=1 
  199. ###   @2='Mary' 
  200. ###   @33=3 
  201. # at 2516 
  202. #180430 22:37:45 server id 1  end_log_pos 2547 CRC32 0x94f89393         Xid = 30 
  203. COMMIT/*!*/; 
  204. SET @@SESSION.GTID_NEXT= 'AUTOMATIC' /* added by MySQLbinlog */ /*!*/; 
  205. DELIMITER ; 
  206. # End of log file 
  207. /*!50003 SET COMPLETION_TYPE=@OLD_COMPLETION_TYPE*/; 
  208. /*!50530 SET @@SESSION.PSEUDO_SLAVE_MODE=0*/; 
  • 梳理操作和binlog的记录关系
  • Apache Flink 漫谈系列 - 流表对偶(duality)性

    梳理操作和binlog的记录关系

    Apache Flink 漫谈系列 - 流表对偶(duality)性

  • 简化一下binlog
  • Apache Flink 漫谈系列 - 流表对偶(duality)性

  • replay binlog会得到如下表数据(按timestamp顺序)
  • Apache Flink 漫谈系列 - 流表对偶(duality)性

  • 表与binlog的关系简单示意如下
  • Apache Flink 漫谈系列 - 流表对偶(duality)性

流表对偶(duality)性

(编辑:鹤壁站长网)

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

推荐文章
    热点阅读