ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

dlt 实战教程:使用 sql_database 源将 SQL 数据库加载到 DuckDB(含 append / replace / merge 与增量加载)

dlt 实战教程:使用 sql_database 源将 SQL 数据库加载到 DuckDB(含 append / replace / merge 与增量加载) dlt 实战教程使用 sql_database 源将 SQL 数据库加载到 DuckDB含 append / replace / merge 与增量加载【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy ️项目地址: https://gitcode.com/GitHub_Trending/dl/dlt本篇教程基于 dltdata load tool官方 SQL Database 教程展开完整演示如何使用dlt init sql_database duckdb初始化项目、以公开的 MySQL Rfam 数据库为数据源、通过sql_database/sql_table源将family、genome等表加载到本地 DuckDB并深入讲解append、replace、merge三种写入策略与基于游标列的增量加载同时结合当前仓库源码dlt/sources/sql_database/说明其底层实现原理。读完本文你将能够独立搭建一个可复现的「SQL 数据库 → dlt 目标库」管道并掌握按表定制写入策略与增量游标的实战技巧。背景dlt 的 SQL Database 核心源dlt 通过sql_database源与sql_table资源支持从 SQL 数据库抽取数据借助 SQLAlchemy 方言可覆盖 PostgreSQL、MySQL、Microsoft SQL Server、Oracle、IBM DB2、SQLite、MariaDB、BigQuery、Snowflake、Redshift、CockroachDB 等 30 数据库见 SQL Database 源索引。数据可被加载到任何 dlt 兼容目标库例如 Postgres、BigQuery、Snowflake、DuckDB 等。本教程将数据源锁定为公开的 MySQL Rfam 数据库只读账号rfamro目标库为本地 DuckDB全程可免费复现。你将从本教程学到初始化并配置一个基础的 SQL 数据库管道实现append、replace、merge三种加载策略基于游标列实现增量加载。0. 前置条件开始前请确保环境满足Python 3.10 或更高版本已创建并激活虚拟环境已安装 dlt参照 安装指南 创建虚拟环境并安装dlt包。1. 创建新的 dlt 项目在当前工作目录下使用dlt init命令初始化项目dlt init sql_database duckdb该 CLI 命令会一次性生成「SQL 数据库 → DuckDB」管道所需的文件与目录。duckdb可替换为任意 受支持的目标库。运行后项目结构如下├── .dlt │ ├── config.toml │ └── secrets.toml ├── sql_database_pipeline.py └── requirements.txt各文件职责sql_database_pipeline.py主脚本用于定义数据管道内含多种预配置的加载示例函数requirements.txt列出项目所需全部 Python 依赖.dlt/存放项目配置文件secrets.toml存放凭据、API Key、Token 等敏感信息config.toml存放 dlt 项目的通用配置。生产环境提示部署到生产环境时把所有配置塞进 TOML 文件可能不便。dlt 强烈建议改用环境变量或其他配置提供器来存放密钥与配置。dlt init生成的sql_database_pipeline.py实际对应仓库内的核心源模板 dlt/_workspace/_templates/_core_source_templates/sql_database_pipeline.py其中预置了load_select_tables_from_database()、load_entire_database()、load_standalone_table_resource()、my_sql_via_pyarrow()、test_pandas_backend_verbatim_decimals()等函数分别演示按表加载、全库加载、独立表资源、pyarrow/pandas 后端与add_limit抽样等场景可作为你后续扩展的参考样板。2. 配置管道脚本生成的项目已经包含sql_database_pipeline.py里面有多组预配置示例。本教程从零编写一个新函数。注意直接运行生成的脚本会执行load_standalone_table_resource()函数请记得在main块中注释掉该函数调用。下面的函数将加载family与genome两张表import dlt from dlt.sources.sql_database import sql_database def load_tables_family_and_genome(): # Create a dlt source that will load tables family and genome source sql_database().with_resources(family, genome) # Create a dlt pipeline object pipeline dlt.pipeline( pipeline_namesql_to_duckdb_pipeline, # Custom name for the pipeline destinationduckdb, # dlt destination to which the data will be loaded dataset_namesql_to_duckdb_pipeline_data # Custom name for the dataset created in the destination ) # Run the pipeline load_info pipeline.run(source) # Pretty print load information print(load_info) if __name__ __main__: load_tables_family_and_genome()逐段解释sql_database源内置两个辅助函数见 源码sql_database()是一个 dlt source 函数会迭代加载传入with_resources()方法中的表本例为family、genomesql_table()是一个 dlt resource 函数加载单张独立表。例如只加载family表时可写sql_table(tablefamily)。dlt.pipeline()创建名为sql_to_duckdb_pipeline、目标为 DuckDB 的 dlt 管道pipeline.run()将数据加载进目标库。源码视角sql_database与sql_table的核心参数阅读 dlt/sources/sql_database/init.py#L38-L58 可以看到sql_database的完整签名除凭据外还支持以下常用参数参数默认值说明schemaNone要加载的数据库 schema 名称区别于默认 schema 时使用metadataNone可选的sqlalchemy.MetaData实例传入后schema参数被忽略table_namesNone要加载的表名列表默认加载 schema 中全部表chunk_size50000每批产出的行数SQLAlchemy 内部会额外开辟两倍于chunk_size的行缓冲区backendsqlalchemy数据产出后端sqlalchemyPython 字典列表、pyarrowArrow 表、pandasDataFrame、connectorxArrow 表最快但忽略chunk_sizereflection_levelfull模式反射级别minimal仅表名/可空性/主键类型由数据推断、full额外反射数据类型并做必要的类型强转、full_with_precision为 decimal/text/binary 等设置精度与 scale并区分 big int 与普通 intdefer_table_reflectNone延迟到产出数据时才连接并反射表结构要求必须显式传table_names适合 Airflow 等执行期才确定 schema 的编排器include_viewsFalse是否把视图与表一起反射显式出现在table_names中的视图名不受此参数影响resolve_foreign_keysFalse将同 schema 内的外键转换为references表提示会带来额外的反射开销engine_kwargsNone直接传给sqlalchemy.create_engine()的关键字参数影响表反射与 SQLAlchemy 后端的数据读取table_loader_classNone自定义表加载器类子类化TableLoaderSQLAlchemy 体系内定制如分页、重试或BaseTableLoader完全不同的后端而sql_tabledlt/sources/sql_database/init.py#L187-L212额外支持incremental、write_disposition、primary_key、merge_key、included_columns、excluded_columns等资源级参数这些参数与apply_hints的效果一一对应。选择表的两种写法差异sql_database(table_names[family, clan])只反射并生成指定表的资源大 schema 下性能更优sql_database().with_resources(family, clan)先反射整个 schema 生成全部资源再过滤出指定表见 配置文档。3. 添加凭据要成功连接 SQL 数据库需要把凭据传入管道。dlt 会自动在生成的 TOML 文件中查找这些信息。将 Rfam 连接信息 填入secrets.toml[sources.sql_database.credentials] drivername mysqlpymysql # databasedialect database Rfam password username rfamro host mysql-rfam-public.ebi.ac.uk port 4497也可以直接粘贴为连接字符串形式sources.sql_database.credentialsmysqlpymysql://rfamromysql-rfam-public.ebi.ac.uk:4497/Rfam连接字符串遵循 SQLAlchemy Database URL 的通用格式dialectdatabase_type://username:passwordserver:port/database_name数据库专属驱动可通过查询参数追加例如 MSSQL 的 ODBC 驱动mssqlpyodbc://username:passwordserver/database?driverODBCDriver17forSQLServer更多凭据格式与连接方式参见 SQL Database 源连接配置。凭据的三种传递方式源码佐证从 dlt/sources/sql_database/helpers.py#L610-L621 的engine_from_credentials可以看到凭据可以是Engine实例、ConnectionStringCredentials对象或原始字符串secrets.toml/ 环境变量推荐不写死在代码中由 dlt 配置系统自动注入脚本内直接传ConnectionStringCredentialsfrom dlt.sources.credentials import ConnectionStringCredentials from dlt.sources.sql_database import sql_database credentials ConnectionStringCredentials( mysqlpymysql://rfamromysql-rfam-public.ebi.ac.uk:4497/Rfam ) source sql_database(credentials).with_resources(family)直接传 SQLAlchemyEngine支持多线程可与 dlt 的parallelize设置协同from dlt.sources.sql_database import sql_table from sqlalchemy import create_engine engine create_engine(mysqlpymysql://rfamromysql-rfam-public.ebi.ac.uk:4497/Rfam) table sql_table(engine, tablechat_message, schemadata)官方建议始终把凭据放在.dlt/secrets.toml不要把敏感信息写进管道代码。4. 安装依赖运行管道前需要安装必要依赖通用依赖sql_database源所需的基础依赖pip install -r requirements.txt数据库专属依赖本教程连接 MySQL 还需安装pymysqlpip install pymysql原因说明dlt 通过 SQLAlchemy 连接源数据库因此还需要对应的 SQLAlchemy 方言例如pymysqlMySQL、psycopg2Postgres、pymssqlMSSQL、snowflake-sqlalchemySnowflake等。5. 运行管道完成第 1–4 步后即可执行python sql_database_pipeline.py该命令会在 dlt 项目目录下生成sql_to_duckdb_pipeline.duckdb文件其中包含已加载的数据。6. 探索数据dlt 内置了可交互浏览数据的浏览器仪表盘基于 Marimopip install marimo随后启动数据浏览应用dlt pipeline sql_to_duckdb_pipeline show你可以浏览已加载的数据、运行查询并查看管道执行细节。Marimo 相关的加载包查看器、管道选择器与 schema 查看器实现在 dlt/helpers/marimo/更多用法可参考 Marimo 数据探索文档。7. append、replace 与 merge三种写入策略再次执行python sql_database_pipeline.py你会发现所有表的数据都重复了一遍。这是因为 dlt 默认每次加载都以append方式向目标表追加数据。可通过pipeline.run()中的write_disposition参数调整该行为append向目标表追加数据默认行为replace用新数据替换目标表中的数据merge基于主键将新数据与目标表既有数据合并。源码佐证在 dlt/sources/sql_database/init.py#L283-L284 中sql_table的write_disposition在配置缺省时被强制回退为append。使用 replace 加载为避免数据逐行重复将write_disposition设为replaceimport dlt from dlt.sources.sql_database import sql_database def load_tables_family_and_genome(): source sql_database().with_resources(family, genome) pipeline dlt.pipeline( pipeline_namesql_to_duckdb_pipeline, destinationduckdb, dataset_namesql_to_duckdb_pipeline_data ) load_info pipeline.run(source, write_dispositionreplace) # Set write_disposition to load the data with replace print(load_info) if __name__ __main__: load_tables_family_and_genome()再次运行管道目标表中的数据将被替换而非追加。使用 merge 加载当需要随新数据更新既有数据时使用merge写入策略这要求为表指定主键主键用于将新数据与目标表既有数据匹配。在上例中pipeline.run(..., write_dispositionreplace)会让所有表都以replace加载。而通过apply_hints方法可以为每张表单独定义写入策略。下例为两张表分别指定不同的主键做 mergeimport dlt from dlt.sources.sql_database import sql_database def load_tables_family_and_genome(): source sql_database().with_resources(family, genome) # specify different loading strategy for each resource using apply_hints source.family.apply_hints(write_dispositionmerge, primary_keyrfam_id) # merge table family on column rfam_id source.genome.apply_hints(write_dispositionmerge, primary_keyupid) # merge table genome on column upid pipeline dlt.pipeline( pipeline_namesql_to_duckdb_pipeline, destinationduckdb, dataset_namesql_to_duckdb_pipeline_data ) load_info pipeline.run(source) print(load_info) if __name__ __main__: load_tables_family_and_genome()apply_hints是资源创建后修改 schema 的强力手段除了write_disposition与primary_key还可以调整incremental、table_name如表名前缀、merge_key等。从源码看显式传入的primary_key会覆盖反射结果dlt/sources/sql_database/init.py#L329-L331。merge_key则常用于按天等粒度去除重叠数据区间。提示使用merge时源表需要主键dlt 会在反射阶段自动获取主键信息也可用apply_hints(primary_key...)手动指定。8. 增量加载通常你并不想在每次加载时都拉取全量数据而只想加载新增或更新的数据。dlt 通过增量加载机制轻松实现这一点。下例为family表配置基于updated列的增量加载import dlt from dlt.sources.sql_database import sql_database def load_tables_family_and_genome(): source sql_database().with_resources(family, genome) # only load rows whose updated value is greater than the last pipeline run source.family.apply_hints(incrementaldlt.sources.incremental(updated)) pipeline dlt.pipeline( pipeline_namesql_to_duckdb_pipeline, destinationduckdb, dataset_namesql_to_duckdb_pipeline_data ) load_info pipeline.run(source) print(load_info) if __name__ __main__: load_tables_family_and_genome()首次运行python sql_database_pipeline.py时family整表被加载之后的每次运行只加载被updated列追踪的新增/更新行。增量加载的底层原理在 dlt/sources/sql_database/helpers.py#L163-L217 的BaseTableLoader._make_query()中可以看到增量逻辑如何翻译成 SQLdlt 根据dlt.sources.incremental(updated)拿到游标列名校验其存在于反射出的表中不存在会抛KeyError游标函数为max默认时生成cursor_column last_value闭区间closed或 last_value开区间open的过滤条件并可结合end_value与range_end生成区间过滤支持row_orderasc | desc控制返回行的排序便于把长周期增量加载按时间与行数切分为多个 chunk支持on_cursor_value_missinginclude | exclude决定游标值为 NULL 的行是否纳入结果。例如带initial_value的增量配置import dlt from dlt.sources.sql_database import sql_table from dlt.common.pendulum import pendulum table sql_table( tablefamily, incrementaldlt.sources.incremental( last_modified, # Cursor column name initial_valuependulum.DateTime(2024, 1, 1, 0, 0, 0) # Initial cursor value ) )首次运行会生成形如SELECT * FROM family WHERE last_modified 2024-01-01T00:00:00Z的查询后续运行中右侧被替换为 dlt 存储在 state 中的上一次游标最大值。若配置range_startopen则生成的排他比较。实用要点游标列选择优先选择时间戳列或自增 ID 列作为游标时区一致性若游标是 datetime 列注意使用带时区与不带时区的 Python 时间对象保持一致pendulum默认带时区标准datetime默认 naive避免隐式类型转换造成数据丢失建议配合full及以上反射级别获取时区提示游标列名含特殊字符如$时需转义例如incremental(example_$column, ...)sql_table参数式增量把incremental直接作为sql_table(...)的参数传入能让抽取阶段只取游标之后的行比先全量抽取再过滤更高效见模板 sql_database_pipeline.py 中load_standalone_table_resource()的写法。9. 进阶后端、模式反射与查询定制选择数据后端backendsql_database/sql_table支持四种backend见 dlt/sources/sql_database/helpers.py#L68 与配置文档sqlalchemy默认按行产出 Python 字典走常规 extract normalize 流程无需额外依赖最稳健但最慢适合小表与初期开发pyarrow产出 Arrow 表完整反射并保留源库原始类型decimal/numeric 无损目标库加载 parquet 时可直接跳过 normalizer官方文档描述可获得数量级的速度提升需要numpypandas产出 DataFrame默认使用 PyArrow dtype注意默认设置下 decimal 会映射为 double 可能损失精度、date/time 会映射为字符串——官方建议源表含 date/time/decimal 列时避免使用该后端connectorx在 Rust 中直接读取数据跳过 SQLAlchemy 行读取默认产出 Arrow 表忽略chunk_size除非设置return_typearrow_stream通常需要独立于 SQLAlchemy 的连接串用backend_kwargs{conn: ...}传入。后端由表加载器注册表驱动TABLE_LOADER_REGISTRY见 helpers.py#L406-L413你也可以通过子类化BaseTableLoader/TableLoader并调用register_table_loader_backend注册自定义后端详见高级指南。控制模式反射级别reflection_level决定从源库反射多少 schema 信息类型定义见 dlt/sources/sql_database/schema_types.py#L21minimal只反射表名、可空性与主键类型由数据推断注意 JSON/ARRAY 列在 minimal 级别下可能被跳过full默认在 minimal 基础上反射数据类型必要时对数据做类型强转full_with_precision为 decimal、text、binary 等设置精度与 scale并区分 big int 与普通 int。旧的detect_precision_hintsTrue参数已废弃等效于reflection_levelfull_with_precisionhelpers.py#L675-L689。定制查询与行级转换query_adapter_callback改写底层SELECT。例如按列过滤from dlt.sources.sql_database import sql_database def query_adapter_callback(query, table): if table.name orders: # Only select rows where the column customer_id has value 1 return query.where(table.c.customer_id1) return query source sql_database(query_adapter_callbackquery_adapter_callback).with_resources(orders)回调签名支持两种形式(query, table)或(query, table, incremental, engine)后者可用于改写增量 WHERE 条件或拼入计算列helpers.py#L219-L232。table_adapter_callback在反射后调整列集配合included_columns/excluded_columns例如把表转成带计算列的子查询再交给增量逻辑使用add_map对逐行数据做转换例如加载前对rfam_acc做哈希伪匿名化或在数据进入目标前删除敏感列详见 SQL Database 使用文档。注意 pyarrow 后端产出的是 Arrow 块而非单行字典add_map中的函数需按块处理。10. 常见问题排查速查KeyError: Cursor column ... does not exist in table ...游标列名与 SQLAlchemy 反射出的列名不一致需检查反射后的列名见排查文档重复数据确认write_disposition是否为append默认按需改为replace或merge时区/精度丢失datetime 游标注意时区一致性decimal 精度问题优先使用pyarrow后端凭据未生效确认secrets.toml中节名严格为[sources.sql_database.credentials]或改用环境变量注入。下一步恭喜完成本教程你已经掌握如何在 dlt 中配置 SQL Database 源并将数据加载到 DuckDB。接下来可以在 workspace 仪表盘 中检视管道与数据使用dataset接口访问已加载数据在 Marimo 笔记本中探索数据并生成报告深入了解 SQL Database 源参考、单表与 pyarrow/connectorx 快速后端配置、schema 与查询重写在进阶教程中学习如何创建自定义源。仓库中与本文相关的可直接查阅资源核心源码dlt/sources/sql_database/init.py、dlt/sources/sql_database/helpers.py、dlt/sources/sql_database/schema_types.py初始化模板dlt/_workspace/_templates/_core_source_templates/sql_database_pipeline.py官方文档SQL Database 教程、SQL Database 源文档。【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy ️项目地址: https://gitcode.com/GitHub_Trending/dl/dlt创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表