- 数据工程
- 大数据
- 批处理
- 流处理
【免费下载链接】seatunnel
SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.
本文围绕 SeaTunnel 官方文档 Phoenix Sink Connector 展开,介绍如何借助
Jdbc连接器将上游数据写入 Apache Phoenix(底层为 HBase),覆盖厚/薄两种 JDBC 驱动连接方式、批量与流式写入、配置参数详解、完整可运行的示例配置,并结合仓库中的源码与端到端测试给出实现层面的印证。读完本文,你将能够独立完成 Phoenix Sink 任务的配置、驱动选型与常见问题排查。
一、Phoenix Sink 是什么
Apache Phoenix 是一个运行在 HBase 之上的 SQL 层,通过 JDBC 驱动将 SQL 请求翻译为 HBase 的读写操作。SeaTunnel 的 Phoenix Sink Connector 本身不是一个独立插件,而是Jdbc 连接器在 Phoenix 方言(Dialect)下的一个落地场景:在配置中声明 Phoenix 的 driver 与 url,Jdbc 连接器便以 Phoenix 方言执行写入。
其核心机制如下:
- 支持Batch 模式与Streaming 模式(与 Jdbc 连接器能力一致);
- 官方文档标注的已测 Phoenix 版本为4.xx 与 5.xx;
- 底层通过 Phoenix 的 JDBC 驱动执行UPSERT 语句(如
UPSERT INTO ... VALUES(?, ?))写入 HBase; - 提供两种 Java JDBC 连接方式:厚驱动(thick)经 ZooKeeper 连接,薄驱动(thin client)经 QueryServer 连接。
graph TD A[SeaTunnel 上游数据] --> B[Jdbc Sink 连接器] B --> C[PhoenixDialect 方言] C --> D[Phoenix JDBC 驱动 thick/thin] D --> E[Phoenix QueryServer / ZooKeeper] E --> F[HBase]关键注意点
- 默认使用薄驱动(thin)的 jar。若需使用厚驱动或其它版本的薄驱动,需要重新编译 jdbc 连接器模块(详见下文“驱动依赖”小节);
- 不支持 exactly-once 语义:Phoenix 尚未支持 XA 事务,因此无法像 MySQL 等数据库那样通过
is_exactly_once=true获得精确一次投递(见 Jdbc 连接器文档 中关于 XA 事务的说明)。
二、两种 JDBC 连接方式对比
Phoenix 官方文档提供了两种连接方式,SeaTunnel 的 Phoenix Sink 同时支持二者,区别在于 driver 类名与 url 格式。
| 连接方式 | driver 值 | url 格式示例 | 适用场景 |
|---|---|---|---|
| 厚驱动(thick client) | org.apache.phoenix.jdbc.PhoenixDriver | jdbc:phoenix:localhost:2182/hbase | 直接通过 ZooKeeper 连接 HBase 集群,驱动需要打入完整的 Phoenix/HBase 客户端依赖 |
| 薄驱动(thin client) | org.apache.phoenix.queryserver.client.Driver | jdbc:phoenix:thin:url=http://localhost:8765;serialization=PROTOBUF | 通过 Phoenix QueryServer(HTTP 服务,默认端口 8765)代理请求,客户端依赖轻量 |
从仓库源码可以印证这一连接策略:Phoenix 方言工厂通过acceptsURL()判断 URL 前缀是否为jdbc:phoenix:来匹配 Phoenix 方言,见 PhoenixDialectFactory.java;方言本体则负责提供行数据转换器与类型映射,见 PhoenixDialect.java。
其中serialization=PROTOBUF是薄客户端与 QueryServer 之间的序列化协议参数;url参数中的主机名在容器化或分布式部署中应填写 QueryServer 所在节点的主机名或服务别名(如测试环境中的seatunnel_e2e_phoenix)。
三、驱动依赖(jar 的准备)
Phoenix Sink 属于 Jdbc 连接器家族,驱动 jar 的放置规则与通用 Jdbc 一致:
- Spark/Flink 引擎:将 Phoenix JDBC 驱动 jar 放入
${SEATUNNEL_HOME}/plugins/目录; - SeaTunnel Zeta 引擎:将驱动 jar 放入
${SEATUNNEL_HOME}/lib/目录。
提示:官方文档明确说明,默认使用薄驱动 jar(例如阿里云 Phoenix 的
ali-phoenix-shaded-thin-client,见 Jdbc.md 附录表 中 Phoenix 一行的 maven 依赖说明)。如果想改用厚驱动(org.apache.phoenix.jdbc.PhoenixDriver)或其它版本的薄驱动,必须重新编译 jdbc 连接器模块(即seatunnel-connectors-v2/connector-jdbc模块),把对应的驱动依赖打入连接器。
四、Sink 配置参数详解
driver [string]
- 厚驱动:
org.apache.phoenix.jdbc.PhoenixDriver - 薄驱动:
org.apache.phoenix.queryserver.client.Driver
该值用于驱动类的加载,必须与所放置的 jar 匹配。
url [string]
- 厚驱动:
jdbc:phoenix:localhost:2182/hbase(2182为 ZooKeeper 端口,/hbase为 HBase 在 ZK 上的根节点,按实际集群调整) - 薄驱动:
jdbc:phoenix:thin:url=http://localhost:8765;serialization=PROTOBUF(8765为 QueryServer 默认 HTTP 端口)
query [string]
写入 SQL,使用?占位符接收上游字段,例如:
upsert into test.sink(age, name) values(?, ?)扩展说明:Jdbc 连接器还支持
database/table+generate_sink_sql自动生成 SQL、primary_keys、schema_save_mode/data_save_mode等高级能力(详见 Jdbc.md 参数表)。针对 Phoenix 场景,最直接、最可控的写法仍是显式指定query的 UPSERT 语句。
通用参数(common options)
Sink 插件的通用参数(如source_table_name、result_table_name的搭配规则)请参考 Sink Common Options:当任务中 source/transform/sink 任一环节存在多个实例时,需要通过result_table_name与source_table_name显式串接数据流。
五、完整配置示例
以下示例来自官方文档并补充了完整结构,展示厚/薄两种驱动下 Phoenix Sink 的写法(均配合Jdbc插件名使用)。
示例一:厚驱动(thick client)
env { parallelism = 1 job.mode = "BATCH" } source { Jdbc { driver = org.apache.phoenix.jdbc.PhoenixDriver url = "jdbc:phoenix:localhost:2182/hbase" query = "select age, name from test.source" } } transform { } sink { Jdbc { driver = org.apache.phoenix.jdbc.PhoenixDriver url = "jdbc:phoenix:localhost:2182/hbase" query = "upsert into test.sink(age, name) values(?, ?)" } }示例二:薄驱动(thin client)
env { parallelism = 1 job.mode = "BATCH" } source { Jdbc { driver = org.apache.phoenix.queryserver.client.Driver url = "jdbc:phoenix:thin:url=http://spark_e2e_phoenix_sink:8765;serialization=PROTOBUF" query = "select age, name from test.source" } } transform { } sink { Jdbc { driver = org.apache.phoenix.queryserver.client.Driver url = "jdbc:phoenix:thin:url=http://spark_e2e_phoenix_sink:8765;serialization=PROTOBUF" query = "upsert into test.sink(age, name) values(?, ?)" } }两个示例中的 URL 主机名均为示例值(
localhost:2182、spark_e2e_phoenix_sink:8765),实际使用时请替换为你的 ZooKeeper / QueryServer 地址。
模式说明:Batch 与 Streaming
- Batch 模式:设置
job.mode = "BATCH",适合定时批量的 UPSERT 回填; - Streaming 模式:设置
job.mode = "STREAMING"并配合checkpoint.interval,持续消费上游实时数据并写入 HBase(受 Phoenix 事务能力限制,流式场景下同样只保证 at-least-once 级别的投递,请结合业务做幂等设计——Phoenix 的 UPSERT 本身以主键为准,天然具备覆盖语义,可在一定程度上缓解重复写入问题)。
六、源码与测试印证:Phoenix 方言在 Jdbc 连接器中的落地
Phoenix 并非独立连接器模块,而是 Jdbc 连接器内部的方言实现,相关代码位于seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/phoenix/目录:
- PhoenixDialectFactory.java:通过
@AutoService注册为JdbcDialectFactory,以jdbc:phoenix:前缀识别 Phoenix 连接,这也解释了为什么配置中 driver/url 必须符合 Phoenix 约定; - PhoenixDialect.java:提供行转换器与类型映射器;
- PhoenixJdbcRowConverter.java:继承
AbstractJdbcRowConverter,负责 SeaTunnel 行数据与 JDBC 参数之间的类型互转; - PhoenixTypeMapper.java:将
ResultSetMetaData中的列信息(类型、精度、scale、可空性)映射为 SeaTunnel 的列定义。
仓库同时提供了完整的端到端验证:测试类 JdbcPhoenixIT.java 使用iteblog/hbase-phoenix-docker:1.0容器启动带 QueryServer 的 Phoenix 环境(端口 8765),建表test.SOURCE/test.SINK(age INTEGER PRIMARY KEY, name VARCHAR(255)),向源表写入 100 行测试数据,随后通过 jdbc_phoenix_source_and_sink.conf 配置完成“Jdbc 源 → Phoenix 薄驱动 UPSERT 写 SINK 表”的整链路验证:
source { Jdbc { driver = org.apache.phoenix.queryserver.client.Driver url = "jdbc:phoenix:thin:url=http://seatunnel_e2e_phoenix:8765;serialization=PROTOBUF" query = "select * from test.SOURCE" } } sink { Jdbc { driver = org.apache.phoenix.queryserver.client.Driver url = "jdbc:phoenix:thin:url=http://seatunnel_e2e_phoenix:8765;serialization=PROTOBUF" query = "upsert into test.SINK(age, name) values(?, ?)" } }测试中还演示了 Phoenix 特有的两点:
- 建表语法:
CREATE TABLE ... (age INTEGER PRIMARY KEY, name VARCHAR(255)),主键列直接写在列定义中,这是 Phoenix 将主键映射为 HBase RowKey 的方式; - 清表方式:Phoenix 不支持
TRUNCATE,测试中使用delete from ... where 1=1完成清空(见 JdbcPhoenixIT.java),实际运维中清空 Phoenix 表数据时同样需要借助 DELETE 语句。
七、常见问题与排查建议
| 现象 | 可能原因 | 处理建议 |
|---|---|---|
| 报 ClassNotFound / 驱动类加载失败 | 驱动 jar 未放入plugins/或lib/,或 driver 类名与 jar 不匹配 | 按引擎类型放置 jar,并核对 driver 取值 |
| 薄驱动连接失败 | QueryServer 未启动、端口不对、主机名无法解析 | 确认 8765 端口服务可用,URL 中主机名改为 QueryServer 可达地址 |
| 厚驱动连接失败 | ZooKeeper 地址/端口或 HBase 根节点错误 | 核对jdbc:phoenix:<zk_host>:<port>/<hbase_rootnode>各段 |
| 需要 exactly-once 语义 | Phoenix 暂不支持 XA 事务 | 关闭is_exactly_once,利用 UPSERT 主键覆盖语义自行保证最终一致 |
| 更换驱动版本失败 | 驱动依赖未随连接器重新编译 | 按文档说明重新编译connector-jdbc模块并打入对应驱动依赖 |
八、变更记录
- 2.2.0-beta(2022-09-26):新增 Phoenix Sink Connector。
延伸阅读
- Jdbc Sink Connector(Phoenix 的载体与完整参数表)
- Jdbc Source Connector(Phoenix 数据读取)
- Sink 通用参数说明
- 连接器能力特性总览(batch / stream / exactly-once 等)
- 数据工程
- 大数据
- 批处理
- 流处理
【免费下载链接】seatunnel
SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.
相关推荐
SeaTunnel Phoenix Sink 实战指南:基于 Jdbc 连接器将数据 UPSERT 写入 HBase
SeaTunnel Phoenix Sink 实战指南:基于 Jdbc 连接器将数据 UPSERT 写入 HBase 本指南以 SeaTunnel 仓库中 Ph
数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Phoenix Sink 连接器实战:基于 Jdbc Connector 向 Apache Phoenix/HBase 写入数据
SeaTunnel Phoenix Sink 连接器实战:基于 Jdbc Connector 向 Apache Phoenix/HBase 写入数据 本文基于
数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Vertica JDBC Sink Connector 实战指南:从依赖配置到 MERGE Upsert 写入
SeaTunnel Vertica JDBC Sink Connector 实战指南:从依赖配置到 MERGE Upsert 写入 Vertica 作为一款面向
数据集成ETL大数据批处理流处理变更数据捕获
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考