【问题标题】:how flink interacts with MySQL for the temporal join with mysqlflink 如何与 MySQL 交互以实现与 mysql 的 temporal join
【发布时间】:2021-08-07 22:24:58
【问题描述】:

我正在阅读

https://ci.apache.org/projects/flink/flink-docs-release-1.13/docs/dev/table/sql/queries/joins/#lookup-join,

它是使用 MySQL 作为时态表连接中的查找表

-- Customers is backed by the JDBC connector and can be used for lookup joins
CREATE TEMPORARY TABLE Customers (
  id INT,
  name STRING,
  country STRING,
  zip STRING
) WITH (
  'connector' = 'jdbc',
  'url' = 'jdbc:mysql://mysqlhost:3306/customerdb',
  'table-name' = 'customers'
);

-- enrich each order with customer information
SELECT o.order_id, o.total, c.country, c.zip
FROM Orders AS o
  JOIN Customers FOR SYSTEM_TIME AS OF o.proc_time AS c
    ON o.customer_id = c.id;

我想知道flink如何与mysql交互,mysql端是否存在temporal join mysql的性能问题。

基本问题是 flink 如何与 mysql 进行时间连接。

【问题讨论】:

    标签: apache-flink


    【解决方案1】:

    您可以在 Table / JDBC 连接器的文档中找到一些相关详细信息:https://ci.apache.org/projects/flink/flink-docs-stable/docs/connectors/table/jdbc/#features。请特别参阅描述查找缓存的部分,其中说

    JDBC 连接器可以在临时连接中用作查找源(又名维度表)。目前仅支持同步查找模式。

    默认情况下,查找缓存未启用。您可以通过设置lookup.cache.max-rows 和lookup.cache.ttl 来启用它。

    查找缓存用于提高临时连接 JDBC 连接器的性能。默认情况下,查找缓存未启用,因此所有请求都发送到外部数据库。启用查找缓存后,每个进程(即 TaskManager)都会持有一个缓存。 Flink 会先查找缓存,只有在缓存缺失时才会向外部数据库发送请求,并根据返回的行更新缓存。当缓存达到最大缓存行lookup.cache.max-rows 或行超过最大存活时间lookup.cache.ttl 时,缓存中最旧的行将过期。缓存的行可能不是最新的,用户可以将 lookup.cache.ttl 调整为较小的值以获得更好的新鲜数据,但这可能会增加发送到数据库的请求数。所以这是吞吐量和正确性之间的平衡。

    【讨论】:

    • 感谢@david-anderson 提供的有用答案,这是我所期望的。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-11-19
    • 2021-11-22
    • 2014-02-05
    • 2012-04-03
    相关资源
    最近更新 更多