尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

dlt sql_database 源进阶用法:query_adapter_callback 自定义查询、计算列与加载前数据变换

dlt sql_database 源进阶用法:query_adapter_callback 自定义查询、计算列与加载前数据变换 dlt sql_database 源进阶用法query_adapter_callback 自定义查询、计算列与加载前数据变换【免费下载链接】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 认证源使用文档系统讲解sql_database/sql_table的高级定制能力如何在提取阶段通过query_adapter_callback对生成的SELECT语句施加列级过滤、改写为自定义 SQL、追加计算列并配合增量加载以及如何通过add_map在数据加载前做脱敏、删列等变换。读完本文你将掌握在 dlt 管线中查询层与数据层两个维度定制 SQL 数据接入的完整实战方案并理解其底层实现机制。一、为什么需要定制 SQL 查询sql_database源与sql_table资源在默认情况下会读取源表的全部记录。源码中查询的构建发生在 helpers.py 的BaseTableLoader._make_query()它基于反射得到的 SQLAlchemyTable对象直接生成table.select()随后在存在增量加载时附加WHERE与ORDER BY子句。# dlt/sources/sql_database/helpers.py节选 def _make_query(self) - SelectAny: table self.table query table.select() ...这意味着默认行为下每一行数据都会被提取、规范化并写入目标。若你只想加载满足特定条件的行例如某个customer_id的订单最直接的办法是在提取之前就把过滤条件下推到数据库执行。这正是query_adapter_callback参数存在的意义它接收 SQLAlchemy 的Select对象与对应的Table允许你在底层SELECT语句中注入WHERE子句让过滤在数据库端完成而非在 Python 端事后丢弃数据。从源码类型定义看该回调有两种签名见 helpers.pyTQueryAdapter Union[ Callable[[SelectAny, Table], SelectClause], Callable[[SelectAny, Table, Incremental[Any], Engine], SelectClause], ]即既支持(query, table)两参数形式也支持(query, table, incremental, engine)四参数扩展形式。调用时BaseTableLoader.make_query()会先尝试以四参数调用若抛出因参数数量不匹配导致的TypeError则回退为两参数调用见 helpers.py。二、按列过滤给 SELECT 注入 WHERE 子句最基础的应用是列级过滤。下面示例通过query_adapter_callback只加载orders表中customer_id 1的行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) # Use the original query for other tables return query source sql_database( query_adapter_callbackquery_adapter_callback ).with_resources(orders)要点说明回调对每个被加载的表都会执行因此必须对表名做判断未命中的表返回原query过滤条件使用 SQLAlchemy 核心表达式table.c.column构造dlt 会将其编译为对应数据库方言的 SQL该能力同样适用于sql_table资源且与增量加载协同工作——回调接收到的query已经包含了由_make_query()生成的增量过滤条件你在其上追加的where会与之合并。仓库测试 test_sql_database_source.py 中的test_query_adapter_callback给出了等价实践它针对chat_channel表执行query.where(table.c.active.is_(True))并分别在独立资源与源两种模式下、四种 backendsqlalchemy / pyarrow / pandas / connectorx下验证结果说明该回调对后端无关、是通用机制。三、编写自定义 SQL 查询3.1 首选方案创建数据库 VIEW当需要执行复杂查询时官方推荐先在源数据库中创建 SQL VIEW再从该视图提取数据。这样 dlt 会自动反射视图的所有列类型并按你在视图中定义的形状读取数据无需任何额外定制。3.2 扩展版 query_adapter_callback完全重写查询若无法创建视图则可以使用扩展版回调彻底重写自动生成的查询。扩展签名可以拿到incremental与engine从而自行构造带增量条件的原生 SQLimport sqlalchemy as sa def query_adapter_callback( query, table, incrementalNone, engineNone ) - TextClause: if incremental and incremental.start_value is not None: t_query sa.text( fSELECT *, 1 as add_int, const as add_text FROM {table.fullname} WHERE f {incremental.cursor_path} :start_value ).bindparams(**{start_value: incremental.start_value}) else: t_query sa.text(fSELECT *, 1 as add_int, const as add_text FROM {table.fullname}) return t_query这段代码完成了三件有意义的事情用sa.text构造原生文本查询返回值类型是TextClauseBaseTableLoader的类型别名SelectClause Union[SelectAny, TextClause]明确允许这种返回改写增量比较符把默认的ge改为严格大于即自定义增量区间语义追加计算列1 as add_int, const as add_text为每一行增加常量列此处也可以自由 JOIN 其他表。注意文本查询中的:start_value占位符必须通过.bindparams()绑定实际值避免 SQL 注入风险。3.3 为新增列显式声明类型table_adapter_callback当你在自定义查询中追加了源表不存在的列时建议用table_adapter_callback显式声明这些列的类型否则 dlt 只能从提取到的数据中推断类型from sqlalchemy.sql import sqltypes def add_new_columns(table) - None: required_columns [ (add_int, sqltypes.BigInteger, {nullable: True}), (add_text, sqltypes.Text, {default: None, nullable: True}), ] for col_name, col_type, col_kwargs in required_columns: if col_name not in table.c: table.append_column(sa.Column(col_name, col_type, **col_kwargs)) # type: ignore[arg-type]table_adapter_callback在反射阶段执行从源码看_execute_table_adapter见 helpers.py会先执行内置的default_table_adapter处理included_columns/excluded_columns过滤与 UUID 转字符串再执行你提供的自定义适配器若适配器返回了非空值例如一个子查询该返回值会替换原表对象继续参与后续流程。将两者组合起来调用sql_tableimport dlt from dlt.sources.sql_database import sql_table table sql_table( tablechat_channel, table_adapter_callbackadd_new_columns, query_adapter_callbackquery_adapter_callback, incrementaldlt.sources.incremental(updated_at), )仓库测试test_custom_sql_query见 test_sql_database_source.py对chat_channel表同时使用new_columns追加新列与自定义query_adapter并在三种 backend 下验证加载成功正是上述模式的回归保障。四、添加计算列并用于增量加载4.1 通过子查询添加计算列你可以把表转换成子查询为其追加计算列。典型场景是取多个时间列的最大值作为增量游标def add_max_timestamp(table): computed_max_timestamp sa.sql.type_coerce( sa.func.greatest(table.c.created_at, table.c.updated_at), sqltypes.DateTime, ).label(max_timestamp) subquery sa.select(*table.c, computed_max_timestamp).subquery() return subquery这里新增了max_timestamp列——它是created_at与updated_at的GREATEST结果——然后整体转换为子查询。之所以必须转为子查询是因为随后要把它用于增量加载dlt 需要在该对象上附加WHERE子句而子查询是 SQLAlchemy 中唯一可被where作用的结构。4.2 将子查询接入增量加载import dlt from dlt.sources.sql_database import sql_table read_table sql_table( tablechat_message, table_adapter_callbackadd_max_timestamp, incrementaldlt.sources.incremental(max_timestamp), )此时 dlt 会用你的子查询替代原始的chat_message表来生成增量查询。_make_query()中构造WHERE时通过self.cursor_column即table.c[max_timestamp]引用游标列所以游标列必须存在于子查询形式的表中。此外你还可以像上文那样叠加query_adapter_callback进一步定制子查询。仓库测试test_computed_column见 test_sql_database_source.py完整复现了该模式add_max_timestamp计算chat_message表created_at/updated_at的最大值随后以max_timestamp作为增量游标执行加载同样在 sqlalchemy / pandas / pyarrow 三种 backend 下通过。五、加载前的数据变换add_map除了在 SQL 层定制你还可以直接操作提取到的数据。sql_table()或sql_database().with_resources()返回的每个资源对象对应一张 SQL 表本质上是逐行产出数据的生成器。通过add_map挂载自定义 Python 函数即可在数据进入规范化与加载阶段之前逐行变换。注意 PyArrow 后端的差异PyArrow 后端不产出单行而是以ndarray形式产出成块数据。此时传给add_map的变换函数需要按ndarray输入来编写见 helpers.py 中TableLoader._convert_result的分块实现pyarrow 分支会按chunk_size分块并转成 Arrow 表。5.1 示例一PII 列伪匿名化伪匿名化是一种确定性的 PII 混淆手段允许通过哈希识别用户而不暴露底层原始信息。下面示例对family表的rfam_acc列做加盐 SHA-256 哈希import dlt import hashlib from dlt.sources.sql_database import sql_database def pseudonymize_name(doc): Pseudonymization is a deterministic type of PII-obscuring. Its role is to allow identifying users by their hash, without revealing the underlying info. # add a constant salt to generate salt WIN57%zZrmk#88c salted_string doc[rfam_acc] salt sh hashlib.sha256() sh.update(salted_string.encode()) hashed_string sh.digest().hex() doc[rfam_acc] hashed_string return doc pipeline dlt.pipeline( # Configure the pipeline ) # using sql_database source to load family table and pseudonymize the column rfam_acc source sql_database().with_resources(family) # modify this source instances resource source.family.add_map(pseudonymize_name) # Run the pipeline. For a large db this may take a while info pipeline.run(source, write_dispositionreplace) print(info)add_map直接作用在源实例的命名资源上source.family对资源内每条记录执行变换函数并返回修改后的字典。关于dlt的伪匿名化最佳实践可进一步阅读 伪匿名化列指南位于docs/website/docs/general-usage/customising-pipelines/目录。5.2 示例二加载前剔除多余列如果某些列完全不需要落入目标端与其加载后再处理不如在资源层直接删除import dlt from dlt.sources.sql_database import sql_database def remove_columns(doc): del doc[rfam_id] return doc pipeline dlt.pipeline( # Configure the pipeline ) # using sql_database source to load family table and remove the column rfam_id source sql_database().with_resources(family) # modify this source instances resource source.family.add_map(remove_columns) # Run the pipeline. For a large db this may take a while info pipeline.run(source, write_dispositionreplace) print(info)提示若只是单纯不需要某些列更高效的做法是在配置层面用included_columns/excluded_columns直接控制反射出的列集合见 configuration.md避免数据被提取后才丢弃add_map更适合需要复杂逻辑哈希、脱敏、拼装新字段的场景。六、部署 sql_database 管线sql_database管线可使用 dlt 提供的任意部署方式运行包括 GitHub Actions、Airflow、Dagster 等完整列表见 部署指南 目录。在 Airflow 上运行在 Airflow 上运行时有三个关键实践使用 dlt Airflow Helper 构建任务通过 Airflow Helper 文档 将sql_database源转换为 DAG 任务。若希望多表提取并行执行可在源到 DAG 转换时设置decompose parallel-isolated运行时反射表结构设置defer_table_reflectTrue。源码 __init__.py 中明确该参数会在产出数据时才连接并反射表结构但要求必须显式传入table_names否则抛出ValueError。这正适合 Airflow 等编排器先构建执行 DAG、后执行任务的场景——注意启用后模式schema在执行期才确定可能覆盖query_adapter_callback的修改或apply_hints配合调度区间做增量设置allow_external_schedulers以便借助 Airflow 的调度区间schedule interval实现回填与增量加载参见 增量游标文档。七、底层实现一览回调何时被调用理解回调的调用时机有助于写出正确的适配器。以sql_table资源为例完整调用链如下sql_table()init.py构建 SQLAlchemyEngine、MetaData反射或按需延迟反射Tabletable_adapter_callback在反射后立即执行_execute_table_adapter此时可增删列、换类型或返回子查询资源被迭代时table_rows()helpers.py依据 backend 查表加载器注册表TABLE_LOADER_REGISTRYsqlalchemy / pyarrow / pandas 共用TableLoaderconnectorx 用ConnectorXTableLoader加载器load_rows()中先调用make_query()内部先由_make_query()生成基础查询含增量过滤与排序再调用query_adapter_callback做最终改写改写后的查询Select或TextClause由各后端执行SQLAlchemy 后端按chunk_size分块转成字典列表 / DataFrame / Arrow 表ConnectorX 后端则编译为 SQL 字符串交由 Rust 执行。该分层设计BaseTableLoader抽象基类 可注册的自定义后端意味着query_adapter_callback对全部四种后端生效——测试文件 test_sql_database_source.py 中对query_adapter_callback、test_computed_column、test_custom_sql_query均做了多 backend 参数化验证。八、小结本文围绕sql_database/sql_table的查询定制与数据变换展开列级过滤query_adapter_callback(query, table)注入WHERE过滤在数据库端完成自定义 SQL扩展四参签名(query, table, incremental, engine)可完全重写查询为sa.text原生 SQL并配合table_adapter_callback显式声明新增列类型计算列table_adapter_callback返回子查询追加计算列如max_timestamp可直接作为增量游标行级变换add_map在加载前逐行PyArrow 后端为逐块ndarray完成脱敏、删列等处理部署Airflow 场景下使用defer_table_reflect、decomposeparallel-isolated与allow_external_schedulers。相关的完整配置backend 选型、连接串、增量配置可继续阅读 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),仅供参考
返回列表