数据工程数据编排ETL任务调度批处理流处理数据集成后端【免费下载链接】mage-ai Build, run, and manage data pipelines for integrating and transforming data.项目地址https://gitcode.com/gh_mirrors/ma/mage-ai点击查看免费下载输出文章mage-ai MySQL 数据源接入指南配置参数、连接方式与源码实现解析本指南围绕 mage-ai及 mage_integrations 子项目中的 MySQL 数据源Data Source展开系统讲解将 MySQL 作为数据集成管道上游时所需的核心配置项、连接认证方式、批量读取机制与底层实现原理。读完本文你将掌握在 mage-ai 中完整配置 MySQL Source、选择直连或 SSH 隧道连接、通过conn_kwargs传入 SSL 等连接参数以及调整batch_fetch_limit控制读取批次规模并了解其背后的源码实现与测试验证方式。概览MySQL 在 mage-ai 数据集成中的定位mage-ai 是一个用于构建、运行和管理数据管道的开源平台其数据集成Data Integration能力由mage_integrations子项目提供。其中sources/mysql目录实现了 MySQL 数据源模块职责是从指定的 MySQL 数据库中读取数据并交付给管道中的后续块Block或目标端。该模块的核心入口位于 sources/mysql/init.py其MySQL类继承自 SQL 数据源的通用基类Source见 sources/sql/base.py因此天然具备模式发现discover、批量数据读取load_data、记录计数count_records、断点/增量读取bookmarks等通用能力并针对 MySQL 语法做了定制例如使用反引号包裹列名保证与 MySQL 的标识符规则兼容通过information_schema.columns发现表结构与主键信息将 DATETIME/TIMESTAMP 等类型映射为DATETIME将整数类型映射为UNSIGNED。提示mage_integrations是 mage-ai 独立维护的数据集成引擎除 MySQL 外还包含 PostgreSQL、BigQuery、Snowflake、Redshift 等数十个数据源与目标端模块。本文仅聚焦 MySQL 数据源本身。必需配置项连接数据库所需的最小参数集官方文档 sources/mysql/README.md对应站点文档 docs/data-integrations/sources/mysql.mdx指出配置该数据源时必须填写以下凭据Key说明示例值database要读取数据的数据库名称demohost数据库主机名mage.abc.us-west-2.rds.amazonaws.comport数据库运行端口通常为 33063306username访问数据库的用户名需具备对指定 schema 的读写权限rootpassword用户访问数据库的密码abc123...connection_method连接 MySQL 服务器的方式取值为direct或ssh_tunneldirectconn_kwargs可选以字典形式传入的额外连接关键字参数{ssl_ca: CARoot.pem, ssl_cert: certificate.pem, ssl_key: key.pem}在 mage-ai 界面或配置文件中填写时这些字段会落入数据源的config字典由 sources/mysql/init.py 中的build_connection()方法逐项取出并传给底层连接对象。仓库中的模板文件 sources/mysql/templates/config.json 给出了完整的字段骨架可直接作为手工配置的起点{ database: , host: , port: 3306, username: , password: , connection_method: direct, ssh_host: , ssh_port: 22, ssh_username: , ssh_password: , ssh_pkey: }从模板可以看到port默认值为3306connection_method默认值为directssh_port默认值为22。这些默认值与源码中的实现保持一致详见下文。配置字段如何被解析与使用在MySQL.build_connection()中配置字段与连接对象参数一一对应return MySQLConnection( databaseself.config[database], hostself.config[host], passwordself.config[password], portself.config.get(port, 3306), usernameself.config[username], connection_methodself.config.get(connection_method, ConnectionMethod.DIRECT), conn_kwargsself.config.get(conn_kwargs), ssh_hostself.config.get(ssh_host), ssh_portself.config.get(ssh_port, 22), ssh_usernameself.config.get(ssh_username), ssh_passwordself.config.get(ssh_password), ssh_pkeyself.config.get(ssh_pkey), verbose0 if self.discover_mode or self.discover_streams_mode else 1, )关键信息database、host、password、username是必填键源码直接使用self.config[xxx]取值缺失会抛KeyErrorport、connection_method、ssh_port使用self.config.get(key, default)形式具备默认值兜底分别为3306、direct、22connection_method与ConnectionMethod枚举绑定取值必须是direct或ssh_tunnel二者之一。连接方式一直连directdirect是默认且最简单的连接方式。底层连接实现位于 connections/mysql/init.py它使用 Python 官方 MySQL 驱动mysql.connector.connect()建立连接return connect( databaseself.database, hosthost, passwordself.password, portport, userself.username, **self.conn_kwargs, )conn_kwargs会以关键字展开**kwargs的方式透传给mysql.connector.connect因此你可以利用 MySQL Connector/Python 支持的任意连接参数最典型的是 SSL 相关配置例如官方 README 给出的{ssl_ca: CARoot.pem, ssl_cert: certificate.pem, ssl_key: key.pem}当你需要从公网直连云数据库如 AWS RDS、阿里云 RDS 等且数据库启用了 TLS/SSL 时通过conn_kwargs传入 CA 证书、客户端证书与私钥即可建立加密连接。连接方式二SSH 隧道ssh_tunnel当 MySQL 位于私有网络如内网、VPC 内部无法直接访问时可以设置connection_method为ssh_tunnel借助中间堡垒机bastion host转发连接。需要配合以下可选参数使用Key说明示例值ssh_host中间堡垒机的主机地址123.45.67.89ssh_port堡垒机端口默认 2222ssh_username连接堡垒机使用的用户名usernamessh_password使用密码认证时设置的密码passwordssh_pkey使用私钥认证时私钥文件的路径/path/to/private/keySSH 隧道底层工作方式从 connections/mysql/init.py 的源码可以看到隧道建立的完整流程if self.connection_method ConnectionMethod.SSH_TUNNEL: ssh_setting dict(ssh_usernameself.ssh_username) if self.ssh_pkey is not None: if os.path.exists(self.ssh_pkey): ssh_setting[ssh_pkey] self.ssh_pkey else: ssh_setting[ssh_pkey] paramiko.RSAKey.from_private_key( io.StringIO(self.ssh_pkey), ) else: ssh_setting[ssh_password] self.ssh_password self.ssh_tunnel SSHTunnelForwarder( (self.ssh_host, self.ssh_port), remote_bind_address(self.host, self.port), local_bind_address(, self.port), **ssh_setting, ) self.ssh_tunnel.start() self.ssh_tunnel._check_is_started() host 127.0.0.1 port self.ssh_tunnel.local_bind_port理解这一实现的几个要点项目使用sshtunnel.SSHTunnelForwarder建立本地端口转发将本地端口映射到 MySQL 服务器地址(self.host, self.port)ssh_pkey存在两种输入形态如果传入的是一个已存在的文件路径则直接作为私钥文件使用如果传入的是私钥内容字符串则通过paramiko.RSAKey.from_private_key()将其解析为密钥对象。这意味着ssh_pkey既可以填文件路径也可以直接填 PEM 格式的私钥文本若未配置ssh_pkey则回退为ssh_password密码认证隧道建立后MySQL 客户端实际连接的目标被改写为127.0.0.1与隧道的本地绑定端口从而经加密通道访问远程数据库连接关闭时close_connection()会先关闭 MySQL 连接再停止 SSH 隧道self.ssh_tunnel.stop()确保资源释放。从代码结构看SSH 隧道模块还依赖paramikoSSH 协议实现与sshtunnel端口转发这两者也是 mage_integrations 的依赖项。可选配置batch_fetch_limit 批量读取控制除连接相关配置外数据源还支持一个可选参数Key说明示例值batch_fetch_limit每个批次拉取的行数默认 50k。如果实例内存较大可以调大该值50000源码中的批量读取实现在 sources/sql/base.py 中fetch_limit属性按以下优先级取批次大小property def fetch_limit(self): config self.config or dict() return ( config.get(SUBBATCH_FETCH_LIMIT_KEY) or config.get(BATCH_FETCH_LIMIT_KEY) or BATCH_FETCH_LIMIT )其中常量定义在 sources/constants.pyBATCH_FETCH_LIMIT 50000 SUBBATCH_FETCH_LIMIT 10000 BATCH_FETCH_LIMIT_KEY batch_fetch_limit SUBBATCH_FETCH_LIMIT_KEY subbatch_fetch_limit即取值顺序为subbatch_fetch_limitbatch_fetch_limit 内置默认值50000。load_data()方法sources/sql/base.py以此为基础进行分批翻页读取通过LIMIT {limit} OFFSET {offset}逐批拉取并循环直到len(rows_temp) limit不足一批或达到自定义_limit为止。因此调大batch_fetch_limit可减少网络往返次数、提升吞吐但会增加单批数据在内存中的占用需要与实例内存匹配若实例内存受限也可调小该值以降低内存峰值。模式发现Discover与类型映射的 MySQL 特化在建立连接后数据源会通过discover过程读取目标库的表与列结构。MySQL.build_discover_query()sources/mysql/init.py查询information_schema.columnsSELECT TABLE_NAME , COLUMN_DEFAULT , COLUMN_KEY , COLUMN_NAME , COLUMN_TYPE , IS_NULLABLE FROM information_schema.columns WHERE table_schema {database}当指定了streams即选定的表列表时查询会追加AND TABLE_NAME IN (...)进行过滤只发现选定表的 schema。发现结果会被转换成 Catalog目录其中COLUMN_KEY为PRI主键或UNIQUE的列会作为key_properties/unique_constraints参与数据同步的排序与冲突处理unique_conflict_method默认取UPDATE具体逻辑见 sources/sql/base.py。列名与类型的 MySQL 适配MySQL类针对 MySQL 方言做了两处关键适配sources/mysql/init.pydef column_type_mapping(self, column_type: str, column_format: str None) - str: if COLUMN_FORMAT_DATETIME column_format: return DATETIME elif COLUMN_TYPE_INTEGER column_type: return UNSIGNED return super().column_type_mapping(column_type, column_format) def update_column_names(self, columns: List[str]) - List[str]: return list(map(lambda column: self.wrap_column_in_quotes(column), columns)) def wrap_column_in_quotes(self, column: str) - str: if not in column: return f{column} return column日期时间类型的列在 SQL 中映射为DATETIME整数类型映射为UNSIGNED所有列名在生成 SQL 时都会被反引号包裹避免与 MySQL 保留字冲突且若列名本身已带反引号则不会重复包裹。测试用例验证discover 行为仓库在 tests/sources/mysql/test_mysql.py 中提供了针对 discover 流程的单元测试。测试模拟了一张demo_users表含主键id、varchar字段、enum字段、timestamp、float等典型类型并断言discover()生成的 catalog 中tap_stream_id为demo_users复制方法replication_method为FULL_TABLEkey_properties与unique_constraints均为[id]来自PRI主键标记各列类型被正确映射int→integer、varchar→string、timestamp→{format: date-time, type: string}、float→number、可空列类型中追加null。该测试印证了上文关于主键发现与类型映射的描述是理解 MySQL 数据源行为的直接参考。完整配置示例汇总综合以上内容一个面向生产场景的 MySQL 数据源配置可以这样组织以 JSON 形式直连 SSL 场景{ database: demo, host: mage.abc.us-west-2.rds.amazonaws.com, port: 3306, username: root, password: abc123..., connection_method: direct, conn_kwargs: { ssl_ca: CARoot.pem, ssl_cert: certificate.pem, ssl_key: key.pem }, batch_fetch_limit: 50000 }内网场景SSH 隧道 私钥认证{ database: demo, host: 10.0.0.5, port: 3306, username: root, password: abc123..., connection_method: ssh_tunnel, ssh_host: 123.45.67.89, ssh_port: 22, ssh_username: username, ssh_pkey: /path/to/private/key, batch_fetch_limit: 50000 }说明私钥认证时ssh_pkey既支持文件路径也支持直接粘贴私钥内容字符串当使用密码认证时则填写ssh_password二者二选一。batch_fetch_limit需根据执行实例的内存情况调整。总结在 mage-ai 中接入 MySQL 数据源的核心要点可归纳为最小配置database、host、port、username、password为必填项默认端口3306两种连接方式direct直连可配合conn_kwargs传 SSL 等驱动参数与ssh_tunnelSSH 隧道支持密码与私钥两种认证批量读取控制batch_fetch_limit默认 50k可依据内存调整源码中还保留了subbatch_fetch_limit默认 10k作为更高优先级配置MySQL 方言适配通过information_schema.columns发现 schema反引号包裹列名并特化 DATETIME / UNSIGNED 类型映射可验证性可通过 tests/sources/mysql/test_mysql.py 的单元测试理解 discover 与类型映射的实际行为。相关参考文件官方数据源说明 sources/mysql/README.md 与站点文档 docs/data-integrations/sources/mysql.mdx、源码实现 sources/mysql/init.py 与 connections/mysql/init.py、通用 SQL 基类 sources/sql/base.py 及配置模板 sources/mysql/templates/config.json。赞分享数据工程数据编排ETL任务调度批处理流处理数据集成后端【免费下载链接】mage-ai Build, run, and manage data pipelines for integrating and transforming data.项目地址https://gitcode.com/gh_mirrors/ma/mage-ai点击查看免费下载相关推荐Apache DolphinScheduler Oracle 数据源接入指南参数配置、ServiceName/SID 连接模式与源码实现解析Apache DolphinScheduler Oracle 数据源接入指南参数配置、ServiceName/SID 连接模式与源码实现解析 本指南以 Apa任务调度大数据后端前端Apache DolphinScheduler Vertica 数据源接入指南参数配置、连接原理与源码实现Apache DolphinScheduler Vertica 数据源接入指南参数配置、连接原理与源码实现 Apache DolphinScheduler 内任务调度大数据后端前端Apache DolphinScheduler 接入 Trino 数据源表单配置、Jdbc 连接参数与源码实现解析Apache DolphinScheduler 接入 Trino 数据源表单配置、Jdbc 连接参数与源码实现解析 本文基于 Apache DolphinSc任务调度大数据后端前端上一篇MatBlazor无障碍访问创建包容性Web应用的完整指南下一篇终极指南KubeSphere日志管理从FluentBit到OpenSearch的完整链路创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
