ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

dlt 使用 filesystem destination 加载 Delta Lake 表:table_format、分区、merge 与存储配置详解

dlt 使用 filesystem destination 加载 Delta Lake 表:table_format、分区、merge 与存储配置详解 dlt 使用 filesystem destination 加载 Delta Lake 表table_format、分区、merge 与存储配置详解【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy ️项目地址: https://gitcode.com/GitHub_Trending/dl/dlt本文基于 dlt 官方文档docs/website/docs/dlt-ecosystem/destinations/delta-iceberg.md并结合当前仓库源码系统讲解如何通过 dlt 的filesystemdestination 写出 Delta Lake 表从依赖安装、table_formatdelta的设置方式、Hive 风格分区、get_delta_tables表访问助手、事务性截断到mergeupsert / insert-only策略与deltalake_storage_options、deltalake_configuration等底层配置的透传机制最终给出可在生产环境中直接复制运行的完整配置与代码示例。Delta 表格式在 dlt 中的定位dlt 支持在使用 filesystem destination 时写出 Delta 表。与直接写 Parquet/CSV 不同Delta 表在 Parquet 文件之上附加了事务日志提供 ACID 语义、schema 演进与时间旅行能力。在仓库中Delta 能力由三个部分协作实现dlt/common/libs/deltalake.py对deltalake库的封装包含write_delta_table、merge_delta_table、truncate_delta_table、get_delta_tables、deltalake_storage_options等核心函数dlt/destinations/impl/filesystem/factory.py在 destination 能力capabilities层面声明supported_table_formats[delta, iceberg]并约束文件格式与 merge 策略dlt/destinations/impl/filesystem/filesystem.pyDeltaLoadFilesystemJob负责在 load 阶段真正执行写入或 merge。工作原理dlt 使用 deltalake 库来写 Delta 表。其数据流为extract 与 normalize 阶段数据被准备成一个或多个 Parquet 文件load 阶段这些 Parquet 文件被暴露为 Arrow 数据结构pyarrow.dataset/RecordBatchReader再交给deltalake写入 Delta 表。从源码可以确认这一流程。DeltaLoadFilesystemJob.run()位于 filesystem.py其核心逻辑是# dlt/destinations/impl/filesystem/filesystem.py节选 with get_local_dataset_reader(self.file_paths) as arrow_rbr: # RecordBatchReader if self._load_table[write_disposition] merge and delta_table is not None: merge_delta_table(...) else: location self._job_client.get_open_table_location(delta, self.load_table_name) write_delta_table( table_or_urilocation if delta_table is None else delta_table, dataarrow_rbr, write_dispositionself._load_table[write_disposition], partition_byself._partition_columns, storage_optionsstorage_options, configurationself._job_client.config.deltalake_configuration, )也就是说dlt 先把 Parquet 文件读成 Arrow 的RecordBatchReader流式读取避免一次性加载全表到内存然后若表已存在且write_disposition merge走merge_delta_table执行 upsert / insert-only否则走write_delta_table由 dlt 根据表的写入语义决定是 append 还是 overwrite。Delta 依赖使用 Delta 表格式需要deltalake包pip install dlt[deltalake]同时还需要pyarrow17.0.0pip install pyarrow17.0.0关于版本要求仓库中有两处佐证pyproject.toml 中deltalake额外依赖声明为deltalake0.25.1与pyarrow16.0.0但 dlt/common/libs/deltalake.py 的ensure_delta_compatible_arrow_data内部会执行assert_min_pkg_version(pkg_namepyarrow, version17.0.0, ...)——因为对RecordBatchReader的cast()需要 pyarrow 17 才可用16 中的cast()存在 bug。因此实际使用delta表格式时pyarrow 必须达到 17.0.0 及以上文档中的要求与源码行为一致。此外如果未安装deltalake导入 dlt/common/libs/deltalake.py 时会抛出MissingDependencyException并提示安装dlt[deltalake]。设置 table_format在定义 resource 时将table_format参数设置为deltadlt.resource(table_formatdelta) def my_delta_resource(): ...也可以在调用 pipeline 的run时传入pipeline.run(my_resource, table_formatdelta)注意使用delta表格式时dlt 会始终强制使用 Parquet 作为loader_file_format任何对loader_file_format的设置都会被忽略。这一点在 factory.py 中由filesystem_loader_file_format_selector实现def filesystem_loader_file_format_selector( preferred_loader_file_format, supported_loader_file_formats, /, *, table_schema ): if table_schema.get(table_format) in (delta, iceberg): return (parquet, [parquet]) # delta/iceberg 强制 parquet return (preferred_loader_file_format, supported_loader_file_formats)即只要表 schema 上标注了table_format为delta或icebergloader 文件格式就被锁定为 parquet测试用例 tests/load/pipeline/test_open_table_pipeline.py 中也会验证 load job 的产物文件均以.parquet结尾。表分区PartitioningDelta 表可以通过一个或多个partition列 hint 进行分区。以下示例按foo列分区dlt.resource( table_formatdelta, columns{foo: {partition: True}}, ) def my_delta_resource(): ...需要注意两点Delta 使用Hive 风格分区即.../fooxxx/...这类目录结构这与 Iceberg 的 spec 式分区不同不支持分区演化partition evolution表创建后无法变更分区列。从源码看分区列的传递链路是normalize 阶段把columns{foo: {partition: True}}记录到表 schemaload 阶段由TableFormatLoadFilesystemJob._partition_columns通过get_columns_names_with_prop(self._load_table, partition)取出最终作为partition_by参数传给write_deltalake见 filesystem.py 与 deltalake.py。写入模式与 schema 演进write_delta_table是一个围绕deltalake.write_deltalake的薄封装其中两个关键行为值得说明见 dlt/common/libs/deltalake.pydlt 写入语义到 Delta write mode 的映射由get_delta_write_mode完成append/merge→appendmerge 处置在写入层面解析为 append真正的 upsert 由DeltaLoadFilesystemJob的 merge 分支处理replace→overwrite。始终传入schema_modemerge即开启 schema 演进允许新增列因此后续 load 中出现新字段时不会失败。写入前还有一次 Arrow 数据兼容化处理ensure_delta_compatible_arrow_schema会把 Delta 不支持的类型如null、time、decimal256转换为 string当存在分区列且其为 dictionary 类型时还会把 dictionary 列转换为其 value 类型delta-rs 不允许 dictionary 分区字段。这套转换逻辑同样被 merge 路径复用。表访问助手函数 get_delta_tables你可以使用get_delta_tables助手函数获取原生 Delta 表对象返回的是deltalake的 DeltaTable 对象from dlt.common.libs.deltalake import get_delta_tables # 获取 DeltaTable 对象的字典 delta_tables get_delta_tables(pipeline) # 对 DeltaTable 对象执行操作 delta_tables[my_delta_table].optimize.compact() delta_tables[another_delta_table].optimize.z_order([col_a, col_b]) # delta_tables[my_delta_table].vacuum() # 等等get_delta_tables的函数签名为get_delta_tables(pipeline, *tables, schema_nameNone, include_dlt_tablesFalse)见 deltalake.py默认返回pipeline.default_schema中所有 Delta 表构成的字典键为表名值为DeltaTable对象传入*tables可按表名过滤若指定表名不存在会抛出ValueError可通过schema_name切换到其他 schemainclude_dlt_tables控制是否包含_dlt_系统表。由于返回的是原生DeltaTable对象你可以直接调用optimize.compact()、optimize.z_order()、vacuum()、history()、version()等 delta-rs API 进行表维护与治理测试用例 tests/load/pipeline/test_open_table_pipeline.py 中即用get_delta_tables(pipeline)[...].version()断言表的版本号如两次 load 后version() 2并调用history()验证版本历史可查。表截断Table Truncation当 dlt 截断 Delta 表时——无论是refreshdrop_data还是replace链中未收到数据的表——它都会执行一次事务性删除transactional delete提交一个没有行的新版本而表本身、schema 与版本历史都完整保留。因此读者永远不会看到一张被删了一半的表物理 Parquet 文件会保留下来以支持时间旅行直到你执行vacuum才会真正清理。源码层面dlt/common/libs/deltalake.py 中truncate_delta_table的实现就是table.delete()注释明确说明“在单次事务提交中删除所有行保留 schema 与版本历史”filesystem.py 的truncate_tables先检测表是否为 open tabledelta / iceberg是则调用_truncate_open_table做事务性清空普通非 Delta表则走_delete_table_files直接删除文件另有_replaced_atomically_in_loadfilesystem.py保证Delta 表 replace处置的表会在 load job 中原子 overwrite不需要预先 truncate。这种“事务性清空而非删文件”的语义正是 Delta 事务日志带来的核心收益也是它区别于普通 Parquet 目录式表的关键差异。Google Cloud Storage 认证在 GCS 上使用 Delta 表格式时并非所有认证方式都受支持Service Account — ✅ 支持Application Default Credentials — ✅ 支持OAuth — ❌ 不支持原因是 Delta 写入 GCS 时由deltalake基于 object_store自行访问存储而 OAuth token 的刷新机制不在该链路中生产环境建议使用 Service Account 或 ADC。merge 支持upsert 与 insert-onlydelta表格式支持 upsert 与 insert-only 两种 merge 策略dlt.resource( write_disposition{disposition: merge, strategy: upsert}, primary_keymy_primary_key, table_formatdelta, ) def my_upsert_resource(): ...策略集合在 factory.py 中声明为supported_merge_strategies[upsert, insert-only]且仅对 delta/iceberg 表开放filesystem_merge_strategies_selector非 open table 返回空列表。从 merge_delta_table 的实现可以看到 dlt 如何把 merge 翻译为 Delta MERGE 语句匹配谓词对顶级表用全部primary_key列构造target.pk source.pk AND ...对子表parent关系使用首个带unique属性的列构造target.id source.idupsertwhen_matched_update_all()when_not_matched_insert_all()insert-only仅when_not_matched_insert_all()merge 前会先调用evolve_delta_table_schema见 deltalake.py按列名比对、为 Delta 表追加缺失的新列通过delta_table.alter.add_columns保证 source 中出现的新字段不会导致 merge 失败。已知限制不支持hard_deletehint不支持删除嵌套表中的记录。这意味着 JSON 列的更新如果涉及元素移除不会被传播。例如先加载{key: 1, nested: [1, 2]}再加载{key: 1, nested: [1]}嵌套表中元素2对应的记录不会被删除。流式执行模式 deltalake_streamed_exec默认情况下dlt 以**流式模式streamed mode**执行 Delta 表 upsert以降低内存压力。若要利用源表统计信息source table statistics推导提前剪枝谓词early pruning predicate可关闭流式执行[destination.filesystem] deltalake_streamed_exec false该配置项定义在 dlt/destinations/impl/filesystem/configuration.pydeltalake_storage_options: Optional[DictStrAny] None Additional storage options passed to deltalake library, overriding credentials-derived values. deltalake_configuration: Optional[DictStrOptionalStr] None Delta table configuration passed to write_deltalake and create_deltalake calls. deltalake_streamed_exec: bool True When true, delta merge operations use streamed execution to reduce memory usage.默认值TrueDeltaLoadFilesystemJob在 merge 分支将其透传给DeltaTable.merge(..., streamed_exec...)见 filesystem.py。可以推断数据量特别大时保持默认流式执行更稳妥在数据可整体载入内存、且希望 merge 更快利用统计信息剪枝时再考虑设为false。存储选项与 Delta 表配置你可以同时配置destination.filesystem.deltalake_storage_options与destination.filesystem.deltalake_configuration向deltalake库传递存储选项和表配置[destination.filesystem] deltalake_configuration {delta.enableChangeDataFeed: true, delta.minWriterVersion: 7} deltalake_storage_options {AWS_S3_LOCKING_PROVIDER: dynamodb, DELTA_DYNAMO_TABLE_NAME: custom_table_name}两者分别映射到deltalake库write_deltalake方法的参数dlt 配置项deltalake 参数作用deltalake_configurationconfigurationDelta 表本身的配置如 CDCdelta.enableChangeDataFeed、writer versiondelta.minWriterVersion等deltalake_storage_optionsstorage_options底层对象存储的行为如 S3 锁机制AWS_S3_LOCKING_PROVIDER、DELTA_DYNAMO_TABLE_NAME关于凭据你不需要在deltalake_storage_options中手写 AK/SKdlt 会把凭据系统中解析出的凭据与你提供的选项合并后再作为storage_options传入。合并逻辑在 deltalake_storage_options()从credentials实现WithObjectStoreRsCredentials的凭据对象调用to_object_store_rs_credentials()取得凭据字典将deltalake_storage_options作为额外选项叠加最终返回{**creds, **extra_options}若两个字典出现同名键会打一条 warning且deltalake_storage_options中的值优先。S3 必读使用 S3 时必须通过 storage options 配置锁locking行为例如上面示例中的AWS_S3_LOCKING_PROVIDER dynamodb配合一个 DynamoDB 表名否则并发提交可能失败。小结一个可运行的完整示例结合本文各部分一个包含分区、merge、CDC 配置的完整 pipeline 大致如下目标为 S3凭据按 dlt 常规方式通过环境变量或credentials参数提供不在此重复import dlt from dlt.common.libs.deltalake import get_delta_tables pipeline dlt.pipeline(destinationfilesystem, destination_namelakehouse) # 配套 dlt/secrets.toml或环境变量 # # [destination.lakehouse] # bucket_url s3://my-bucket # deltalake_configuration {delta.enableChangeDataFeed: true, delta.minWriterVersion: 7} # deltalake_storage_options {AWS_S3_LOCKING_PROVIDER: dynamodb, DELTA_DYNAMO_TABLE_NAME: dlt_locks} dlt.resource( table_formatdelta, columns{event_date: {partition: True}}, write_disposition{disposition: merge, strategy: upsert}, primary_keyid, ) def my_delta_resource(): yield [{id: 1, event_date: 2026-01-01, payload: {k: v}}, {id: 2, event_date: 2026-01-02, payload: {k: w}}] info pipeline.run(my_delta_resource()) # 表维护compact / z-order / vacuum 均走原生 DeltaTable API delta_tables get_delta_tables(pipeline, my_delta_resource) delta_tables[my_delta_resource].optimize.compact()延伸阅读官方文档Delta Lake destination、filesystem destination核心源码dlt/common/libs/deltalake.py、dlt/destinations/impl/filesystem/factory.py、dlt/destinations/impl/filesystem/filesystem.py配置定义dlt/destinations/impl/filesystem/configuration.py端到端测试tests/load/pipeline/test_open_table_pipeline.py覆盖 delta/iceberg 的核心读写、分区、版本历史与子表处理【免费下载链接】dltdata load tool (dlt) is an open source Python library that makes data loading easy ️项目地址: https://gitcode.com/GitHub_Trending/dl/dlt创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表