您好!
我们用flink sql对kafka数据进行处理,版本为1.7。处理过程中要关联一张mysql
表,因为kafka数据量大,所以我们将mysql数据读入缓存,通过读取缓存实现关联。我从文档中了解到,flink sql 进行Lookup
Join,用到了'lookup.cache.max-rows' 、 'lookup.cache.ttl'
两个参数,会优先读缓存,如果缓存中不存在,则读mySQL表。但由于kafka数据量很大,关联读mysql表,flink进程经常卡顿。我想将mysql表全部加载到缓存,如果关联不到不要再去读外部表,并且只想用flink
sql实现,暂不考虑用scala/java编写flink程序,如何实现呢?
我试着通过lookup.cache.max-rows参数值设为比mysql表记录数大的值,没有达到效果。
以下是我的flink sql中mysql 表建表语句:
CREATE TABLE table_test
col1 string,
col2 string,
PRIMARY KEY (col1) NOT ENFORCED
--声明主键,用于lookup
WITH(
'connector' = 'mysql',
'jdbc-url' = 'jdbc:mysql://x.x.x.x:x/dbname',
'table.identifier' = 'table_test',
'username' = 'xxx',
'password' = 'xxx',
-缓存优化,减少查询压力
'lookup.cache.max-rows' = '100000',
'lookup.cache.ttl' = '3600s',
'lookup.jdbc.async' = 'true',
--开启异步查询
'lookup.jdbc.read.thread-size' = '16',
'lookup.jdbc.read.batch.size' = '256',
'lookup.jdbc.read.batch.queue-size' = '4096'
);