【问题标题】:how to read data from mysql to flink parallelly?如何从mysql并行读取数据到flink?
【发布时间】:2017-09-27 10:06:37
【问题描述】:

如何从mysql并行读取数据到flink?我想构建一个sourceFunction每隔一段时间从mysql并行读取数据,如何实现?

【问题讨论】:

    标签: mysql database streaming apache-flink flink-streaming


    【解决方案1】:

    这个问题的答案包括两个方面:

    1. 从 MySQL(或任何其他 JDBC 源)并行读取
    2. 定期从 MySQL(或任何其他 JDBC 源)读取数据

    从 MySQL 并行读取

    为了从 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)被索引。尽管如此,在单个数据库实例上执行多个查询会导致开销。因此,找到提供最佳性能的并行性并非易事。

    定期从 MySQL 读取

    这与并行读取类似。同样,您需要对查询进行分区。但是现在您希望根据描述记录时间的属性来执行此操作。因此,在每个间隔中,您都想询问自上次间隔以来插入的行。同样,这将通过时间属性上的范围谓词来完成。

    Q at T1: SELECT * FROM sourceT WHERE rowtime < T1;
    Q at T2: SELECT * FROM sourceT WHERE rowtime < T2;
    

    和以前一样,这只有在表在rowtime 属性上建立索引时才有效。否则,您将执行全表扫描,并且随着插入更多数据,查询将变得越来越慢。

    定期从 MySQL 并行读取

    为此,您“只需”结合这两种方法并向每个查询添加两个谓词。您实际上所做的是将表划分为分离的部分,并随着时间的推移并行读取它们。

    但是,正如我之前指出的,确切的分区取决于您的数据和用例。此外,您需要创建适当的索引以避免全表扫描。另请注意,使用上述方法,您不会看到在读取后修改的行的任何更新。

    【讨论】:

    • 非常感谢您的回答!另外,我想问一下如何在flink中拆分查询并将不同的查询分配给不同的源工作人员?非常感谢!
    • 您可以使用并行子任务的数量和当前子任务的id来计算每个并行实例的分区。这为您提供了分区的数量以及每个子任务负责的分区。此信息可在RuntimeContext 上获得,RichFunction 可以请求该信息。
    • 谢谢先生!另外,我想问一下如何设置sourceFunction的并行度?当我设置“StreamExecutionEnvironment.setParallelism(N)”时,源的并行度也会是N吗?
    • 仅当您将其实现为ParallelSourceFunction。如果你扩展 SourceFunction 接口,它的并行度总​​是 1。
    猜你喜欢
    • 1970-01-01
    • 2016-05-05
    • 2016-11-08
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多