Flink Connector之插件机制
Connector是Flink与外部存储连接的桥梁,Flink在对接外部存储和数据格式上
对于一个外部存储系统的库表来说,需要处理三方面的定义:
- connector: 定义了流、批、读、写的实现
- Format: 定义了数据的解析
- Schema: 定义了表的字段信息
Flink使用Java SPI机制实现插件,方便扩展和管理。Flink更专注于计算,
而各种不同的Connector专注于实现批量读写、流式读写、数据一致性协议、分区、offset和时间戳生成等逻辑
本文主要介绍Flink Connector插件机制
Java SPI
(Service provider interface,服务提供者接口),是Java提供的一套实现扩展的API。
Java SPI是符合接口隔离原则的设计,实现调用者和实现者的解耦,不同的实现方实现了可插拔:
- 接口定义了交互的”协议”
- 根据不同的实现方案,定义该接口的不同实现类
- 调用者根据实际的使用需要,启用、扩展或者替换框架的实现策略
常见的SPI的例子:
- 数据库驱动加载: JDBC加载不同类型的数据库的驱动
- 日志门面实现类的加载: Slf4J加载不同提供商的日志实现类
使用介绍:
- 服务提供者提供接口的实现类,实现类必须含义不带参数的改造方法
- 在Jar包的
META-INF/services目录下创建一个以“接口全限定名”为命名的文件,内容为实现类的全限定名; - 将接口的实现类所在的Jar包放在
classpath中 java.util.ServiceLoder通过扫描META-INF/services目录下的配置文件找到实现类,加载类到JVM
ServiceLoader
ServiceLoader可以跨越jar包获取META-INF下的配置文件- 通过反射方法Class.forName()加载类对象,并用instance()方法将类实例化
- 把实例化后的类缓存到providers对象中,(
LinkedHashMap<String,S>类型)
然后返回实例对象
性能损失
- 延迟加载,需要遍历全部和加载全部实现类
- 非线程安全
TableFactory的定义
flink-table-common模块中定义了TableFactory的接口、调用者和工具类,
TableFactory
1 | package org.apache.flink.table.factories; |
以flink-hbase为例, 在META-INF/services/中定义了文件org.apache.flink.table.factories.TableFactory
文件的内容为:
1 | org.apache.flink.addons.hbase.HBaseTableFactory |
这个类就是TableFactory的实现
TableFactory查找
Flink 在 Connector的实现上,也是采用SPI机制

graph LR A[SPI查找实现类]-->B[是否TableFactory的子类] B-->C[是否满足必要属性] C-->D[是否满足属性匹配] D-->E[创建TableFactory实现类]
入口方法: org.apache.flink.table.factories.TableFactoryService#find(java.lang.Class<T>, java.util.Map<java.lang.String,java.lang.String>)
调用样例:
1 | TableFactoryService.find(classOf[TableFactory[_]],sinkProperties) |
最终调用查找方法
1 | private static <T extends TableFactory> T findSingleInternal( |
查找实现类
1 | private static List<TableFactory> discoverFactories(Optional<ClassLoader> classLoader) { |
过滤
有三层过滤:
- filterByFactoryClass: 判断是否为TableFactory的实现类
- filterByContext: 判断必要属性是否匹配: 来源于TableFactory.requiredContext
- filterBySupportedProperties: 判断属性是否支持:来源于TableFactory.supportedProperties
TableFactory的分类

对于一个外部存储系统的库表来说,需要处理三方面的定义:
- connector: 定义了流、批、读、写的实现
- Format: 定义了数据的解析
- Schema: 定义了表的字段信息
Schema信息可以在注册库表的时候直接增加,connector、format信息通过TableFactory实现可插拔

classDiagram
TableFactory <|-- StreamTableSourceFactory~T~
TableFactory <|-- StreamTableSinkFactory~T~
TableFactory <|-- BatchTableSourceFactory~T~
TableFactory <|-- BatchTableSinkFactory~T~
TableFactory <|-- TableFormatFactory~T~
StreamTableSourceFactory <|-- HBaseTableFactory
StreamTableSinkFactory <|-- HBaseTableFactory
class TableFactory{
+requiredConext() Map~StringString~
+supportedProperties() List~String~
}
class StreamTableSourceFactory{
+createStreamTableSink(Map) StreamTableSink~T~
+createTableSink(Map) TableSink~T~
}
class StreamTableSinkFactory{
+createStreamTableSink(Map) StreamTableSink~T~
+createTableSink(Map) TableSink~T~
}
class BatchTableSourceFactory{
+createBatchTableSource(Map) BatchTableSource~T~
+createTableSource(Map) TableSource~T~
}
class BatchTableSinkFactory{
+createBatchTableSink(Map) BatchTableSink~T~
+createTableSink(Map) TableSink~T~
}
class TableFormatFactory{
+supportsSchemaDerivation()
+supportedProperties() List~String~
}
class HBaseTableFactory{
}
TableFactory可以分为:
- StreamTableSourceFactory
- StreamTableSinkFactory
- BatchTableSourceFactory
- BatchTableSinkFactory
- TableFormatFactory
常见的Connector样例
援引一张1.10的优化设计图

kafka-DDL定义
1 | CREATE TABLE orders_kafka ( |
MySQL-DDL定义
1 | CREATE TABLE agg_result ( |
Connector设计要点
- 自定义Factory,根据需要实现StreamTableSourceFactory和StreamTableSinkFactory
- 根据需要继承ConnectorDescriptorValidator,定义自己的connector参数(with 后面跟的那些)
- Factory中的requiredContext、supportedProperties都比较重要,框架中对Factory的过滤和检查需要他们
- 需要自定义个TableSink,根据你需要连接的中间件选择是AppendStreamTableSink、Upsert、Retract重写consumeDataStream方法
- 自定义一个SinkFunction,在invoke方法中实现将数据写入到外部中间件。
SourceFunction
SourceFunction是定义Flink Source的根接口,其源码如下。
1 |
|
SourceFunction接口定义了run()方法,该方法用于源源不断地产生源数据,因此重写的时候一般都写成循环,用标志位控制是否结束。cancel()方法则用来打断run()方法中的循环,终止产生数据的过程。
SourceFunction中还嵌套定义了SourceContext接口,它表示这个Source对应的上下文,用来发射数据。其中起主要作用的是前三个方法:
collect():发射一个不带自定义时间戳的元素。如果流程序的时间特征(TimeCharacteristic)是处理时间(ProcessingTime),元素没有时间戳;如果是摄入时间(IngestionTime),元素会附带系统时间;如果是事件时间(EventTime),那么初始没有时间戳,但一旦要做与时间戳相关的操作(如窗口)时,就必须用TimestampAssigner设定一个。
collectWithTimestamp():发射一个带有自定义时间戳的元素。该方法对于时间特征为事件时间的程序是绝对必须的,如果为处理时间就会被直接忽略,如果为摄入时间就会被系统时间覆盖。
emitWatermark():发射一个水印,仅对于事件时间有效。一个带有时间戳t的水印表示不会有任何t’ <= t的事件再发生,如果发生,会被当做迟到事件忽略掉。
SourceFunction还有一些其他实现,如:
ParallelSourceFunction,表示该Source可以按照设置的并行度并发执行。
RichSourceFunction,继承自富函数RichFunction,表示该Source可以感知到运行时上下文(RuntimeContext,如Task、State、并行度的信息),以及可以自定义初始化和销毁逻辑(通过open()/close()方法)。
RichParallelSourceFunction,以上两者的综合。
SinkFunction
SinkFunction是自定义Sink的根接口,其源码如下。
1 | public interface SinkFunction<IN> extends Function, Serializable { |
它的定义比SourceFunction要简单,只有一个invoke()方法,对收集来的每条数据都会调用它来处理。SinkFunction也有对应的上下文对象Context,可以从中获得当前处理时间、当前水印和时间戳。它也有衍生出来的富函数版本RichSinkFunction。
Flink内部提供了一个最简单的实现DiscardingSink。顾名思义,就是将所有汇集的数据全部丢弃。
1 |
|
参考文献