Flink写出数据到Hbase的Sink

2024-05-10 13:32
文章标签 数据 flink sink hbase 写出

本文主要是介绍Flink写出数据到Hbase的Sink,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

文章目录

    • 一、MyHbaseSink
        • 1、继承RichSinkFunction<输入的数据类型>类
        • 2、实现open方法,创建连接对象
        • 3、实现invoke方法,批次写入数据到Hbase
        • 4、实现close方法,关闭连接
    • 二、HBaseUtil工具类

一、MyHbaseSink

1、继承RichSinkFunction<输入的数据类型>类
public class MyHbaseSink extends RichSinkFunction<Tuple2<String, Double>> {private transient Integer maxSize = 1000;private transient Long delayTime = 5000L;public MyHbaseSink() {}public MyHbaseSink(Integer maxSize, Long delayTime) {this.maxSize = maxSize;this.delayTime = delayTime;}private transient Connection connection;private transient Long lastInvokeTime;private transient List<Put> puts = new ArrayList<>(maxSize);
2、实现open方法,创建连接对象
    // 创建连接@Overridepublic void open(Configuration parameters) throws Exception {super.open(parameters);// 获取全局配置文件,并转为ParameterToolParameterTool params =(ParameterTool) getRuntimeContext().getExecutionConfig().getGlobalJobParameters();//创建一个Hbase的连接connection = HBaseUtil.getConnection(params.getRequired("hbase.zookeeper.quorum"),params.getInt("hbase.zookeeper.property.clientPort", 2181));// 获取系统当前时间lastInvokeTime = System.currentTimeMillis();}
3、实现invoke方法,批次写入数据到Hbase
    @Overridepublic void invoke(Tuple2<String, Double> value, Context context) throws Exception {String rk = value.f0;//创建put对象,并赋rk值Put put = new Put(rk.getBytes());// 添加值:f1->列族, order->属性名 如age, 第三个->属性值 如25put.addColumn("f1".getBytes(), "order".getBytes(), value.f1.toString().getBytes());puts.add(put);// 添加put对象到list集合//使用ProcessingTimelong currentTime = System.currentTimeMillis();//开始批次提交数据if (puts.size() == maxSize || currentTime - lastInvokeTime >= delayTime) {//获取一个Hbase表Table table = connection.getTable(TableName.valueOf("database:table"));table.put(puts);//批次提交puts.clear();lastInvokeTime = currentTime;table.close();}}
4、实现close方法,关闭连接
    @Overridepublic void close() throws Exception {connection.close();}

二、HBaseUtil工具类

  • Hbase的工具类,用来创建Hbase的Connection
public class HBaseUtil {/*** @param zkQuorum zookeeper地址,多个要用逗号分隔* @param port     zookeeper端口号* @return connection*/public static Connection getConnection(String zkQuorum, int port) throws Exception {Configuration conf = HBaseConfiguration.create();conf.set("hbase.zookeeper.quorum", zkQuorum);conf.set("hbase.zookeeper.property.clientPort", port + "");Connection connection = ConnectionFactory.createConnection(conf);return connection;}
}

这篇关于Flink写出数据到Hbase的Sink的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/976574

相关文章

Mysql如何将数据按照年月分组的统计

《Mysql如何将数据按照年月分组的统计》:本文主要介绍Mysql如何将数据按照年月分组的统计方式,具有很好的参考价值,希望对大家有所帮助,如有错误或未考虑完全的地方,望不吝赐教... 目录mysql将数据按照年月分组的统计要的效果方案总结Mysql将数据按照年月分组的统计要的效果方案① 使用 DA

鸿蒙中Axios数据请求的封装和配置方法

《鸿蒙中Axios数据请求的封装和配置方法》:本文主要介绍鸿蒙中Axios数据请求的封装和配置方法,本文给大家介绍的非常详细,对大家的学习或工作具有一定的参考借鉴价值,需要的朋友参考下吧... 目录1.配置权限 应用级权限和系统级权限2.配置网络请求的代码3.下载在Entry中 下载AxIOS4.封装Htt

Python获取中国节假日数据记录入JSON文件

《Python获取中国节假日数据记录入JSON文件》项目系统内置的日历应用为了提升用户体验,特别设置了在调休日期显示“休”的UI图标功能,那么问题是这些调休数据从哪里来呢?我尝试一种更为智能的方法:P... 目录节假日数据获取存入jsON文件节假日数据读取封装完整代码项目系统内置的日历应用为了提升用户体验,

Java利用JSONPath操作JSON数据的技术指南

《Java利用JSONPath操作JSON数据的技术指南》JSONPath是一种强大的工具,用于查询和操作JSON数据,类似于SQL的语法,它为处理复杂的JSON数据结构提供了简单且高效... 目录1、简述2、什么是 jsONPath?3、Java 示例3.1 基本查询3.2 过滤查询3.3 递归搜索3.4

MySQL大表数据的分区与分库分表的实现

《MySQL大表数据的分区与分库分表的实现》数据库的分区和分库分表是两种常用的技术方案,本文主要介绍了MySQL大表数据的分区与分库分表的实现,文中通过示例代码介绍的非常详细,对大家的学习或者工作具有... 目录1. mysql大表数据的分区1.1 什么是分区?1.2 分区的类型1.3 分区的优点1.4 分

Mysql删除几亿条数据表中的部分数据的方法实现

《Mysql删除几亿条数据表中的部分数据的方法实现》在MySQL中删除一个大表中的数据时,需要特别注意操作的性能和对系统的影响,本文主要介绍了Mysql删除几亿条数据表中的部分数据的方法实现,具有一定... 目录1、需求2、方案1. 使用 DELETE 语句分批删除2. 使用 INPLACE ALTER T

Python Dash框架在数据可视化仪表板中的应用与实践记录

《PythonDash框架在数据可视化仪表板中的应用与实践记录》Python的PlotlyDash库提供了一种简便且强大的方式来构建和展示互动式数据仪表板,本篇文章将深入探讨如何使用Dash设计一... 目录python Dash框架在数据可视化仪表板中的应用与实践1. 什么是Plotly Dash?1.1

Redis 中的热点键和数据倾斜示例详解

《Redis中的热点键和数据倾斜示例详解》热点键是指在Redis中被频繁访问的特定键,这些键由于其高访问频率,可能导致Redis服务器的性能问题,尤其是在高并发场景下,本文给大家介绍Redis中的热... 目录Redis 中的热点键和数据倾斜热点键(Hot Key)定义特点应对策略示例数据倾斜(Data S

Python实现将MySQL中所有表的数据都导出为CSV文件并压缩

《Python实现将MySQL中所有表的数据都导出为CSV文件并压缩》这篇文章主要为大家详细介绍了如何使用Python将MySQL数据库中所有表的数据都导出为CSV文件到一个目录,并压缩为zip文件到... python将mysql数据库中所有表的数据都导出为CSV文件到一个目录,并压缩为zip文件到另一个

SpringBoot整合jasypt实现重要数据加密

《SpringBoot整合jasypt实现重要数据加密》Jasypt是一个专注于简化Java加密操作的开源工具,:本文主要介绍详细介绍了如何使用jasypt实现重要数据加密,感兴趣的小伙伴可... 目录jasypt简介 jasypt的优点SpringBoot使用jasypt创建mapper接口配置文件加密