尧图网络 高端网站定制 · 原创设计
免费咨询热线
400-888-6620
免费获取方案
Apache Druid Delta Lake 扩展实战:通过 DeltaInputSource 从 Lakehouse 表批量摄入数据
数据库OLAP大数据后端【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址https://gitcode.com/gh_mirrors/druid6/druid点击查看免费下载Delta Lake 是构建 Lakehouse 架构的开放存储框架而 Apache Druid 是高性能实时分析数据库两者结合可以把存储在 Delta 表中的数据直接摄入 Druid 进行实时分析。本文以仓库中的druid-deltalake-extensions扩展为核心讲解其基于 Delta Kernel 的底层工作原理、安装加载方式、delta输入源的完整配置方法以及 8 种 Delta 过滤器的实战用法读完后你可以直接在本地 Druid 集群中把 Delta Lake 表的最新快照摄入为 Druid 数据源。扩展概述为什么需要 Delta Lake 连接器Delta Lake 提供了事务性、可伸缩的数据湖能力支持 Spark、Flink 等多样化的计算引擎在同一份数据上工作。但 Druid 并不能直接读取 Delta 表——Delta 表中的数据以版本化 Parquet 文件的形式存储并带有事务日志_delta_log普通 Parquet 输入源无法感知其协议。为此Druid 官方提供了社区扩展 druid-deltalake-extensions其中实现了DeltaInputSource即type: delta输入源。它的作用正如 delta-lake.md 所述Delta Lake is an open source storage framework that enables building a Lakehouse architecture with various compute engines. DeltaLakeInputSource lets you ingest data stored in a Delta Lake table into Apache Druid.要使用该扩展需要将druid-deltalake-extensions添加到 Druid 的已加载扩展列表中具体加载方式参见 Loading extensions社区扩展的加载说明见同页的 Loading community extensions 小节。工作原理从最新快照到 Druid InputRow 的数据链路该扩展并没有重新实现 Delta 协议而是直接基于 Delta Kernel API从 Delta Lake 3.0.0 引入的官方内核抽象与 Delta 表交互。从 DeltaInputSource.java 的源码可以梳理出完整的摄入流程定位表通过Table.forPath(engine, tablePath)打开指定路径的 Delta 表引擎由DefaultEngine.create(conf)创建内部基于 HadoopConfiguration。取最新快照table.getLatestSnapshot(engine)获取当前最新快照及其完整 Schema。值得注意的是代码在调用该 API 前会把上下文类加载器临时切换为LogStore的类加载器这是针对 Delta Kernel 3.2.0 在实例化LogStore时的已知问题对应 delta-io/delta 的 issue 3299所做的 workaround详见源码注释。列剪枝优化pruneSchema()根据InputRowSchema的ColumnsFilter从快照 Schema 中筛选需要的列构造物理读取 Schema从而在扫描阶段就只读必要列。构建 Scan 并应用过滤通过ScanBuilder把用户配置的 Delta 过滤器翻译为 Delta Kernel 的PredicatescanBuilder.withFilter(...)并配合裁剪后的读取 Schema 构建Scan。枚举数据文件scan.getScanFiles(engine)返回当前快照中需要读取的扫描文件列表每个文件对应一个DeltaSplit其state字段保存快照状态的 JSON 序列化结果files字段保存扫描文件列表见 DeltaSplit.java。读取 Parquet 数据对每个扫描文件engine.getParquetHandler().readParquetFiles(...)按物理读取 Schema和剩余谓词读取 Parquet再经Scan.transformPhysicalData转换为列式批数据。转换为 Druid 行DeltaInputSourceReader.java 中的DeltaInputSourceIterator逐批消费列式数据每个 KernelRow被包装为 DeltaInputRow.java最终通过MapInputRowParser解析成 Druid 的InputRow。文档对这一步的概括是Delta 输入源读取配置的 Delta 表并基于可选的 Delta 过滤器提取该表最新快照中的底层 Delta 文件——这些 Delta Lake 文件本身就是带版本信息的 Parquet 文件因此DeltaInputSource.needsFormat()直接返回false即无需也不允许再额外指定输入格式格式固定为 Parquet。可拆分性并行摄入的基础DeltaInputSource实现了SplittableInputSourceDeltaSplitcreateSplits()会把最新快照的每个扫描文件封装成一个独立的InputSplitestimateNumSplits()返回分片数量withSplit()则为每个分片构造只含单个 split 的新输入源。这套机制让index_parallel任务可以把不同 Delta 文件分发到不同任务并行处理是批式摄入吞吐的关键。版本支持根据 delta-lake.md 的 Version support 小节该扩展使用Delta Lake 3.0.0 引入的 Delta Kernel其兼容Apache Spark 3.5.x更旧的 Delta Lake 版本不受支持如需使用本扩展请升级到 Delta Lake 3.0.x 或更高版本。从仓库当前状态看pom.xml 中声明的delta-kernel.version为3.2.0依赖包括delta-kernel-api、delta-kernel-defaults与delta-storage三个内核模块。同时DeltaInputSource.java 的 Javadoc 明确指出目前 Delta Kernel 的 Table API 只支持读取最新快照。安装与加载扩展与大多数 Druid 社区扩展一样druid-deltalake-extensions通过pull-deps工具下载。在替换VERSION为期望的 Druid 版本后执行如下命令命令来自原文档保持原样java \ -cp lib/* \ -Ddruid.extensions.directoryextensions \ -Ddruid.extensions.hadoopDependenciesDirhadoop-dependencies \ org.apache.druid.cli.Main tools pull-deps \ --no-default-hadoop \ -c org.apache.druid.extensions.contrib:druid-deltalake-extensions:VERSION要点说明-Ddruid.extensions.directory指定扩展安装目录默认为extensions下载的 jar 会被放到这里-Ddruid.extensions.hadoopDependenciesDir指定 Hadoop 依赖目录--no-default-hadoop表示不拉取默认 Hadoop 依赖-c后的坐标由 groupIdorg.apache.druid.extensions.contrib、artifactIddruid-deltalake-extensions和版本号组成。下载完成后还需把该扩展加入druid.extensions.loadList配置参见 Loading extensions并重启相关 Druid 服务使其生效。若使用包含全部社区扩展的发行包该扩展已随包分发只需确认其在加载列表中。使用 Delta 输入源启用扩展后即可在批式摄入任务index_parallel的ioConfig.inputSource中使用type: delta。核心属性如下表源自 input-sources.md属性描述是否必填type固定为delta是tablePathDelta 表所在位置本地路径或对象存储路径是filter用于在快照内过滤数据文件的 JSON 对象否示例一读取整个快照以下 spec 读取/delta-table/foo表中的全部记录... ioConfig: { type: index_parallel, inputSource: { type: delta, tablePath: /delta-table/foo }, }在源码层面tablePath为空时会抛出InvalidInputtablePath cannot be null.且该路径直接传给 Delta Kernel 的Table.forPath()因此可以是本地文件系统路径也可以是 Delta Kernel 支持的对象存储路径。示例二带过滤器的读取以下 spec 只读取name Employee4 and age 30的数据... ioConfig: { type: index_parallel, inputSource: { type: delta, tablePath: /delta-table/foo, filter: { type: and, filters: [ { type: , column: name, value: Employee4 }, { type: , column: age, value: 30 } ] } }, }摄入任务的其余部分dataSchema、tuningConfig等与普通批式任务一致可参考 native-batch.md 中index_parallel的完整 spec 结构。Delta 过滤器详解过滤器的作用是在快照层面剪除不需要的数据文件从而减少 Druid 需要摄入的文件数量。输入源共提供 8 种过滤器and、or、not、、、、、。从 DeltaFilter.java 可以看到该接口通过 Jackson 注解注册了全部子类型type字段即 JSON 中的过滤器名每个过滤器最终通过getFilterPredicate(snapshotSchema)翻译成 Delta Kernel 的Predicate表达式树。各过滤器参数and过滤器逻辑与两个条件都必须为真属性描述是否必填type固定为and是filtersDelta 过滤器谓词列表要求恰好两个过滤器是or过滤器逻辑或满足其一即可属性描述是否必填type固定为or是filtersDelta 过滤器谓词列表要求恰好两个过滤器是not过滤器逻辑非属性描述是否必填type固定为not是filter被取反的 Delta 过滤器要求恰好一个是比较类过滤器、、、、参数一致属性描述是否必填type分别固定为、、、、是column应用过滤器的表列名是value过滤器使用的值是过滤的语义与保证需要特别注意过滤的语义边界这是该扩展最重要的使用前提原文档与源码 Javadoc 均明确说明对分区列过滤保证生效。当过滤器作用于分区表的分区列时Delta Kernel 可以精确剪枝只读取匹配分区的文件对非分区列过滤best-effort尽力而为。Delta Kernel 只依赖建表时收集的统计信息进行剪枝因此 Druid 连接器可能摄入不符合过滤条件的数据。若要确保 Delta Kernel 能剪除不必要的列值请只在分区列上使用过滤器。过滤器实现细节类型推断DeltaFilterUtils.java 的dataTypeToLiteral()会根据快照 Schema 中列的数据类型把字符串值转换为对应的 Delta 字面量支持String、Integer、Short、Long、Float、Double、Date若列不存在或类型不支持会抛出InvalidInput。数值列若传入非数字值会提示value must be a number。组合限制从 DeltaAndFilter.java 的源码看and/or目前只允许恰好两个谓词多余或不足都会抛出InvalidInput源码注释提到未来可以通过递归展平支持更复杂的表达式树。not则要求恰好一个谓词。翻译方式以为例DeltaEqualsFilter会构造new Predicate(, [Column(column), literal])and直接构造 Kernel 的And(left, right)谓词。数据类型映射Delta Kernel 的列式Row需要转换为 Druid 的InputRow。从 DeltaInputRow.java 的getValue()可以看到支持的 Delta 类型及转换规则Delta Kernel 类型转换结果BooleanTypebooleanByteType/ShortType/IntegerType对应整数类型DateType由DeltaTimeUtils.getSecondsFromDate(...)转换为 epoch 秒LongTypelongTimestampType由DeltaTimeUtils.getMillisFromTimestamp(...)转换为 epoch 毫秒FloatType/DoubleType对应浮点类型StringType字符串BinaryType字节数组按字符转换后以字符串形式返回DecimalType以decimal.longValue()转换为long其他类型抛出InvalidInputUnsupported data type其中Date/Timestamp的时间换算逻辑集中在 DeltaTimeUtils.java这决定了 Delta 表中的时间列进入 Druid 后的数值语义设计timestampSpec时应与其保持一致。行转换完成后DeltaInputRow委托给MapInputRowParser完成 Druid 维度/指标/时间戳的解析因此下游的timestampSpec、dimensionsSpec、metricsSpec用法与其他输入源完全一致。已知限制综合 delta-lake.md 的 Known limitations 小节与源码注释使用本扩展时需注意以下限制仅支持最新快照该扩展依赖 Delta Kernel API只能读取 Delta 表的最新快照无法读取任意历史快照任意快照读取能力由上游跟踪见 delta-io/delta 的 issue 2581。非分区列过滤是 best-effort对非分区列应用过滤器时可能摄入不匹配的数据详见上文过滤的语义与保证。and/or组合受限当前实现要求恰好两个谓词无法表达超过两个条件的组合除非嵌套and/or从代码结构看嵌套是可行的因为每个过滤器本身也是DeltaFilter。数据格式固定为 Parquet输入源不接收外部inputFormatDelta 表底层文件必须是版本化 Parquet 文件。版本要求需要 Delta Lake 3.0.0与 Spark 3.5.x 兼容更旧版本不受支持。深入阅读源码与测试若想进一步验证上述行为可参考仓库中的以下位置输入源实现DeltaInputSource.java、DeltaInputSourceReader.java、DeltaInputRow.java、DeltaSplit.java过滤器实现DeltaFilter.java 及filter包下的DeltaAndFilter、DeltaOrFilter、DeltaNotFilter、DeltaEqualsFilter、DeltaGreaterThanFilter、DeltaGreaterThanOrEqualsFilter、DeltaLessThanFilter、DeltaLessThanOrEqualsFilter、DeltaFilterUtils扩展装配DeltaLakeDruidModule.java测试用例extensions-contrib/druid-deltalake-extensions/src/test/java/org/apache/druid/delta/下的DeltaInputSourceTest、DeltaInputSourceSerdeTest、RowSerdeTest、DeltaTimeUtilsTest以及各过滤器测试其中PartitionedDeltaTable与NonPartitionedDeltaTable分别构造了分区/非分区表场景用于验证过滤行为输入源文档input-sources.md 的 Delta Lake input source 小节扩展文档delta-lake.md综上druid-deltalake-extensions是一个基于 Delta Kernel 的轻量连接器它让 Druid 得以以标准delta输入源消费 Delta Lake 表的最新快照通过扫描级过滤和列剪枝控制摄入规模并通过可拆分输入源支撑并行批式摄入。理解其仅最新快照分区列过滤才保证生效等边界是把它稳定用于生产摄入任务的关键。赞分享数据库OLAP大数据后端【免费下载链接】druidApache Druid: a high performance real-time analytics database.项目地址https://gitcode.com/gh_mirrors/druid6/druid点击查看免费下载相关推荐Apache Druid Thrift 扩展实战从实时流到 Hadoop 批量的 Thrift 数据摄取与解析Apache Druid Thrift 扩展实战从实时流到 Hadoop 批量的 Thrift 数据摄取与解析 Apache Druid 的 druid th数据库数据分析OLAP大数据实时分析数据仓库后端Apache Druid PostgreSQL 元数据存储与 PostgreSQL 批量摄入实战指南Apache Druid PostgreSQL 元数据存储与 PostgreSQL 批量摄入实战指南 Apache Druid 的 Coordinator、Ov数据库OLAP大数据后端用 Lucky 内网穿透三步搞定在家外的任何地方访问内网服务用 Lucky 内网穿透三步搞定在家外的任何地方访问内网服务 Lucky 是一款面向软硬路由的公网管理工具集成端口转发、动态域名DDNS、反向代理、网络后端网络通信创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED

相关推荐

储能运维工程师培训机构推荐:从报名学习到考试拿证,报考全攻略

储能运维工程师培训机构推荐:从报名学习到考试拿证,报考全攻略

储能是新型电力系统的重要支柱,随着新能源装机快速增长,储能电站如雨后春笋般涌现,储能运维工程师成为新能源领域最紧缺的人才之一。储能运维工程师是做什么的?前景怎么样?怎么考证?本文给你一份完整的储能…

📅 2026/9/23 5:41:41
使用 Laradock 将 PHP 应用部署到 Google Cloud Run:一条命令从开发镜像到 Serverless 上线

使用 Laradock 将 PHP 应用部署到 Google Cloud Run:一条命令从开发镜像到 Serverless 上线

使用 Laradock 将 PHP 应用部署到 Google Cloud Run:一条命令从开发镜像到 Serverless 上线 【免费下载链接】laradock Full PHP development environment for Docker. Run Laravel, Symfony, CodeIgniter, Phalcon, WordPress, Drupal, Magento, Moodle, or any PH…

📅 2026/9/23 5:41:41
安防系统工程师培训机构推荐:从报名学习到考试拿证,报考全攻略

安防系统工程师培训机构推荐:从报名学习到考试拿证,报考全攻略

从视频监控到门禁系统,从入侵报警到智慧安防,安防系统已成为楼宇、园区、城市的”安全基础设施”。安防系统工程师作为安防行业的技术人才,需求持续稳定。本文给你一份完整的安防系统工程师报考全攻略。 一、安防系统工程师是做什么的&#x…

📅 2026/9/23 5:41:41
MORE NEWS

更多资讯

📰

GayPon:LGBTQ+垂直团购平台的信任生态与冷启动实战

1. 项目缘起与需求判断1.1 这个点子是怎么冒出来的先说清楚,GayPon不是什么标新立异的恶搞,而是把两个已经被验证的商业模式做了一次精准拼接:左边是Groupon(本地生活团购),右边是Gay(LGBTQ人群…

📰

AARRR模型实战指南:用户增长分析的核心指标与实操流程

1. 为什么AARRR模型至今仍是用户增长分析的底层框架第一次接触AARRR模型是在一个电商项目的数据复盘会上。当时运营团队报上来一堆指标——日活、注册量、下单转化率、复购率、分享次数——数据铺满三块大屏,但没人能说清楚问题到底出在哪个环节。后来一位从硅谷回来…

📰

嵌入式开发者的Claude Code安装配置与实战指南

1. 嵌入式开发者的AI编程工具链现状嵌入式软件工程师过去十年的工作流基本没怎么变过:交叉编译工具链、串口终端、调试器、示波器,再加上一堆芯片原厂的SDK和参考手册。写代码这件事本身,在嵌入式领域一直是个"重手工活"——寄存器…

📰

设计师配色网站实战指南:从选色到落地的完整方法论

配色这件事,说它是设计师的“日常刚需”一点都不夸张。不管你是做UI、做品牌、做电商详情页,还是偶尔帮朋友P个海报,只要涉及到视觉输出,颜色选不对,后面排版再精致也白搭。我见过太多刚入行的朋友,拿到需求…

📰

技术博文合规创作规范与可执行标题设计指南

我无法根据“海南综合网(影视、动漫)”这一标题生成符合要求的博文内容。原因如下:该标题指向性模糊,无明确技术载体、实现路径或可操作对象。“海南综合网”并非公开可查的标准化平台名称,亦非通用技术术语(如Nginx镜…

📰

CUA实战指南:从零搭建AI自动操作电脑的完整流程

最近有个词在AI圈子里刷屏频率高得吓人:CUA。第一次看到的人可能会把它当成某个新梗,甚至有人直接念成“夸”。可你要是真把它当梗,就容易错过一件挺重要的事——CUA全称是Computer Use Agent,翻译过来就是“计算机使用智能体”&a…

TODAY

今日更新

THIS WEEK

本周精选

THIS MONTH

本月热门

读完文章,想聊聊您的网站?

告诉我们您的行业与需求,资深顾问一对一梳理方案与报价,全程免费。

📞 💬