拓冰建站拓冰建站
首页 / 资讯中心 / 正文

dlt 分块回填(Backfill in Chunks):用 sql_database 源按指定块大小分批加载 MySQL 数据

dlt 分块回填Backfill in Chunks用 sql_database 源按指定块大小分批加载 MySQL 数据【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy ️项目地址: https://gitcode.com/GitHub_Trending/dl/dlt本文基于 dlt 官方示例文档 backfill_in_chunks.md 展开并结合仓库内dlt/sources/sql_database的源码实现进行深入讲解。你将学会如何用sql_database源连接 MySQL 并只加载指定表、如何通过apply_hints声明主键与增量游标、如何用chunk_size与add_limit把一次超大回填拆成多次小规模 pipeline run以及如何借助pipeline.dataset()的 datasets 访问器断言加载进度。读完本文你可以把一次跑几个小时、容易在临时存储上 OOM 的大回填改造成每次只加载固定块、随时能看到目标端进度的可观测分批回填。1. 为什么要分块回填一次性大回填的三大痛点把历史数据从源数据库例如 MySQL全量回填到目标端例如 DuckDB时最直接的做法是一次 pipeline run 加载整张表。但这种方式在实践中会踩到三个典型问题内存压力整表数据一次性进入 extract/normalize/load 链路在临时存储或内存受限的环境如无服务器函数、ephemeral 容器中极易 OOM 或写满磁盘。回填通常意味着很长时间没有结果一旦失败就要从头再来。无可观测进度单次大任务在完成前目标端看不到任何中间结果无法判断任务是否健康、还需多久。失败成本高中间任意一步出错整批数据都要重跑且很难定位问题发生在哪一段数据上。dlt 给出的解法正是本文的主题把回填拆成多个固定大小的块chunk每个 pipeline run 只加载一个块。每次运行都能在目标端立即看到新增的一批数据块与块之间由增量游标incremental cursor衔接天然支持断点续传。从源码结构看这一能力由三层协作实现sql_database/sql_table源负责按chunk_size产出数据批次见 dlt/sources/sql_database/init.pyadd_limit负责限制单次运行产出的块数见 dlt/extract/source.py 与 dlt/extract/items_transform.py 中的LimitItem而Incremental负责跨运行记录游标位置见 dlt/extract/incremental/init.py。2. 前置准备依赖与连接串2.1 安装依赖示例使用sql_database源读取 MySQL并用到 pandas 读取目标端数据做断言因此至少需要pip install dlt[duckdb] dlt[sql_database] pandas其中dlt[sql_database]会带入 SQLAlchemy 及对应数据库驱动连接 MySQL 需要pymysql或mysqlclient。DuckDB 作为目标端由dlt[duckdb]提供pandas 用于把目标表读回 DataFrame 做唯一性断言。2.2 连接串格式sql_database源接受 SQLAlchemy 风格连接串。示例中使用的是 EMBL-EBI 公开的 Rfam 数据库只读、无密码RFAM_CONNECTION_STRING mysqlpymysql://rfamromysql-rfam-public.ebi.ac.uk:4497/Rfam连接串的通用格式为dialectdriver://user:passwordhost:port/database。sql_database的credentials参数除字符串外也接受ConnectionStringCredentials对象或现成的sqlalchemy.Engine实例见 dlt/sources/sql_database/init.py#L38-L58。生产环境建议把连接串放进.dlt/secrets.toml或环境变量而不是写死在脚本里当不传credentials时dlt 会自动从配置系统解析默认值为dlt.secrets.value。3. 创建 sql_database 源只选一张表 设置块大小3.1 最小调用import dlt from dlt.sources.sql_database import sql_database source sql_database( RFAM_CONNECTION_STRING, table_names[family], chunk_size1000, )三个关键参数的含义如下credentials数据库连接信息字符串 / 凭据对象 / Engine。table_names[family]只加载family一张表。不传该参数时源会反射reflect整个 schema 下所有表并分别为每张表创建资源——对于只需要回填一张表的大库来说这会白白付出大量反射开销因此显式传table_names是分块回填的标准做法。chunk_size1000每批产出的行数。源码中默认值为50000见 dlt/sources/sql_database/init.py#L44此处调小为 1000 是为了演示每块 1000 行的效果。3.2 源码视角chunk_size 如何影响底层行为阅读sql_database源码可以发现chunk_size不只是一批多少行那么简单它直接影响 SQLAlchemy 引擎的流式读取配置engine.execution_options(stream_resultsTrue, max_row_buffer2 * chunk_size)见 dlt/sources/sql_database/init.py#L134。stream_resultsTrue让查询结果以流式游标返回而非一次性载入内存max_row_buffer2 * chunk_size则限制了内部缓冲的行数。也就是说chunk_size在源层面同时决定了产出的批次大小和内存中的缓冲上限——这正是它能控制回填内存占用的原因之一。在表加载器层面dlt/sources/sql_database/helpers.pyTableLoader._load_rows执行查询时也会设置yield_perself.chunk_size并在_convert_result中通过result.partitions(sizeself.chunk_size)把游标结果切分成若干批次逐批 yield见 dlt/sources/sql_database/helpers.py#L300-L347。因此chunk_size贯穿引擎缓冲 → 查询取数 → 批次产出整条链路。另外sql_database还支持backend参数sqlalchemy默认、pyarrow、pandas、connectorx用于指定批次的数据形态。其中sqlalchemy每批产出 list of dictpyarrow/connectorx产出 Arrow 表pandas产出 DataFrame。注意源码中特别注明connectorx后端会忽略chunk_size见 dlt/sources/sql_database/init.py#L77大表需要自行处理因此分块回填场景建议使用默认sqlalchemy后端。4. 用 apply_hints 声明主键与增量游标4.1 为什么必须声明这两个 hint分块回填的核心是每次运行从上次停下的位置继续。要做到这一点dlt 需要知道两件事按什么列排序取数——增量游标列cursor column示例中为created时间列保证每块数据是按时间升序的一段连续区间。块与块之间如何去重——主键列示例中为rfam_id。由于增量游标采用闭区间策略range_start默认为closed相邻两块会在游标值上存在少量重叠示例中第二块得到 1999 行而非 2000 行正是这个原因主键让 dlt 能对重叠行做去重确保不丢数据。4.2 示例中的 apply_hints 调用source.family.apply_hints( primary_keyrfam_id, incrementaldlt.sources.incremental( cursor_pathcreated, initial_valueNone, row_orderasc, ), )逐项说明source.familysql_database生成的源会为每张表创建一个同名的资源resource这里直接取family表对应的资源。primary_keyrfam_id声明主键。在sql_table资源中主键会被写入表 hints加载时用于唯一性去重见 dlt/sources/sql_database/init.py#L329-L331其中primary_key最终被并进资源 hints。incrementaldlt.sources.incremental(...)启用增量加载通过apply_hints绑定到资源上apply_hints支持incremental参数见 dlt/extract/hints.py#L441-L457。dlt.sources.incremental返回的Incremental实例支持以下核心参数见 dlt/extract/incremental/init.py#L125-L148参数默认值说明cursor_path必填游标字段名JSON 路径必须是表中的一个简单列名initial_valueNone首次运行、state 中无记录时使用的初始游标值None表示从表头开始last_value_funcmax决定游标推进函数可选max/min或自定义 callableprimary_keyNone用于去重的键此处由apply_hints另行声明row_orderNone声明数据源按asc/desc有序返回用于提前停止拉取end_valueNone与initial_value配合加载固定区间如某个月设置后为无状态过滤on_cursor_value_missingraise游标值缺失/为 None 时的行为raise/include/excluderange_startclosed过滤区间起点是否闭区间open会跳过与上次游标值相同的行并关闭去重逻辑range_endopen过滤区间终点是否闭区间4.3 源码视角增量游标如何生成 SQL在表加载器BaseTableLoader._make_query中dlt/sources/sql_database/helpers.py#L163-L217增量游标会被编译为 SQL 的 WHERE 与 ORDER BY 子句游标值从 state 中取出self.last_value, self.end_value incremental.get_current_range()见 dlt/sources/sql_database/helpers.py#L149使用max游标函数时生成cursor_column last_value闭区间之类的过滤条件按row_order生成ORDER BY cursor_column ASC/DESC保证块内顺序确定。这说明增量过滤发生在 SQL 层而非 Python 层每个块只从源库拉取大于等于上次游标值的若干行效率很高。同时注意游标列必须能直接映射为表的一列否则BaseTableLoader会抛出KeyError见 dlt/sources/sql_database/helpers.py#L142-L148。5. 用 add_limit 限制单次运行的块数5.1 关键一行source.add_limit(1)source.add_limit(1)注释中写得很清楚块大小 1000、limit 为 1意味着每次 pipeline run 加载 1000 行。add_limit是DltSource的方法会把 limit 应用到源中所有已选中的非 transformer 资源上见 dlt/extract/source.py#L526-L555。5.2 源码视角add_limit 与 chunk_size 的乘法关系add_limit的本质是在资源管道中插入一个LimitItem变换器见 dlt/extract/resource.py#L453-L489。关键实现在LimitItem.limit()dlt/extract/items_transform.py#L197-L206def limit(self, chunk_size: int) - Optional[int]: if self.max_items in (None, -1): return None return self.max_items * (1 if self.count_rows else chunk_size)也就是说add_limit(N)限制的是产出批次数而不是行数。表加载器在构建查询时调用limit self.limit.limit(self.chunk_size)见 dlt/sources/sql_database/helpers.py#L168-L172把它换算成 SQL 的LIMIT max_items * chunk_size。所以chunk_size1000add_limit(1)→ 每次运行 SQLLIMIT 1000加载恰好 1000 行chunk_size1000add_limit(5)→ 每次运行加载 5000 行5 个块不调用add_limit→ 单次运行加载整表退化为一次性大回填。add_limit还支持max_time按运行秒数截断常用于流式/长任务保护和count_rowsTrue按行数而非批次计数注意行计数模式下最后一个批次不会被裁剪可能多出部分行见 dlt/extract/items_transform.py#L208-L233。此外要留意LimitItem的placement_affinity 1.1见 dlt/extract/items_transform.py#L175它被设计为紧跟在 incremental 变换器之后执行——也就是说先按游标过滤/排序再按 limit 截断顺序保证每个块都是按游标排序后的前 N 行。6. 创建 Pipeline 并分块回填6.1 Pipeline 定义pipeline dlt.pipeline( pipeline_namerfam, destinationduckdb, dataset_namerfam_data, dev_modeTrue, )pipeline_namerfampipeline 的名字。state含增量游标就是按 pipeline 名持久化的多次运行必须使用同一个名字增量才能衔接。destinationduckdb目标端为本地 DuckDB无需额外配置即可运行演示生产环境可换成 postgres、bigquery、snowflake、filesystem 等任意 dlt 支持的目标端。dataset_namerfam_data目标 schema/数据集名。dev_modeTrue每次运行前重建 dataset先删后建。注意dev_mode 只适合演示/测试用于回填的正式 pipeline 应去掉它否则每次运行都会清空目标数据。6.2 定义唯一性断言函数def _assert_unique_row_count(df: pd.DataFrame, num_rows: int) - None: Assert that a dataframe has the correct number of unique rows assert len(df) num_rows assert len(set(df.rfam_id.tolist())) num_rows该函数同时检查两件事行数正确、且rfam_id无重复。它通过pipeline.dataset().family.df()把目标表整体读回内存再统计因此示例源码注释也特别提醒这种校验方式只适合小表/测试阶段真正做大规模回填前验证用不要把整表读回内存作为回填后的常规校验手段。6.3 完整回填流程# 第一次运行目标端 family 表应恰好包含前 1000 行 pipeline.run(source) _assert_unique_row_count(pipeline.dataset().family.df(), 1000) # 第二次运行应为 1999 行增量游标闭区间导致 1 行重叠用于防止丢行 pipeline.run(source) _assert_unique_row_count(pipeline.dataset().family.df(), 1999) # ... pipeline.run(source) _assert_unique_row_count(pipeline.dataset().family.df(), 2998) # ... pipeline.run(source) _assert_unique_row_count(pipeline.dataset().family.df(), 3997) # 最后一次运行加载到表尾行数应为全表行数 pipeline.run(source) _assert_unique_row_count(pipeline.dataset().family.df(), TOTAL_TABLE_ROWS)其中TOTAL_TABLE_ROWS 4178是对应的全表行数常量。观察行数序列1000 → 1999 → 2998 → 3997 → 4178可以发现两个关键事实前四次每次净增约 9991000 行增量游标从上次位置继续块与块之间无缝衔接第 2 块开始比理想值少 1 行相邻块的游标区间存在重叠闭区间设计重叠行由主键去重换取任何情况下都不会因边界值而丢行的保证最后一次第 5 次运行只剩 181 行limit 是上限而非固定值当源表剩余行数不足一个块时加载到表尾自然结束Incremental感知到游标到达末端后后续运行会变成空跑不再产出数据。示例末尾的注释给出了生产环境的迁移建议见 docs/website/docs/examples/backfill_in_chunks.md实际回填时chunk_size和limit通常会大得多例如块 50 万行用循环反复调用pipeline.run(source)而不是手写多次通过程序化方式判断表是否已完全加载例如每次运行后检查pipeline.dataset().family的行数/游标状态加载完毕即跳出循环。7. 用 datasets 访问器检查加载进度7.1pipeline.dataset()是什么pipeline.dataset()返回一个dlt.Dataset对象用于查询目标端已落库的数据见 dlt/pipeline/pipeline.py#L2184-L2194。示例中的调用链是pipeline.dataset().family.df()pipeline.dataset()定位到dataset_namerfam_data对应的数据集.family数据集中的family表对应的关系对象Relation.df()把该表读成 pandas DataFrame见 dlt/dataset/relation.py#L183-L185。Dataset访问器是回填进度的第一观测窗口每次pipeline.run完成后立即用它读取目标表即可确认这一块数据是否真的写进去了、行数是否符合预期。这是分块回填相比一次性大回填在可观测性上的核心优势。7.2 其他可用的读取接口除了df()Dataset.Relation还提供arrow()读成 Arrow 表、fetchall()/fetchmany()/fetchone()读成元组列表/单行、iter_df()/iter_arrow()流式迭代等方法见 dlt/dataset/relation.py#L183-L214。在生产环境的回填循环中可以用更轻量的方式例如只查询COUNT(*)或游标列的MAX值来判断是否已加载完成避免整表读回内存。8. 生产级分块回填模板把示例改造成适合生产环境的循环版本核心骨架如下import dlt from dlt.sources.sql_database import sql_database TOTAL_TABLE_ROWS 4178 # 若为动态表需每次运行后从目标端重新确认 CONNECTION_STRING mysqlpymysql://user:passhost:3306/db source sql_database( CONNECTION_STRING, table_names[family], chunk_size100000, # 生产环境通常远大于 1000 ) source.family.apply_hints( primary_keyrfam_id, incrementaldlt.sources.incremental( cursor_pathcreated, initial_valueNone, row_orderasc ), ) source.add_limit(1) # 每次运行加载 1 个块 pipeline dlt.pipeline( pipeline_namerfam_backfill, destinationpostgres, # 或其他目标端 dataset_namerfam_data, ) # 循环回填每轮加载一个块直到目标端行数达到全表行数 while True: info pipeline.run(source) print(info) rows pipeline.dataset().family.df().shape[0] # 生产环境建议改用 COUNT(*) 查询 if rows TOTAL_TABLE_ROWS: break几点生产注意事项去掉dev_modeTrue否则每轮运行都会重建 dataset回填永远无法累积块大小与 limit 的权衡块太小则运行次数过多、调度开销大块太大则单次运行内存/耗时回升。建议结合目标端写入性能实测调整幂等与断点续传增量游标存在 pipeline state 中任意一轮失败后修复问题直接重跑即可从断点继续重复块由主键去重不要在循环内用df()做全表校验读回整表会抵消分块带来的内存收益生产环境用SELECT COUNT(*)或游标MAX(created)判断完成度Airflow 等编排场景可把每轮回填封装为一个 task/DAG run如需延迟表反射可参考sql_database的defer_table_reflectTrue参数要求必须显式传table_names见 dlt/sources/sql_database/init.py#L140-L143。9. 小结分块回填的三个杠杆回顾本文dlt 的分块回填本质上是三个参数的组合拳参数作用示例值chunk_size控制每批/每块的行数同时决定引擎缓冲与 SQLLIMIT的粒度1000演示/100000生产add_limit(N)限制单次 pipeline run 产出的块数与chunk_size相乘得到单次 SQL 的LIMIT1incremental游标 主键跨运行记忆游标位置保证块间无缝衔接、重叠去重、断点续传cursor_pathcreated,primary_keyrfam_id再配合pipeline.dataset()在每次运行后核对目标端行数你就能把一次跑数小时、失败全重来的大回填变成每次只搬一块、每块都可见、随时可续传的稳健流程。示例完整源码与说明位于 docs/website/docs/examples/backfill_in_chunks.md相关实现可继续阅读 dlt/sources/sql_database/init.py、dlt/sources/sql_database/helpers.py 与 dlt/extract/incremental/init.py。【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy ️项目地址: https://gitcode.com/GitHub_Trending/dl/dlt创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门