免费获取学习方案
ARTICLE DETAIL

资讯详情

深耕编程基础知识与建站技术分享的一线实战洞察。

SeaTunnel Phoenix Sink Connector 实战指南:基于 JDBC 实现 HBase 高性能 UPSERT 写入

SeaTunnel Phoenix Sink Connector 实战指南:基于 JDBC 实现 HBase 高性能 UPSERT 写入 数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载本文围绕 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 与 urlJdbc 连接器便以 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_oncetrue获得精确一次投递见 Jdbc 连接器文档 中关于 XA 事务的说明。二、两种 JDBC 连接方式对比Phoenix 官方文档提供了两种连接方式SeaTunnel 的 Phoenix Sink 同时支持二者区别在于 driver 类名与 url 格式。连接方式driver 值url 格式示例适用场景厚驱动thick clientorg.apache.phoenix.jdbc.PhoenixDriverjdbc:phoenix:localhost:2182/hbase直接通过 ZooKeeper 连接 HBase 集群驱动需要打入完整的 Phoenix/HBase 客户端依赖薄驱动thin clientorg.apache.phoenix.queryserver.client.Driverjdbc:phoenix:thin:urlhttp://localhost:8765;serializationPROTOBUF通过 Phoenix QueryServerHTTP 服务默认端口 8765代理请求客户端依赖轻量从仓库源码可以印证这一连接策略Phoenix 方言工厂通过acceptsURL()判断 URL 前缀是否为jdbc:phoenix:来匹配 Phoenix 方言见 PhoenixDialectFactory.java方言本体则负责提供行数据转换器与类型映射见 PhoenixDialect.java。其中serializationPROTOBUF是薄客户端与 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/hbase2182为 ZooKeeper 端口/hbase为 HBase 在 ZK 上的根节点按实际集群调整薄驱动jdbc:phoenix:thin:urlhttp://localhost:8765;serializationPROTOBUF8765为 QueryServer 默认 HTTP 端口query [string]写入 SQL使用?占位符接收上游字段例如upsert into test.sink(age, name) values(?, ?)扩展说明Jdbc 连接器还支持database/tablegenerate_sink_sql自动生成 SQL、primary_keys、schema_save_mode/data_save_mode等高级能力详见 Jdbc.md 参数表。针对 Phoenix 场景最直接、最可控的写法仍是显式指定query的 UPSERT 语句。通用参数common optionsSink 插件的通用参数如source_table_name、result_table_name的搭配规则请参考 Sink Common Options当任务中 source/transform/sink 任一环节存在多个实例时需要通过result_table_name与source_table_name显式串接数据流。五、完整配置示例以下示例来自官方文档并补充了完整结构展示厚/薄两种驱动下 Phoenix Sink 的写法均配合Jdbc插件名使用。示例一厚驱动thick clientenv { 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 clientenv { parallelism 1 job.mode BATCH } source { Jdbc { driver org.apache.phoenix.queryserver.client.Driver url jdbc:phoenix:thin:urlhttp://spark_e2e_phoenix_sink:8765;serializationPROTOBUF query select age, name from test.source } } transform { } sink { Jdbc { driver org.apache.phoenix.queryserver.client.Driver url jdbc:phoenix:thin:urlhttp://spark_e2e_phoenix_sink:8765;serializationPROTOBUF query upsert into test.sink(age, name) values(?, ?) } }两个示例中的 URL 主机名均为示例值localhost:2182、spark_e2e_phoenix_sink:8765实际使用时请替换为你的 ZooKeeper / QueryServer 地址。模式说明Batch 与 StreamingBatch 模式设置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.SINKage 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:urlhttp://seatunnel_e2e_phoenix:8765;serializationPROTOBUF query select * from test.SOURCE } } sink { Jdbc { driver org.apache.phoenix.queryserver.client.Driver url jdbc:phoenix:thin:urlhttp://seatunnel_e2e_phoenix:8765;serializationPROTOBUF 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 11完成清空见 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-beta2022-09-26新增 Phoenix Sink Connector。延伸阅读Jdbc Sink ConnectorPhoenix 的载体与完整参数表Jdbc Source ConnectorPhoenix 数据读取Sink 通用参数说明连接器能力特性总览batch / stream / exactly-once 等赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐SeaTunnel Phoenix Sink 实战指南基于 Jdbc 连接器将数据 UPSERT 写入 HBaseSeaTunnel 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大数据批处理流处理变更数据捕获上一篇OneDrive Free Client选择性同步只同步你需要的文件和文件夹下一篇终极指南如何让Intel无线网卡在Mac上获得原生Wi-Fi体验创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表