Flink自定義Sink端實(shí)現(xiàn)過程講解
Sink介紹
在Fink官網(wǎng)中sink端只是給出了常規(guī)的write api.在我們實(shí)際開發(fā)場景中需要將flink處理的數(shù)據(jù)寫入kafka,hbase kudu等外部系統(tǒng)。
UML關(guān)系
自定義Sink需要實(shí)現(xiàn)父類的接口和繼承抽象類。

上面是Sink的繼承關(guān)系
Flink addSink
// 方法需要SinkFunction的對象
public DataStreamSink<T> addSink(SinkFunction<T> sinkFunction) {
// read the output type of the input Transform to coax out errors about MissingTypeInfo
transformation.getOutputType();
// configure the type if needed
if (sinkFunction instanceof InputTypeConfigurable) {
((InputTypeConfigurable) sinkFunction).setInputType(getType(), getExecutionConfig());
}
StreamSink<T> sinkOperator = new StreamSink<>(clean(sinkFunction));
DataStreamSink<T> sink = new DataStreamSink<>(this, sinkOperator);
getExecutionEnvironment().addOperator(sink.getTransformation());
return sink;
}
SinkFunction
// SinkFunction是一個(gè)接口
public interface SinkFunction<IN> extends Function, Serializable {
//公共方法
default void invoke(IN value, Context context) throws Exception {
invoke(value);
}
}
RichSinkFunction
@Public
public abstract class RichSinkFunction<IN> extends AbstractRichFunction implements SinkFunction<IN> {
private static final long serialVersionUID = 1L;
}
其他繼承接口SinkFunction的類:
案例
自定義HbaseSink
public class HbaseSink extends RichSinkFunction<Tuple2<Integer, String>> {
Logger logger = LoggerFactory.getLogger(HbaseSink.class);
org.apache.hadoop.conf.Configuration configuration;
Connection connection;
@Override
public void open(Configuration parameters) throws Exception {
super.open(parameters);
//獲取hbase 的鏈接信息
configuration = HBaseConfiguration.create();
configuration.set("hbase.zookeeper.quorum", "hadoop101,hadoop102,hadoop103");
//創(chuàng)建conn
connection = ConnectionFactory.createConnection(configuration);
logger.info("創(chuàng)建鏈接成功");
}
@Override
public void invoke(Tuple2<Integer, String> value, Context context) throws Exception {
//往habse 里面插入數(shù)據(jù)
SimpleDateFormat format = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
Table table = connection.getTable(TableName.valueOf("torder_count"));
Put put = new Put(value.f1.getBytes(StandardCharsets.UTF_8));
put.addColumn("info".getBytes(), // 列族
"order_total".getBytes(StandardCharsets.UTF_8), //特征字段
value.f0.toString().getBytes()); //屬性值
put.addColumn("info".getBytes(), "insert_time".getBytes(), format.format(new Date(System.currentTimeMillis())).getBytes());
table.put(put);
table.close();
logger.info("=====一條數(shù)據(jù)寫入成功======,時(shí)間:"+value.f1+", 值:"+value.f0);
}
@Override
public void close() throws Exception {
super.close();
connection.close();
}
通過以上案例我們熟悉了addSink函數(shù)的操作。
到此這篇關(guān)于Flink自定義Sink端實(shí)現(xiàn)過程講解的文章就介紹到這了,更多相關(guān)Flink自定義Sink內(nèi)容請搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!
相關(guān)文章
SpringMVC框架和SpringBoot項(xiàng)目中控制器的響應(yīng)結(jié)果深入分析
這篇文章主要介紹了SpringMVC框架和SpringBoot項(xiàng)目中控制器的響應(yīng)結(jié)果,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)吧2022-12-12
Java實(shí)現(xiàn)向Word文檔添加文檔屬性
這篇文章主要介紹了Java實(shí)現(xiàn)向Word文檔添加文檔屬性的相關(guān)資料,需要的朋友可以參考下2023-01-01
JAVA實(shí)現(xiàn)基于皮爾遜相關(guān)系數(shù)的相似度詳解
這篇文章主要介紹了JAVA實(shí)現(xiàn)基于皮爾遜相關(guān)系數(shù)的相似度詳解,具有一定參考價(jià)值,需要的朋友可以了解下。2017-11-11
SpringBoot3整合EasyExcel動(dòng)態(tài)實(shí)現(xiàn)表頭重命名
這篇文章主要為大家詳細(xì)介紹了SpringBoot3整合EasyExcel如何通過WriteHandler動(dòng)態(tài)實(shí)現(xiàn)表頭重命名,文中的示例代碼講解詳細(xì),有需要的可以了解下2025-03-03
java使用CKEditor實(shí)現(xiàn)圖片上傳功能
這篇文章主要為大家詳細(xì)介紹了java使用CKEditor實(shí)現(xiàn)圖片上傳功能,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下2017-07-07
springbootAOP定義切點(diǎn)獲取/修改請求參數(shù)方式
這篇文章主要介紹了springbootAOP定義切點(diǎn)獲取/修改請求參數(shù)方式,具有很好的參考價(jià)值,希望對大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2024-08-08
spring?cloud中Feign導(dǎo)入jar失敗的問題及解決方案
這篇文章主要介紹了spring?cloud中Feign導(dǎo)入jar失敗的問題及解決方案,具有很好的參考價(jià)值,希望對大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教2022-03-03

