FlinkSQL(JDBC)支持谓词下推、支持复杂SQL查询、并行度、CLOB
·
生产环境中使用FlinkSQL(JDBC)做离线数据同步,有以下几种场景得不到很好的解决
1. 不需要同步源表所有数据
2. 同步的表需要根据源库中另一张表做筛选
3. 如果源表中没有一个递增的数值类型字段作为分区采集字段,无法实现多并行度采集
4. Oracle中CLOB类型字段的解析
记录一下解决方案
正常的表做离线全量采集的写法:

问题1解决方案:
修改源码增加自定义的数据过滤方式

问题2解决方案:
修改源码增加自定义的SQL查询,增加custom-sql属性

问题3解决方案:
修改源码增加多并行度的实现,增加scan.parallel.enabled、scan.parallel.column、scan.parallel.number 三个属性,确定开启并行采集

源码修改:
1. JdbcConnectorOptions 增加属性
DATA_FILTER、SCAN_PARALLEL_ENABLED、SCAN_PARALLEL_COLUMN、SCAN_PARALLEL_NUMBER、CUSTOM_SQL、CLOB_COLUMNS
2. JdbcDynamicTableFactory增加属性的赋值
3. JdbcDynamicTableSource修改
getScanRuntimeProvider方法中的dialect.getSelectFromStatement,如下图:原本直接拼接表名的地方,现在修改为拼接customSql,增加dataFilter的拼接

增加以下代码支持多并行度

增加以下代码支持Oracle的Clob类型字段的解析

更多推荐
所有评论(0)