【发布时间】:2017-09-27 10:06:37
【问题描述】:
如何从mysql并行读取数据到flink?我想构建一个sourceFunction每隔一段时间从mysql并行读取数据,如何实现?
【问题讨论】:
标签: mysql database streaming apache-flink flink-streaming
如何从mysql并行读取数据到flink?我想构建一个sourceFunction每隔一段时间从mysql并行读取数据,如何实现?
【问题讨论】:
标签: mysql database streaming apache-flink flink-streaming
这个问题的答案包括两个方面:
为了从 MySQL 中并行读取,您需要发送多个不同的查询。查询的组合方式必须使其结果的并集等同于预期结果。例如,您可以使用范围谓词在数字属性之间拆分查询:
Q1: SELECT * FROM sourceT WHERE num < 10;
Q2: SELECT * FROM sourceT WHERE num >= 10 AND num < 20;
Q3: SELECT * FROM sourceT WHERE num >= 20;
还有其他方法可以对查询进行分区。但是为了真正获得一些东西,DBMS 必须能够比查询整个表的单个查询更有效地处理多个查询。所以通常,您希望确保您分区的属性(上面示例中的num)被索引。尽管如此,在单个数据库实例上执行多个查询会导致开销。因此,找到提供最佳性能的并行性并非易事。
这与并行读取类似。同样,您需要对查询进行分区。但是现在您希望根据描述记录时间的属性来执行此操作。因此,在每个间隔中,您都想询问自上次间隔以来插入的行。同样,这将通过时间属性上的范围谓词来完成。
Q at T1: SELECT * FROM sourceT WHERE rowtime < T1;
Q at T2: SELECT * FROM sourceT WHERE rowtime < T2;
和以前一样,这只有在表在rowtime 属性上建立索引时才有效。否则,您将执行全表扫描,并且随着插入更多数据,查询将变得越来越慢。
为此,您“只需”结合这两种方法并向每个查询添加两个谓词。您实际上所做的是将表划分为分离的部分,并随着时间的推移并行读取它们。
但是,正如我之前指出的,确切的分区取决于您的数据和用例。此外,您需要创建适当的索引以避免全表扫描。另请注意,使用上述方法,您不会看到在读取后修改的行的任何更新。
【讨论】:
RuntimeContext 上获得,RichFunction 可以请求该信息。
ParallelSourceFunction。如果你扩展 SourceFunction 接口,它的并行度总是 1。