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数据存储至数据盘

查看磁盘及分区 [root@MySQL tmp]# fdisk -lDisk /dev/sda: 21.5 GB, 21474836480 bytes255 heads, 63 sectors/track, 2610 cylindersUnits = cylinders of 16065 * 512 = 8225280 bytesSector size (logical/physical)

SQL Server中,查询数据库中有多少个表,以及数据库其余类型数据统计查询

sqlserver查询数据库中有多少个表 sql server 数表:select count(1) from sysobjects where xtype='U'数视图:select count(1) from sysobjects where xtype='V'数存储过程select count(1) from sysobjects where xtype='P' SE

数据时代的数字企业

1.写在前面 讨论数据治理在数字企业中的影响和必要性,并介绍数据治理的核心内容和实践方法。作者强调了数据质量、数据安全、数据隐私和数据合规等方面是数据治理的核心内容,并介绍了具体的实践措施和案例分析。企业需要重视这些方面以实现数字化转型和业务增长。 数字化转型行业小伙伴可以加入我的星球,初衷成为各位数字化转型参考库,星球内容每周更新 个人工作经验资料全部放在这里,包含数据治理、数据要

如何在Java中处理JSON数据?

如何在Java中处理JSON数据? 大家好,我是免费搭建查券返利机器人省钱赚佣金就用微赚淘客系统3.0的小编,也是冬天不穿秋裤,天冷也要风度的程序猿!今天我们将探讨在Java中如何处理JSON数据。JSON(JavaScript Object Notation)作为一种轻量级的数据交换格式,在现代应用程序中被广泛使用。Java通过多种库和API提供了处理JSON的能力,我们将深入了解其用法和最佳

两个基因相关性CPTAC蛋白组数据

目录 蛋白数据下载 ①蛋白数据下载 1,TCGA-选择泛癌数据  2,TCGA-TCPA 3,CPTAC(非TCGA) ②蛋白相关性分析 1,数据整理 2,蛋白相关性分析 PCAS在线分析 蛋白数据下载 CPTAC蛋白组学数据库介绍及数据下载分析 – 王进的个人网站 (jingege.wang) ①蛋白数据下载 可以下载泛癌蛋白数据:UCSC Xena (xena

53、Flink Interval Join 代码示例

1、概述 interval Join 默认会根据 keyBy 的条件进行 Join 此时为 Inner Join; interval Join 算子的水位线会取两条流中水位线的最小值; interval Join 迟到数据的判定是以 interval Join 算子的水位线为基准; interval Join 可以分别输出两条流中迟到的数据-[sideOutputLeftLateData,

中国341城市生态系统服务价值数据集(2000-2020年)

生态系统服务反映了人类直接或者间接从自然生态系统中获得的各种惠益,对支撑和维持人类生存和福祉起着重要基础作用。目前针对全国城市尺度的生态系统服务价值的长期评估还相对较少。我们在Xie等(2017)的静态生态系统服务当量因子表基础上,选取净初级生产力,降水量,生物迁移阻力,土壤侵蚀度和道路密度五个变量,对生态系统供给服务、调节服务、支持服务和文化服务共4大类和11小类的当量因子进行了时空调整,计算了

【计算机网络篇】数据链路层(12)交换机式以太网___以太网交换机

文章目录 🍔交换式以太网🛸以太网交换机 🍔交换式以太网 仅使用交换机(不使用集线器)的以太网就是交换式以太网 🛸以太网交换机 以太网交换机本质上就是一个多接口的网桥: 交换机的每个接口考研连接计算机,也可以理解集线器或另一个交换机 当交换机的接口与计算机或交换机连接时,可以工作在全双工方式,并能在自身内部同时连通多对接口,使每一对相互通信的计算机都能像

使用Jsoup抓取数据

问题 最近公司的市场部分布了一个问题,到一个网站截取一下医院的数据。刚好我也被安排做。后来,我发现为何不用脚本去抓取呢? 抓取的数据如下: Jsoup的使用实战代码 结构 Created with Raphaël 2.1.0 开始 创建线程池 jsoup读取网页 解析Element 写入sqlite 结束