欢迎光临
我们一直在努力

Flink 写入 HBase 终极实战:从同步单条到高性能批量写入与幂等性设计

Flink 写入 HBase 终极实战:从同步单条到高性能批量写入与幂等性设计

在这里插入图片描述

1. 引言:为什么实时数据需要存入 HBase?

在实时计算链路中,Flink 处理完窗口聚合、事件过滤或数据增强后,结果通常需要落地到在线存储系统,供下游 API 查询、报表生成或实时推荐。HBase 作为一款分布式的、面向列的 NoSQL 数据库,以其海量存储能力(可轻松扩展至数十亿行)和稳定的毫秒级随机读延迟,成为了实时结果存储的优选方案之一。

然而,在生产环境中,许多开发者把 Flink 与 HBase 整合时屡屡踩坑:

  • 性能瓶颈:每条数据都 table.put() 然后 close(),导致吞吐量卡在每秒几百条,完全无法应对真实流量。
  • 数据热点:RowKey 使用时间戳或随机递增数字,导致写入压力全部集中在某一个 Region,集群利用率严重不均。
  • 数据重复:当 Flink 发生故障恢复时,由于 Checkpoint 回滚,可能导致部分数据重复写入,而 HBase 本身不支持事务,如何避免数据冗余?

本文将从最小可用 Demo 起步,逐步演进到生产级高性能 Sink,并深入剖析 HBase 写入原理、RowKey 防热点设计、以及 Flink Checkpoint 与 HBase 版本号的联动去重策略。读完本文,你将收获:

  • 一个可直接运行的 Flink + HBase 项目,包含完整的 docker-compose 环境。
  • 三种写入方案的性能对比与选型决策树。
  • 应对数据重复的幂等性设计方法论。

2. 前置知识:HBase 写入路径与 Flink Sink 模型

2.1 HBase 写入的底层旅程

理解 HBase 写入过程,是进行性能调优的前提。一条 Put 请求从客户端到持久化,大致经历以下路径:

阶段组件说明
1. 客户端缓存 本地 Write Buffer Put 先进入客户端的写缓冲区(默认 2MB),而非立即发送。
2. 批量发送 RPC 通信 缓冲区满或手动 flush 时,将一批 Put 通过 RPC 发送给对应的 RegionServer。
3. WAL 预写日志 RegionServer 数据先顺序写入 HDFS 上的 WAL(Write-Ahead Log),保证故障恢复。
4. MemStore 内存 RegionServer 数据写入内存中的 MemStore,提供快速访问。
5. Flush 到 HFile RegionServer MemStore 达到阈值(默认 128MB)时,异步刷写为 HFile 存储到 HDFS。

核心启示:第 1 步的客户端缓冲和第 2 步的批量发送是提升吞吐的关键。如果每条数据都立即触发一次 RPC,吞吐量将受到网络 RTT 的严重限制。

2.2 Flink Sink 的两种模式

Flink 支持两种 Sink 接口:

接口生命周期管理适用场景
SinkFunction 无 open/close,每条数据独立处理 极简输出,无状态资源
RichSinkFunction 支持 open(初始化)、close(释放资源)、invoke(处理记录) 需要连接池、批量缓冲、事务管理的场景

我们将全程基于 RichSinkFunction,在 open 中建立 HBase 连接,close 中清理资源,避免每条数据重复创建连接。


3. 核心剖析:三种写入方案深度对比

本节展示三种实现,并给出典型测试环境(单并行度,千兆网络,HBase 2.4,Flink 1.13)下的相对性能参考值。

3.1 方案一:同步单条写入(纯教学,禁止生产)

直接获取 Table,每次一条 put,立即关闭。此方式完全无法利用客户端缓冲。

// 仅示意,不展开完整代码
val table = conn.getTable(tableName)
val put = new Put(rowKey)
table.put(put)
table.close()

性能估算:≈ 800 ~ 1200 TPS(受限于 RPC 往返延迟)。

3.2 方案二:复用 Table + 手动批量提交

在 open 中获取一次 Table,invoke 中将 Put 攒到 List,达到阈值后调用 table.put(List<Put>)。

class BatchTableSink extends RichSinkFunction[Event] {
var table: Table = _
val batchSize = 100
val puts = new java.util.ArrayList[Put]()

override def open(parameters: Configuration): Unit = {
val conn = ConnectionFactory.createConnection(hbaseConf)
table = conn.getTable(TableName.valueOf("user_events"))
}

override def invoke(value: Event, context: SinkFunction.Context): Unit = {
val put = new Put(rowKey)
put.addColumn(...)
puts.add(put)
if (puts.size() >= batchSize) {
table.put(puts) // 批量提交
puts.clear()
}
}

override def close(): Unit = {
if (!puts.isEmpty) table.put(puts)
table.close()
}
}

性能估算:≈ 10000 ~ 14000 TPS(批量减少 RPC 次数,提升明显)。

3.3 方案三:BufferedMutator(生产级,官方推荐)

BufferedMutator 是 HBase 专门为高吞吐写入设计的组件,异步缓存、自动 flush,并提供异常监听。

class BufferedMutatorSink extends RichSinkFunction[Event] {
var mutator: BufferedMutator = _

override def open(parameters: Configuration): Unit = {
val conn = ConnectionFactory.createConnection(hbaseConf)
mutator = conn.getBufferedMutator(TableName.valueOf("user_events"))
mutator.setExceptionListener(new BufferedMutator.ExceptionListener {
override def onException(e: Exception, mutator: BufferedMutator): Unit = {
// 生产环境可在此实现重试逻辑
println(s"Write error: ${e.getMessage}")
}
})
}

override def invoke(value: Event, context: SinkFunction.Context): Unit = {
val put = new Put(rowKey)
put.addColumn(...)
mutator.mutate(put) // 异步缓存,立即返回
}

override def close(): Unit = {
mutator.close() // 自动 flush 剩余数据并释放资源
}
}

性能估算:≈ 15000 ~ 20000+ TPS(异步非阻塞,吞吐最高,且延迟更稳定)。

本节小结:生产环境请无条件选择 BufferedMutator。方案二可作为理解批量机制的过渡,但方案一必须杜绝。


4. 手把手实操:从零搭建可运行项目

4.1 一键启动 HBase 环境(Docker Compose)

为了让读者最快上手,我提供一个 docker-compose.yml,它启动一个包含 ZooKeeper 和 HBase 的单机容器集群(使用官方镜像 apache/hbase,更可靠):

version: '3'
services:
zookeeper:
image: zookeeper:3.8
container_name: zkhbase
ports:
"2181:2181"
environment:
ZOO_MY_ID: 1

hbase-master:
image: apache/hbase:2.5.0
container_name: hbasemaster
ports:
"16010:16010" # Web UI
"16000:16000" # Master RPC
environment:
HBASE_CONF_zk_quorum: zookeeper:2181
HBASE_CONF_hbase_rootdir: /tmp/hbaseroot
depends_on:
zookeeper
volumes:
./hbasedata:/tmp/hbaseroot # 持久化数据(可选)
command: ["/opt/hbase/bin/hbase", "master", "start"]

hbase-regionserver:
image: apache/hbase:2.5.0
container_name: hbasers
ports:
"16020:16020" # RegionServer RPC
environment:
HBASE_CONF_zk_quorum: zookeeper:2181
HBASE_CONF_hbase_rootdir: /tmp/hbaseroot
depends_on:
hbasemaster
command: ["/opt/hbase/bin/hbase", "regionserver", "start"]

使用方法:

docker-compose up -d
# 等待约 30 秒,访问 http://localhost:16010 查看 HBase Web UI

4.2 Maven 依赖(pom.xml)

<properties>
<flink.version>1.13.6</flink.version>
<scala.binary.version>2.12</scala.binary.version>
<hbase.version>2.5.0</hbase.version>
</properties>

<dependencies>
<!– Flink –>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-scala_${scala.binary.version}</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-clients_${scala.binary.version}</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>

<!– HBase Client –>
<dependency>
<groupId>org.apache.hbase</groupId>
<artifactId>hbase-client</artifactId>
<version>${hbase.version}</version>
</dependency>
<dependency>
<groupId>org.apache.hbase</groupId>
<artifactId>hbase-common</artifactId>
<version>${hbase.version}</version>
</dependency>

<!– 日志 –>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-log4j12</artifactId>
<version>1.7.32</version>
<scope>runtime</scope>
</dependency>
</dependencies>

4.3 完整主程序(使用 BufferedMutator + 幂等性 RowKey)

package sink

import org.apache.flink.configuration.Configuration
import org.apache.flink.streaming.api.functions.sink.{RichSinkFunction, SinkFunction}
import org.apache.flink.streaming.api.scala._
import org.apache.hadoop.hbase.{HBaseConfiguration, TableName}
import org.apache.hadoop.hbase.client.{BufferedMutator, ConnectionFactory, Put}
import org.apache.hadoop.hbase.util.Bytes
import common.Event // 假设 Event(user, url, timestamp) 已定义

import java.util.UUID

object HBaseBufferedMutatorSinkDemo {
def main(args: Array[String]): Unit = {
val env = StreamExecutionEnvironment.getExecutionEnvironment
// 生产环境务必开启 Checkpoint,用于故障恢复
env.enableCheckpointing(10000L) // 10秒一次
env.setParallelism(1) // 便于观察,生产可调大

val dataStream: DataStream[Event] = env.fromElements(
Event("Mary", "./home", 100L),
Event("Sum", "./cart", 500L),
Event("King", "./prod", 1000L),
Event("King", "./root", 200L)
)

val hbaseSink = new RichSinkFunction[Event] {
private var mutator: BufferedMutator = _

override def open(parameters: Configuration): Unit = {
val conf = HBaseConfiguration.create()
// Docker Compose 中 ZK 的服务名,若本地则用 localhost
conf.set("hbase.zookeeper.quorum", "localhost:2181")
conf.set("hbase.client.write.buffer", "4194304") // 4MB 缓冲区

val conn = ConnectionFactory.createConnection(conf)
mutator = conn.getBufferedMutator(TableName.valueOf("user_events"))
mutator.setExceptionListener(new BufferedMutator.ExceptionListener {
override def onException(e: Exception, mutator: BufferedMutator): Unit = {
// 生产环境可记录日志并触发告警,此处简化为打印
System.err.println(s"BufferedMutator write error: ${e.getMessage}")
// 可选:实现重试,但注意幂等性
}
})
println("HBase BufferedMutator opened.")
}

override def invoke(value: Event, context: SinkFunction.Context): Unit = {
// ====== 关键:幂等性 RowKey 设计 ======
// 使用 "用户_时间戳_UUID" 保证唯一,即使同一个用户同一毫秒重复写入也不会覆盖
// 若想利用 HBase 版本号去重(见下节),可固定 RowKey 并设置不同的时间戳版本
// 此处采用唯一 RowKey,天然幂等(重复写入会产生多行,但业务可接受或下游去重)
val rowKey = s"${value.user}_${System.currentTimeMillis()}_${UUID.randomUUID().toString.take(8)}"
val put = new Put(Bytes.toBytes(rowKey))
put.addColumn(Bytes.toBytes("info"), Bytes.toBytes("url"), Bytes.toBytes(value.url))
put.addColumn(Bytes.toBytes("info"), Bytes.toBytes("ts"), Bytes.toBytes(value.timestamp))

mutator.mutate(put)
}

override def close(): Unit = {
if (mutator != null) {
mutator.close() // 自动 flush
println("HBase BufferedMutator closed.")
}
}
}

dataStream.addSink(hbaseSink)
env.execute("Flink HBase BufferedMutator Demo")
}
}

4.4 建表与运行验证

进入 HBase Shell(可通过 docker exec -it hbase-master hbase shell),创建表:

create 'user_events', 'info'

运行 Flink 主程序,然后扫描表:

scan 'user_events', {LIMIT => 10}

应该看到多条记录,RowKey 各不相同,每个 RowKey 只有一列 info:url。

4.5 常见错误与调试表

错误信息可能原因解决
Connection refused: localhost/127.0.0.1:2181 ZooKeeper 未启动,或 Docker 映射端口不正确 检查 docker ps,确保 ZK 容器运行;若在容器外访问,将 localhost 改为宿主机 IP 或 host.docker.internal(Mac/Win)。
TableNotFoundException: user_events 表未提前创建 在 HBase Shell 中执行 create 命令。
RetriesExhaustedWithDetailsException 写入超时或 Region 不可用 检查 HBase Master 和 RegionServer 是否正常;增大超时参数 hbase.rpc.timeout。
数据写入后扫描为空 缓冲区未 flush 调用 mutator.flush() 强制执行,或等待缓冲满(4MB)。可主动在 invoke 中按条数计数调 flush。

5. 进阶思考:Flink Checkpoint 与 HBase 幂等性设计

5.1 HBase 不支持事务,如何保证 Exactly-Once?

Flink 的端到端 Exactly-Once 需要 Sink 支持事务或幂等性。HBase 不支持分布式事务(如两阶段提交),因此我们无法像 Kafka 那样实现事务性写入。但我们可以通过幂等性设计,在 At-Least-Once 语义下实现数据最终一致(即重复写入不会产生脏数据)。

5.2 两种幂等性实现路径

路径 A:使用唯一 RowKey(已在上文代码中实现)

每次写入的 RowKey 都包含一个随机部分(如 UUID),使得即使同一批数据重算,新 RowKey 也与旧的不同。这会导致数据重复存储(多行相同内容),但不会覆盖旧行。下游查询时可根据业务时间戳取最新一条,或任务层做去重。

优点:实现简单,无冲突。
缺点:存储量增加,且下游需要额外去重逻辑。

路径 B:固定 RowKey + HBase 版本号(列时间戳)

如果业务要求同一个 Key(如 Mary)只保留最新一条记录,则可固定 RowKey,并利用 HBase 单元格支持多版本的特性。每次写入时,手动指定一个递增的版本号(如系统时间戳或单调递增序列),HBase 默认只返回最新版本。这样,即使同一 RowKey 重复写入,只是新增一个版本,查询时只看到最新的,视觉上“幂等”。

// 固定 RowKey = user
val put = new Put(Bytes.toBytes(value.user))
// 显式设置版本号为当前时间戳(毫秒),确保唯一性
put.addColumn(Bytes.toBytes("info"), Bytes.toBytes("url"),
System.currentTimeMillis(), Bytes.toBytes(value.url))
mutator.mutate(put)

此时,如果 Flink 故障恢复导致同一条数据重复写入,由于时间戳相同(或相差很小),HBase 会将它们视为同一版本的多个写入(如果时间戳完全一致,可能会覆盖),实际效果取决于具体时间戳值。更稳健的做法是使用单调递增的批次 ID 作为版本号,并将该 ID 存入 Checkpoint,恢复时从 Checkpoint 读取上次写入的批次 ID。

5.3 结合 Checkpoint 的批次版本号方案(高级)

  • 在 Flink 的 open 方法中从状态恢复一个 AtomicLong batchId。
  • 每次 Checkpoint 完成后,batchId 自增,并将该值作为所有 Put 的版本号。
  • 这样,同一批次(Checkpoint)内的所有数据共享一个版本号,如果该批次因故障重放,版本号相同,HBase 会以相同版本覆盖写入,等同于“幂等”。
  • 该方案虽然复杂,但能真正做到存储层面的无重复。建议对数据准确性极高的场景采用。

    5.4 性能调优补充

    • 缓冲区大小:hbase.client.write.buffer 默认 2MB,可调至 4~8MB 以提升吞吐,但会增加内存占用。
    • 并行度与连接数:每个 Sink Subtask 会持有一个 BufferedMutator,连接数 = 并行度。建议并行度不超过 RegionServer 数量 * 2。
    • 超时参数:生产环境调大 hbase.rpc.timeout(默认 60s)和 hbase.client.retries.number(默认 15),避免网络抖动导致任务失败。

    6. 总结:写入方案决策树与最终建议

    #mermaid-svg-FXXMfjNrkIWa3kWk{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}@keyframes edge-animation-frame{from{stroke-dashoffset:0;}}@keyframes dash{to{stroke-dashoffset:0;}}#mermaid-svg-FXXMfjNrkIWa3kWk .edge-animation-slow{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 50s linear infinite;stroke-linecap:round;}#mermaid-svg-FXXMfjNrkIWa3kWk .edge-animation-fast{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 20s linear infinite;stroke-linecap:round;}#mermaid-svg-FXXMfjNrkIWa3kWk .error-icon{fill:#552222;}#mermaid-svg-FXXMfjNrkIWa3kWk .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-FXXMfjNrkIWa3kWk .edge-thickness-normal{stroke-width:1px;}#mermaid-svg-FXXMfjNrkIWa3kWk .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-FXXMfjNrkIWa3kWk .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-FXXMfjNrkIWa3kWk .edge-thickness-invisible{stroke-width:0;fill:none;}#mermaid-svg-FXXMfjNrkIWa3kWk .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-FXXMfjNrkIWa3kWk .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-FXXMfjNrkIWa3kWk .marker{fill:#333333;stroke:#333333;}#mermaid-svg-FXXMfjNrkIWa3kWk .marker.cross{stroke:#333333;}#mermaid-svg-FXXMfjNrkIWa3kWk svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-FXXMfjNrkIWa3kWk p{margin:0;}#mermaid-svg-FXXMfjNrkIWa3kWk .label{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;color:#333;}#mermaid-svg-FXXMfjNrkIWa3kWk .cluster-label text{fill:#333;}#mermaid-svg-FXXMfjNrkIWa3kWk .cluster-label span{color:#333;}#mermaid-svg-FXXMfjNrkIWa3kWk .cluster-label span p{background-color:transparent;}#mermaid-svg-FXXMfjNrkIWa3kWk .label text,#mermaid-svg-FXXMfjNrkIWa3kWk span{fill:#333;color:#333;}#mermaid-svg-FXXMfjNrkIWa3kWk .node rect,#mermaid-svg-FXXMfjNrkIWa3kWk .node circle,#mermaid-svg-FXXMfjNrkIWa3kWk .node ellipse,#mermaid-svg-FXXMfjNrkIWa3kWk .node polygon,#mermaid-svg-FXXMfjNrkIWa3kWk .node path{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#mermaid-svg-FXXMfjNrkIWa3kWk .rough-node .label text,#mermaid-svg-FXXMfjNrkIWa3kWk .node .label text,#mermaid-svg-FXXMfjNrkIWa3kWk .image-shape .label,#mermaid-svg-FXXMfjNrkIWa3kWk .icon-shape .label{text-anchor:middle;}#mermaid-svg-FXXMfjNrkIWa3kWk .node .katex path{fill:#000;stroke:#000;stroke-width:1px;}#mermaid-svg-FXXMfjNrkIWa3kWk .rough-node .label,#mermaid-svg-FXXMfjNrkIWa3kWk .node .label,#mermaid-svg-FXXMfjNrkIWa3kWk .image-shape .label,#mermaid-svg-FXXMfjNrkIWa3kWk .icon-shape .label{text-align:center;}#mermaid-svg-FXXMfjNrkIWa3kWk .node.clickable{cursor:pointer;}#mermaid-svg-FXXMfjNrkIWa3kWk .root .anchor path{fill:#333333!important;stroke-width:0;stroke:#333333;}#mermaid-svg-FXXMfjNrkIWa3kWk .arrowheadPath{fill:#333333;}#mermaid-svg-FXXMfjNrkIWa3kWk .edgePath .path{stroke:#333333;stroke-width:2.0px;}#mermaid-svg-FXXMfjNrkIWa3kWk .flowchart-link{stroke:#333333;fill:none;}#mermaid-svg-FXXMfjNrkIWa3kWk .edgeLabel{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-FXXMfjNrkIWa3kWk .edgeLabel p{background-color:rgba(232,232,232, 0.8);}#mermaid-svg-FXXMfjNrkIWa3kWk .edgeLabel rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-FXXMfjNrkIWa3kWk .labelBkg{background-color:rgba(232, 232, 232, 0.5);}#mermaid-svg-FXXMfjNrkIWa3kWk .cluster rect{fill:#ffffde;stroke:#aaaa33;stroke-width:1px;}#mermaid-svg-FXXMfjNrkIWa3kWk .cluster text{fill:#333;}#mermaid-svg-FXXMfjNrkIWa3kWk .cluster span{color:#333;}#mermaid-svg-FXXMfjNrkIWa3kWk div.mermaidTooltip{position:absolute;text-align:center;max-width:200px;padding:2px;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:12px;background:hsl(80, 100%, 96.2745098039%);border:1px solid #aaaa33;border-radius:2px;pointer-events:none;z-index:100;}#mermaid-svg-FXXMfjNrkIWa3kWk .flowchartTitleText{text-anchor:middle;font-size:18px;fill:#333;}#mermaid-svg-FXXMfjNrkIWa3kWk rect.text{fill:none;stroke-width:0;}#mermaid-svg-FXXMfjNrkIWa3kWk .icon-shape,#mermaid-svg-FXXMfjNrkIWa3kWk .image-shape{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-FXXMfjNrkIWa3kWk .icon-shape p,#mermaid-svg-FXXMfjNrkIWa3kWk .image-shape p{background-color:rgba(232,232,232, 0.8);padding:2px;}#mermaid-svg-FXXMfjNrkIWa3kWk .icon-shape .label rect,#mermaid-svg-FXXMfjNrkIWa3kWk .image-shape .label rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-FXXMfjNrkIWa3kWk .label-icon{display:inline-block;height:1em;overflow:visible;vertical-align:-0.125em;}#mermaid-svg-FXXMfjNrkIWa3kWk .node .label-icon path{fill:currentColor;stroke:revert;stroke-width:revert;}#mermaid-svg-FXXMfjNrkIWa3kWk :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}

    否/学习测试

    可容忍少量重复

    完全不可重复

    需要将 Flink 数据写入 HBase

    是否生产环境?

    使用 RichSinkFunction + 手动批量或单条

    使用 BufferedMutator

    数据重复可否容忍?

    唯一 RowKey + 下游去重

    固定 RowKey + 版本号 + Checkpoint 批次ID

    最终选型建议

    场景推荐方案原因
    学习/快速验证 方案一或二,单条或手动批量 代码简单,便于理解 API。
    生产环境(高吞吐,容忍少量重复) BufferedMutator + 唯一 RowKey 吞吐最高,实现简单,下游根据时间戳去重。
    生产环境(财务/计费,零重复) BufferedMutator + 固定 RowKey + 版本号(批次ID) 需结合 Flink 状态,实现稍复杂,但保证数据精确一次。

    最后,三条黄金法则:

  • 绝不用单条同步写入 —— 那是在浪费集群资源。
  • RowKey 必须防热点 —— 加盐、哈希、反转,至少选一种。
  • 监控 HBase 写入延迟和异常 —— 设置告警,避免数据积压。
  • 现在,你已经拥有了从入门到生产的全套技能,开始动手打造你的高性能 Flink-HBase 数据管道吧!如果实践过程中遇到新问题,欢迎在评论区探讨。


    附录:快速测试性能的简易方法

    在 invoke 中加一个计数器,每 10000 条打印一次耗时,即可粗略估算 TPS。生产环境建议使用 JMH 或 Flink 自带的 Metrics 系统进行准确测量。

    赞(0)
    未经允许不得转载:171主机测评 » Flink 写入 HBase 终极实战:从同步单条到高性能批量写入与幂等性设计
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址