ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

dlt 1.18 版本亮点解读:命名目标、Databricks Iceberg 与流水线可观测性增强

2026/9/18 16:29:39 拓冰建站 浏览量
dlt 1.18 版本亮点解读:命名目标、Databricks Iceberg 与流水线可观测性增强 dlt 1.18 版本亮点解读命名目标、Databricks Iceberg 与流水线可观测性增强【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy ️项目地址: https://gitcode.com/GitHub_Trending/dl/dltdltdata load tool1.18 是一次聚焦「配置灵活性与生产级可靠性」的版本更新它引入了可按名称引用的命名目标Named Destinations与dlt.destination工厂让开发、预发、生产环境的切换不再依赖修改代码同时为 Databricks 目标补齐了 Iceberg 表格式与分区/聚簇 hint 的真正落地并重构了信号处理实现优雅停机、新增add_metrics()自定义指标采集、支持从 HTTP URL 直接读取文件。读完本文你将掌握上述每一项新能力的配置方式、代码用法与其背后的源码实现可直接在真实 pipeline 中落地。按名称引用目标Named Destinations 与dlt.destination工厂在 1.18 之前pipeline 的目标通常在代码中直接指定如destinationduckdb环境切换往往意味着改动代码或依赖环境变量注入。1.18 引入的命名目标机制让你可以在 workspace 配置文件.dlt/config.toml中为目标起一个自定义名称代码中只引用这个名字目标类型与相关设置由 dlt 自动从配置解析。配置方式在项目根目录的.dlt/config.toml中加入一个以[destination.自定义名称]为前缀的段并用destination_type字段声明实际的目标类型[destination.custom_name] destination_type duckdb代码中使用import dlt dlt.resource def my_data(): yield {id: 1} pipeline dlt.pipeline(destinationcustom_name) pipeline.run(my_data)pipeline.run(my_data)时dlt 会读取配置段[destination.custom_name]解析出真实的destination_type此处为 duckdb并实例化对应目标段内的其他键如凭据、数据集名等也会一并作为目标配置生效。这样一来切换环境只需修改配置文件预发环境把destination_type换成postgres、生产环境换成databrickspipeline 代码保持不动。源码层面的机制从源码结构看命名目标与自定义目标工厂共享同一套配置规范dlt/destinations/impl/destination/factory.py中的destination类Destination[CustomDestinationClientConfiguration, DestinationClient]内部维护destination_type属性并通过dlt/common/reflection/ref.py的object_from_ref与ImportTrace机制解析类型引用当引用无法导入时会抛出UnknownCustomDestinationCallable异常并列出尝试过的完整限定名便于排查配置笔误。配置段的解析复用dlt.common.configuration的with_config与get_fun_spec体系因此命名目标同样支持凭据注入如dlt.secrets.value与分层配置覆盖。Databricks 目标Iceberg 表格式与分区/聚簇 hint 全面落地1.18 之前databricks_adapter传入的分区partitioning与聚簇clusteringhint 会被静默忽略从 1.18 起这些 hint 被真正强制执行且无效组合例如将分区与聚簇混用在同一张表上会在早期被拒绝避免生成错误的表布局。完整示例import dlt from dlt.destinations.adapters import databricks_adapter pipeline dlt.pipeline(pipeline_nameto_databricks, destinationdatabricks) dlt.resource def my_data(): for i in range(10): yield { event_id: i, year: 2024, month: (i % 3) 1, customer_id: i % 5, } # Enable clustering databricks_adapter(my_data, clusterAUTO) # Enable partitioning databricks_adapter(my_data, partition[year, month]) # Use Iceberg table format databricks_adapter(my_data, table_formatICEBERG) pipeline.run(my_data)适配器参数详解databricks_adapter定义于dlt/destinations/impl/databricks/databricks_adapter.py其核心参数如下参数类型说明clusterUnion[TColumnNames, Literal[AUTO]]列名、列名列表或AUTO由 Databricks 自动决定最优聚簇键partitionTColumnNames列名或列名列表按列值把表拆分为独立文件table_formatLiteral[DELTA, ICEBERG]表格式默认DELTA设为ICEBERG可享受更好的 schema 演进与 time travel 能力table_commentOptional[str]表描述table_tagsList[Union[str, Dict[str, str]]]表标签可混合字符串与键值对如[production, {environment: prod}]table_propertiesOptional[Dict[str, ...]]通过TBLPROPERTIES写入的表属性如{delta.appendOnly: True}column_hintsTTableSchemaColumns列级 hint支持column_comment、column_tagsinsert_apiOptional[TDatabricksInsertApi]append写入方式的后端copy_into或zerobus其中table_format若取值不是DELTA或ICEBERG适配器会直接抛出ValueError表格格式 hint 会以小写形式delta/iceberg存储到表级 hint 中与目标端 DDL 生成逻辑对齐。DDL 生成原理在dlt/destinations/impl/databricks/databricks.py中建表 SQL 的组装逻辑会依次检查聚簇与分区 hint表级 hintx-databricks-cluster为AUTO时生成CLUSTER BY AUTO子句否则按列 hint 收集聚簇列生成CLUSTER BY (col1,col2)列 hint 中partition: True的列会收集为PARTITIONED BY (col1,col2)table_format iceberg时使用USING ICEBERG子句并会拦截与 Iceberg 不兼容的表属性报错提示移除相关属性后再使用table_formaticeberg。最终这些子句被拼接进建表 SQL对已存在表的ALTER场景若配置了聚簇子句则生成ALTER TABLE ... CLUSTER BY ...追加聚簇。由此可见「分区与聚簇互斥」的校验发生在 hint 合法性检查与 SQL 生成两个层面确保表布局不被错误组合污染。流水线优雅停机信号处理全面重构长时间运行或生产环境的 pipeline 最怕被CtrlC粗暴打断导致中间状态损坏。1.18 重做了 dlt 的信号处理中断行为变为「优雅关机」第一次CtrlC触发干净的 shutdown正在进行的抽取/加载任务会尽力排空drain后再退出第二次中断强制立即退出。该行为同时适用于基于进程process与基于线程threaded两种 loader生产环境运行的可靠性显著提升。实现机制核心实现位于dlt/common/runtime/signals.py的_signal_receiver采用「两阶段升级」策略通过intercepted_signals()安装处理器覆盖SIGINT、SIGTERM等 POSIX 信号第一次收到某类信号时_signal_counts[sig] 1处理器设置优雅关闭标志并唤醒等待中的线程同时用信号处理器安全的低级日志函数打印「收到信号正在优雅关闭可能需要时间排空」之类的提示第二次收到同类信号时_signal_counts[sig] 2处理器委托给原始处理器若是SIG_DFL则恢复默认行为并重新raise_signal(sig)从而触发立即终止如KeyboardInterrupt。注释中特别说明CPython 只在主线程分发信号处理器工作线程必须通过raise_if_signalled()或信号感知的sleep()协作式地观察关闭标志——这正是「第一次优雅、第二次强制」能在多线程 loader 下稳定工作的关键约定。更详细的运维建议可参考仓库文档 running-in-production。用add_metrics()采集自定义指标以往要统计批次数、API 页数这类抽取层信号要么侵入数据本身要么额外写旁路计数器。1.18 提供add_metrics()可直接把自定义指标采集器挂到 resource 上数据项原样流过、不被修改。示例统计 resource 产生了多少个批次import dlt def batch_counter(items, meta, metrics): metrics[batch_count] metrics.get(batch_count, 0) 1 dlt.resource def purchases(): for i in range(3): yield [{id: i}] purchases purchases.add_metrics(batch_counter) pipeline dlt.pipeline(metrics_demo, destinationduckdb) load_info pipeline.run(purchases) trace pipeline.last_trace load_id load_info.loads_ids[0] resource_metrics trace.last_extract_info.metrics[load_id][0][resource_metrics][purchases] # type: ignore print(Custom metrics:, resource_metrics.custom_metrics)输出Custom metrics: {batch_count: 3}自定义指标与性能、转换统计一起存放在resource_metrics之下可在 trace 与 dashboard 中查看。底层实现add_metrics定义于dlt/extract/resource.py签名是add_metrics(metrics_f: MetricsFunctionWithMeta, insert_at: int None)默认把采集器追加到 resource 管道pipe末尾也可用insert_at指定插入位置每个指标采集函数被包装成MetricsItem见dlt/extract/items_transform.py其__call__调用self._metrics_f(item, meta, self._custom_metrics)后原样返回 item因此对数据流是零侵入的每个 resource 实例维护_custom_metrics字典dlt/extract/resource.py中初始化custom_metrics属性可读同时MetricsItem也持有自己的_custom_metrics抽取阶段dlt/extract/extract.py会调用_get_all_resource_custom_metrics(resource_name)把 resource 级与每个 pipe 步骤step上的自定义指标用update_dict_nested深度合并随ExtractInfo写入 trace在dlt/extract/decorators.py中还有get_resource().custom_metrics的便捷访问入口方便在资源装饰器内部直接读取。一个值得注意的细节示例中每yield一个列表一个批次触发一次计数3 个批次得到batch_count: 3——说明指标函数作用于「流经 pipe 的数据项/批次」与add_map、add_filter等转换步骤共用同一管道模型语义清晰且可组合。HTTP 文件系统资源直接读取公开数据集filesystem源此前主要面向本地路径与 S3 等对象存储1.18 起它可以直接从HTTP URL读取文件基于 fsspec 的能力扩展加载公开托管的数据集不再需要额外搭建基础设施。示例import dlt from dlt.sources.filesystem import filesystem, read_csv pipeline dlt.pipeline(metrics_demo, destinationduckdb) fs filesystem(bucket_urlhttps://example.com/data/, file_glob*.csv) pipeline.run(fs | read_csv())相关说明filesystem源由dlt/sources/filesystem/__init__.py导出底层通过dlt/common/storages/fsspec_filesystem.py的fsspec_filesystem、glob_files与FileItem系列类型对接 fsspec 协议族。源函数支持bucket_url、credentials、file_glob、client_kwargs、incremental增量加载等参数对 HTTP 这种「每次文件获取都需要一次服务端调用」的协议配置项requires a server call per file默认为False仅在http文件系统场景下需要开启细节可查阅dlt/sources/filesystem/settings.py与源函数 docstring。组合用法filesystem(...) | read_csv()表明它仍是标准 source/resource 管道可继续叠加read_jsonl、read_parquet等其他 reader或进一步接入增量参数实现按文件变更的周期加载。结语与社区除上述功能外1.18 还包含若干质量改进与 bug 修复并向多位新贡献者致谢覆盖 3000 号段左右的多个 PR。从整体看本版本的主线是「让目标配置更灵活、让表布局更可控、让运行更可观测」命名目标把环境切换成本降到配置文件级别Databricks 的 Iceberg/分区/聚簇让数仓表设计落到 dlt 层优雅停机与add_metrics则分别补强了生产运行的可靠性观测面。建议在升级后对照dlt/version.py中的版本号与docs/website/docs/release-notes/目录下的完整变更记录核对行为变化尤其是涉及 Databricks hint 的既有 pipeline 可能需要重新验证建表 SQL 是否符合预期。【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy ️项目地址: https://gitcode.com/GitHub_Trending/dl/dlt创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考