ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

SeaTunnel Databend Source 连接器实战指南:JDBC 批式读取、SQL 优先级与类型映射全解析

SeaTunnel Databend Source 连接器实战指南:JDBC 批式读取、SQL 优先级与类型映射全解析 SeaTunnel Databend Source 连接器实战指南JDBC 批式读取、SQL 优先级与类型映射全解析【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelSeaTunnel 的 Databend Source 连接器基于 Databend JDBC 驱动将 Databend 中的数据以批处理方式读取并转换为 SeaTunnel 行支持单表读取、一次性查询与完整 SQL 语句三种入口。本文将完整讲解该连接器的依赖安装、数据类型映射、全部源选项语义与 SQL 选择优先级并结合仓库源码剖析其读取实现原理帮助你快速搭建 Databend → SeaTunnel 的数据集成管道。连接器概述与核心能力Databend Source 是 SeaTunnel 官方提供的批式源连接器插件标识为Databend实现位于 connector-databend 模块。它的核心职责是通过 JDBC 在 Databend 上执行 SQL 查询把结果集的每一行转换为SeaTunnelRow交给下游算子处理。从源码看连接器的核心类包括DatabendSourceFactory工厂类定义OptionRule校验规则并构建 SQL 语句DatabendSource源实现继承AbstractSingleSplitSource声明有界BOUNDEDDatabendSourceReader读取器负责建立 JDBC 连接、执行查询、逐行转换。支持的引擎引擎支持情况Spark✅Flink✅SeaTunnel Zeta✅功能特性矩阵特性支持情况批处理✅流处理❌并行度✅支持用户自定义分片❌支持多表读❌连接器当前不支持在同一 source 块中读取多张表也不支持table_list参数如果需要读取多张表必须为每张表分别配置一个 Databend source。从源码结构看DatabendSource继承的是单分片基类AbstractSingleSplitSource读取逻辑在单个 split 内完成并行度提升主要依赖执行引擎的任务调度层面这与其不支持用户自定义分片的特性声明是一致的。依赖与驱动安装Databend 连接器依赖connector-databend插件包和 Databend 官方 JDBC 驱动安装方式与运行引擎有关Spark / Flink 引擎下载 Databend JDBC driver jar 包放到${SEATUNNEL_HOME}/plugins/目录SeaTunnel Zeta 引擎下载 Databend JDBC driver jar 包放到${SEATUNNEL_HOME}/lib/目录。此外连接器插件本身可以通过install-plugin.sh脚本或从 Maven 中央仓库获取。数据源支持版本依赖Databend1.2.x 及以上版本org.apache.seatunnel:connector-databend驱动加载逻辑位于 DatabendSourceReader.open()通过Class.forName(com.databend.jdbc.DatabendDriver)显式加载驱动类若驱动未正确放入对应目录会在 open 阶段抛出DatabendConnectorException错误码为CONNECT_FAILED。数据类型映射连接器支持 Databend 常见数据类型到 SeaTunnel 类型的映射官方映射关系如下Databend 数据类型SeaTunnel 数据类型BOOLEANBOOLEANTINYINTTINYINTSMALLINTSMALLINTINTINTBIGINTBIGINTFLOATFLOATDOUBLEDOUBLEDECIMALDECIMALSTRINGSTRINGVARCHARSTRINGCHARSTRINGTIMESTAMPTIMESTAMPDATEDATETIMETIMEBINARYBYTES值得说明的是DatabendSourceReader中的 convertDatabendTypeToSeaTunnelType() 在运行时还会做更细粒度的推断字符串族VARCHAR、STRING、TEXT、CHAR均映射为STRING_TYPE整数族TINYINT/UINT8/INT8映射为字节类型SMALLINT/UINT16/INT16映射为短整型INT/INTEGER/UINT32/INT32映射为整型BIGINT/UINT64/INT64映射为长整型浮点族FLOAT/FLOAT32映射为FLOAT_TYPEDOUBLE/FLOAT64映射为DOUBLE_TYPE时间族DATE映射为本地日期类型TIMESTAMP/DATETIME映射为本地日期时间类型二进制族BINARY/BLOB映射为字节数组类型当通过 JDBCjava.sql.Types推断 DECIMAL/NUMERIC 但精度信息缺失时默认回退为DECIMAL(38, 18)对无法识别的 SQL 类型会记录 warn 日志并回退为STRING_TYPE。此外DatabendTypeConverter 定义了反方向的 SeaTunnel → Databend 类型映射主要用于 Sink 场景建表其中ARRAY映射为ARRAY(STRING)、MAP映射为MAP(STRING, STRING)两套映射共同构成了连接器完整的类型转换体系。源选项详解基础配置项定义在 DatabendOptions 与 DatabendSourceOptions 中完整的必填/可选约束由DatabendSourceFactory.optionRule()统一校验。名称类型是否必须默认值描述urlString是-Databend JDBC 连接 URL必须以jdbc:databend://开头usernameString是-Databend 数据库用户名passwordString是-Databend 数据库密码databaseString否-Databend 数据库名称默认使用连接 URL 中指定的数据库名tableString否-Databend 表名称queryString否-Databend 查询语句设置后覆盖 database 和 table 的设置sqlString否-自定义 SQL 语句若同时配置sql和query优先使用sqlfetch_sizeInteger否1每次从 Databend 拉取的记录数读取大量数据时可适当调大设为0使用 JDBC 驱动默认值sslBoolean否false是否使用 SSL 连接 Databendjdbc_configMap否-额外的 JDBC 连接配置如负载均衡策略等common-options否-源插件常用参数详见源通用选项选项背后的源码约束URL 前缀强校验DatabendSourceFactory.optionRule()使用Conditions.startsWith(URL, jdbc:databend://)做条件校验URL 不满足该前缀时配置会直接校验失败见 DatabendSourceFactory.java必填项url、username、password为必填其余均为可选fetch_size 语义源码默认值为1DatabendOptions.FETCH_SIZE的defaultValue(1)。读取器在 open() 中只有fetch_size 0时才调用statement.setFetchSize(fetchSize)并设置FETCH_FORWARD方向设为0或未配置时则完全交给 JDBC 驱动自行决定ssl 传递ssl参数会被写入 JDBC 连接的 Propertiesproperties.setProperty(ssl, ...)随DriverManager.getConnection一起传给驱动见 DatabendSource.createReader()。SQL 入口的优先级规则必须配置sql、query或同时配置database和table三者之一。当多个入口同时配置时实际读取 SQL 的优先级为sqlquerySELECT * FROM database.table该规则直接对应DatabendSourceFactory.buildSqlStatement()的实现源码先检查sql再检查query最后拼接database.table若三者均未配置则抛出SQL_OPERATION_FAILED异常并提示 Either SQL, query, or both database and table must be specified。任务示例单表读取最基本的用法是databasetable组合适合整表导入env { parallelism 2 job.mode BATCH } source { Databend { url jdbc:databend://localhost:8000 username root password database default table users } } sink { Console {} }该配置等价于执行SELECT * FROM default.users结果输出到 Console sink 便于验证。使用自定义查询通过query可以执行一次性即席查询并覆盖database/table的配置source { Databend { url jdbc:databend://localhost:8000 username root password query SELECT id, name, age FROM default.users WHERE age 18 } }使用 SSL连接启用了 SSL 的 Databend 集群时配置ssl true并建议同步调大fetch_sizesource { Databend { url jdbc:databend://databend.example.com:8000/default username root password sql SELECT * FROM default.users ssl true fetch_size 1000 } }注意此处 URL 中直接携带了数据库名default此时可以省略database选项。在查询中过滤和投影Databend 支持丰富的表达式可以直接在query中使用任意合法表达式在数据进入 SeaTunnel 之前完成列裁剪与行过滤减少网络传输与下游处理压力source { Databend { url jdbc:databend://localhost:8000 username root password query SELECT id, name, age FROM default.users WHERE age 18 AND starts_with(name, A) ORDER BY id } }该示例展示了标准 SQL 列投影只取id、name、age三列、条件过滤年龄不低于 18 且姓名以A开头以及排序的完整组合用法。底层读取原理与 Schema 推断理解读取链路有助于排查问题与做性能调优。整个读取流程如下工厂构建 SQLDatabendSourceFactory.createSource()按优先级确定最终执行的 SQL 字符串Catalog 取 Schema工厂会尝试通过DatabendCatalog.getTable()获取表的 CatalogTable 与字段 Schema如果取不到例如使用复杂query时会记录 warn 日志并降级为从查询结果推断 Schema的空 Schema 方案打开读取器DatabendSourceReader.open()加载驱动 → 建立连接 →prepareStatement(sql)→ 按fetch_size设置抓取大小 →executeQuery()Schema 推断与行转换如果初始 rowType 为空读取器从ResultSetMetaData推断列名与列类型调用convertDatabendTypeToSeaTunnelType随后在internalPollNext()中通过DatabendUtil.convertToSeaTunnelRow()逐行转换并output.collect(row)输出结束与释放读取完最后一行后调用readerContext.signalNoMoreElement()通知引擎并在close()中依次关闭ResultSet、PreparedStatement、Connection。由于DatabendSource.getBoundedness()返回BOUNDED源码该连接器天然只支持批处理作业job.mode应设置为BATCH这也与不支持流处理的特性声明一致。注意事项与常见问题多表读取当前版本不支持table_list一个 source 块只能读一张表多表场景需配置多个 Databend source 块驱动缺失Spark/Flink 与 Zeta 引擎的驱动放置目录不同plugins/vslib/放错位置会在运行时抛出ClassNotFoundExceptioncom.databend.jdbc.DatabendDriver或连接失败异常fetch_size 调优默认值1意味着每次从数据库拉取一条记录适合小数据量验证读取大表时建议调大如 1000 及以上以减少与数据库的交互次数设为0则交给 JDBC 驱动默认策略SSL 配置开启ssl true前请确认 Databend 服务端已启用 TLS 且 JDBC 驱动版本与之兼容SQL 优先级陷阱同时配置sql与query时只会执行sql误配置可能导致实际读取的语句与预期不符建议同一 source 块中只保留一个读取入口。本文基于当前仓库的 中文版 Databend 源连接器文档 与connector-databend模块源码整理官方文档还维护了该连接器的变更记录Changelog可结合 文档 changelog 目录 了解各版本的能力演进。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表