工作流自动化流程编排任务调度数据工程后端【免费下载链接】prefectPrefect is a workflow orchestration framework for building resilient data pipelines in Python.项目地址https://gitcode.com/GitHub_Trending/pr/prefect点击查看免费下载prefect-sqlalchemy是 Prefect 官方提供的数据库集成包帮助你通过 SQLAlchemy 在 Prefect flow 中安全、高效地读写各类数据库。本指南基于当前仓库中 prefect-sqlalchemy 集成文档 及其源码实现系统讲解安装注册、连接信息建模、凭据 Block 的保存与加载、同步/异步数据库操作以及底层的结果缓存与生命周期管理机制读完即可在真实 pipeline 中落地使用。集成包概览prefect-sqlalchemy位于仓库的 src/integrations/prefect-sqlalchemy 目录对外暴露四个核心构件见 prefect_sqlalchemy/init.py构件类型用途SqlAlchemyConnectorBlock使用同步驱动的数据库连接块AsyncSqlAlchemyConnectorBlock使用异步驱动的数据库连接块ConnectionComponentsPydantic 模型将连接参数组装为 SQLAlchemy engine URLSyncDriver/AsyncDriver枚举预定义的同步 / 异步数据库驱动方言该包依赖sqlalchemy1.4.31,3与prefect3.6.17要求 Python 3.11兼容 3.11 ~ 3.13见 pyproject.toml。安装与注册 Block 类型安装与当前已装 Prefect 版本兼容的prefect-sqlalchemy若未安装 Prefect 会一并安装最新版pip install prefect[sqlalchemy]升级到两者最新版本pip install -U prefect[sqlalchemy]安装后注册prefect_sqlalchemy模块中新增的 Block 类型使其在 Prefect UI 和 CLI 中可用prefect block register -m prefect_sqlalchemy注册完成后即获得SQLAlchemy Connector与Async SQLAlchemy Connector两种 Block 类型。仓库中的 test_block_standards.py 通过BlockStandardTestSuite对这两个 Block 逐一验证其符合 Prefect Block 标准。构建连接信息ConnectionComponents无论同步还是异步连接器其connection_info字段都接受两种形式ConnectionComponents对象——把连接参数结构化组装成 engine URLURL 字符串——直接传入形如postgresqlpsycopg2://user:passlocalhost:5432/db的连接串。字段详解ConnectionComponents定义在 credentials.py字段如下字段类型默认值说明driverAsyncDriver/SyncDriver/str必填使用的驱动名databaseOptional[str]None数据库名usernameOptional[str]None认证用户名passwordOptional[SecretStr]None认证密码以SecretStr存储避免明文泄露hostOptional[str]None数据库主机地址portOptional[int]None连接端口queryOptional[Dict[str, str]]None传给方言 / DBAPI 的连接参数若需向 Python DBAPI 传非字符串参数请改用connect_argscreate_url 的组装逻辑create_url()方法credentials.py#L130将上述字段拼装为sqlalchemy.engine.url.URL先取出枚举成员的实际字符串值driver.value再从SecretStr中解出明文密码最后仅将非None的参数传给URL.create。测试 test_credentials.py 验证了两种组装场景from prefect_sqlalchemy.credentials import ConnectionComponents, SyncDriver from sqlalchemy.engine.url import make_url # 最小参数仅 driver database components ConnectionComponents( driverSyncDriver.POSTGRESQL_PSYCOPG2, databasemy.db ) assert components.create_url() make_url(postgresqlpsycopg2:///my.db) # 完整参数生成带认证信息的 URL components ConnectionComponents( driverSyncDriver.POSTGRESQL_PSYCOPG2, databasemy.db, usernamemyusername, passwordmypass, port1234, hostlocalhost, ) assert components.create_url() make_url( postgresqlpsycopg2://myusername:mypasslocalhost:1234/my.db )可见所需参数随驱动而异例如 PostgreSQL 通常需要driverusernamepasswordhostportdatabase而 SQLite 只需driver和database。直接使用 URL 字符串当connection_info传入字符串时连接器会通过check_make_url校验器database.py#L23-L31调用make_url验证 URL 合法性非法输入如plskeepmydata会抛出ValueError: Invalid URL该行为由 test_database.py 的test_connector_init_fails_with_invalid_url用例覆盖。驱动枚举SyncDriver 与 AsyncDriver两个枚举类定义在 credentials.py覆盖 PostgreSQL、MySQL、SQLite、Oracle、MSSQL 等主流数据库。同步驱动SyncDriver枚举成员驱动字符串适用数据库POSTGRESQL_PSYCOPG2postgresqlpsycopg2PostgreSQLPOSTGRESQL_PG8000postgresqlpg8000PostgreSQLPOSTGRESQL_PSYCOPG2CFFIpostgresqlpsycopg2cffiPostgreSQLPOSTGRESQL_PYPOSTGRESQLpostgresqlpypostgresqlPostgreSQLPOSTGRESQL_PYGRESQLpostgresqlpygresqlPostgreSQLMYSQL_MYSQLDBmysqlmysqldbMySQLMYSQL_PYMYSQLmysqlpymysqlMySQLMYSQL_MYSQLCONNECTORmysqlmysqlconnectorMySQLMYSQL_CYMYSQLmysqlcymysqlMySQLMYSQL_OURSQLmysqloursqlMySQLMYSQL_PYODBCmysqlpyodbcMySQLSQLITE_PYSQLITEsqlitepysqliteSQLiteSQLITE_PYSQLCIPHERsqlitepysqlcipherSQLiteORACLE_CX_ORACLEoraclecx_oracleOracleORACLE_ORACLEDBoracleoracledbOracleMSSQL_PYODBCmssqlpyodbcSQL ServerMSSQL_MXODBCmssqlmxodbcSQL ServerMSSQL_PYMSSQLmssqlpymssqlSQL Server异步驱动AsyncDriver枚举成员驱动字符串适用数据库POSTGRESQL_ASYNCPGpostgresqlasyncpgPostgreSQLSQLITE_AIOSQLITEsqliteaiosqliteSQLiteMYSQL_ASYNCMYmysqlasyncmyMySQLMYSQL_AIOMYSQLmysqlaiomysqlMySQLORACLE_ORACLEDB_ASYNCoracleoracledb_asyncOracle保存与加载凭据 Block使用 Block 的load方法前必须先通过代码或 Prefect UI 保存一个 Block。PostgreSQL 示例from prefect_sqlalchemy import SqlAlchemyConnector, ConnectionComponents, SyncDriver connector SqlAlchemyConnector( connection_infoConnectionComponents( driverSyncDriver.POSTGRESQL_PSYCOPG2, usernameUSERNAME-PLACEHOLDER, passwordPASSWORD-PLACEHOLDER, hostlocalhost, port5432, databaseDATABASE-PLACEHOLDER, ) ) connector.save(BLOCK_NAME-PLACEHOLDER)保存后在任何 flow / task 中按名称加载from prefect_sqlalchemy import SqlAlchemyConnector SqlAlchemyConnector.load(BLOCK_NAME-PLACEHOLDER)SQLite 最小示例SQLite 仅需driver和database两个参数from prefect_sqlalchemy import SqlAlchemyConnector, ConnectionComponents, SyncDriver connector SqlAlchemyConnector( connection_infoConnectionComponents( driverSyncDriver.SQLITE_PYSQLITE, databaseDATABASE-PLACEHOLDER.db ) ) connector.save(BLOCK_NAME-PLACEHOLDER)在 Flow 中使用数据库同步示例在 flow 中通过 task 完成建表、写入与分页读取。要点execute执行不返回数据的操作如CREATE、INSERT、UPDATE、DELETE调用即执行execute_many批量执行带参数序列的操作fetch_many按指定size流式取回数据对同一操作重复调用fetch_*不会重复执行 SQL而是返回上一批结果之后的下一批数据详见下文结果缓存机制上下文管理器with语句确保引擎与连接资源在退出后正确释放。from prefect import flow, task from prefect_sqlalchemy import SqlAlchemyConnector task def setup_table(block_name: str) - None: with SqlAlchemyConnector.load(block_name) as connector: connector.execute( CREATE TABLE IF NOT EXISTS customers (name varchar, address varchar); ) connector.execute( INSERT INTO customers (name, address) VALUES (:name, :address);, parameters{name: Marvin, address: Highway 42}, ) connector.execute_many( INSERT INTO customers (name, address) VALUES (:name, :address);, seq_of_parameters[ {name: Ford, address: Highway 42}, {name: Unknown, address: Highway 42}, ], ) task def fetch_data(block_name: str) - list: all_rows [] with SqlAlchemyConnector.load(block_name) as connector: while True: # 对同一操作的重复 fetch* 调用会跳过重新执行返回下一批结果 new_rows connector.fetch_many(SELECT * FROM customers, size2) if len(new_rows) 0: break all_rows.append(new_rows) return all_rows flow def sqlalchemy_flow(block_name: str) - list: setup_table(block_name) all_rows fetch_data(block_name) return all_rows if __name__ __main__: sqlalchemy_flow(BLOCK-NAME-PLACEHOLDER)该完整流程与 test_database.py 中的test_flow_without_initialized_engine高度一致setup_table建表并写入 3 行数据fetch_data以每批 2 行循环取完所有数据最终断言结果为[[(Marvin, Highway 42), (Ford, Highway 42)], [(Unknown, Highway 42)]]。异步支持AsyncSqlAlchemyConnector对于使用异步驱动的异步工作流如AsyncDriver.SQLITE_AIOSQLITE、AsyncDriver.POSTGRESQL_ASYNCPG应使用AsyncSqlAlchemyConnector而非同步版本。其 API 与同步版一一对应但方法均为协程需配合async with与await使用from prefect import flow, task from prefect_sqlalchemy import AsyncSqlAlchemyConnector import asyncio task async def setup_table(block_name: str) - None: async with await AsyncSqlAlchemyConnector.load(block_name) as connector: await connector.execute( CREATE TABLE IF NOT EXISTS customers (name varchar, address varchar); ) await connector.execute( INSERT INTO customers (name, address) VALUES (:name, :address);, parameters{name: Marvin, address: Highway 42}, ) await connector.execute_many( INSERT INTO customers (name, address) VALUES (:name, :address);, seq_of_parameters[ {name: Ford, address: Highway 42}, {name: Unknown, address: Highway 42}, ], ) task async def fetch_data(block_name: str) - list: all_rows [] async with await AsyncSqlAlchemyConnector.load(block_name) as connector: while True: # 对同一操作的重复 fetch* 调用会跳过重新执行返回下一批结果 new_rows await connector.fetch_many(SELECT * FROM customers, size2) if len(new_rows) 0: break all_rows.append(new_rows) return all_rows flow async def sqlalchemy_flow(block_name: str) - list: await setup_table(block_name) all_rows await fetch_data(block_name) return all_rows if __name__ __main__: asyncio.run(sqlalchemy_flow(BLOCK-NAME-PLACEHOLDER))底层机制深度解析以下机制均可在 database.py 源码中直接印证。结果缓存与流式读取fetch_one/fetch_many/fetch_all并不每次都执行 SQL。_get_result_set会用hash_objects对操作与参数计算哈希首次遇到该输入时创建连接并执行结果集缓存在_unique_results字典中再次以相同输入调用时直接复用上次的游标返回下一条 / 下一批数据。这种设计让fetch_many(..., size2)的分页循环不必重写 SQL。调用reset_connections()会关闭所有缓存的结果集与连接使后续fetch_*从头开始返回数据。测试test_fetch_one、test_fetch_many、test_fetch_all均验证了重复调用返回后续数据 → reset 后从头开始的行为。fetch_size 默认值fetch_size字段默认值为1database.py#L114-L116。fetch_many未显式传size或传 0 时使用块上配置的fetch_sizesize size or self.fetch_size。测试 test_database.py 通过fetch_size2的 fixture 验证了sizeNone/1/2三种取值的行为。引擎与连接的生命周期延迟初始化block_initialization只负责把ConnectionComponents转为 URL 并校验驱动类型真正的引擎在首次调用get_engine()时才创建测试test_delay_start验证了_engine初始为None执行操作后才成为Engine实例。get_engine若已有引擎则直接复用日志Reusing existing engine.否则通过create_engine/create_async_engine创建。get_connectionbeginTrue默认时使用engine.begin()开启事务任一操作失败则整体回滚beginFalse时使用engine.connect()。get_client统一入口按client_typeengine或connection分发到上述两个方法非法值抛出ValueError。上下文管理器SqlAlchemyConnector.__exit__调用close()先reset_connections()关闭全部连接再engine.dispose()释放引擎异步版本AsyncSqlAlchemyConnector.__aexit__对应await self.close()。可序列化__getstate__/__setstate__使 Block 可被 pickle在跨 task / flow 传递时引擎与结果集会被置空反序列化后按需重建。驱动类型强制校验同步连接器初始化时若检测到异步驱动如把SQLITE_AIOSQLITE传给SqlAlchemyConnector会抛出ValueError提示改用AsyncSqlAlchemyConnector反之异步连接器遇到同步驱动同样报错。这一行为由 test_async_dispatch.py 的TestDriverTypeEnforcement用例覆盖。使用建议与限制在单个 task / flow 内加载并消费 Block官方文档与类 docstring 均提示若把 Block 跨 task / flow 传递其连接与游标状态可能丢失推荐在 task 内load并配合上下文管理器使用。密码以SecretStr存储ConnectionComponents.password使用 Pydantic 的SecretStr序列化与展示时不会暴露明文。额外连接参数需要向 DBAPI 传非字符串参数时使用块的connect_args字段query仅用于字符串键值对。注册 Block 后才能用 UI 管理执行prefect block register -m prefect_sqlalchemy后可在 UI 中创建、保存和分享连接块。参考资料集成官方文档docs/integrations/prefect-sqlalchemy/index.mdxSDK 参考——凭据类docs/integrations/prefect-sqlalchemy/api-ref/prefect_sqlalchemy-credentials.mdxSDK 参考——数据库类docs/integrations/prefect-sqlalchemy/api-ref/prefect_sqlalchemy-database.mdx源码实现prefect_sqlalchemy/database.py、prefect_sqlalchemy/credentials.py测试用例tests/test_database.py、tests/test_credentials.py、tests/test_async_dispatch.py、tests/test_block_standards.py打包配置src/integrations/prefect-sqlalchemy/pyproject.toml赞分享工作流自动化流程编排任务调度数据工程后端【免费下载链接】prefectPrefect is a workflow orchestration framework for building resilient data pipelines in Python.项目地址https://gitcode.com/GitHub_Trending/pr/prefect点击查看免费下载相关推荐如何为LEGION_Y7000Series_Hackintosh制作自定义DSDT/SSDT补丁如何为LEGION_Y7000Series_Hackintosh制作自定义DSDT/SSDT补丁 在Lenovo LEGION Y7000/Y7000P系列笔记20K Vocab Builder 提示词深度拆解自适应词汇学习 GPT 的 4 步闭环设计20K Vocab Builder 提示词深度拆解自适应词汇学习 GPT 的 4 步闭环设计 本仓库是一个 GPTs 泄露提示词合集20K Vocab Bu提示工程3分钟上手Prefect数据库连接从SQLAlchemy到异步查询全攻略3分钟上手Prefect数据库连接从SQLAlchemy到异步查询全攻略 你是否还在为数据 pipeline 中的数据库连接稳定性发愁作为数据工程师我们都工作流自动化流程编排任务调度数据工程后端创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考