首页/新闻资讯/正文详情

Apache Druid Delta Lake 扩展实战:通过 DeltaInputSource 从 Lakehouse 表批量摄入数据

发布时间:2026/9/23 5:41:17 来源:云帆数科 栏目:资讯中心
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),仅供参考

相关推荐

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

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

claude-code:面向开发者的终端原生AI编程CLI工作流
claude-code:面向开发者的终端原生AI编程CLI工作流

1. 项目概述:这不是一个“工具”,而是一套面向开发者的终端智能协作工作流“claude-code”这个名称乍看像某个独立软件,但实际它根本不是传统意义上的可执行程序——它没有安装包、不提供GUI界面、也不走应用商店分发。我第一次在GitHub上看到… · 2026/9/23 5:41:10

家庭系统源码拆解:版本升级API全变,面试必问的底层逻辑
家庭系统源码拆解:版本升级API全变,面试必问的底层逻辑

家庭系统源码拆解:版本升级API全变,面试必问的底层逻辑 版本升级后 API 全变了,这种绝望感谁懂?刚把老接口封装好,新版文档出来一看,方法名全换,参数结构重组,之前的代码直接报废。这不仅是业务开发的噩梦,更是面试必问的底层架构题。很多候… · 2026/9/23 5:41:10

从Function Calling到技能包:构建可复用的智能体技能系统
从Function Calling到技能包:构建可复用的智能体技能系统

这几年做大模型应用,我有一个特别深的感触:真正难的不是把模型接进来,而是让模型稳定地干杂活。你写一个 agent,要它查资料、算数据、调接口、整理报告,如果每个能力都临时写死在 prompt 里,一两个功能还行… · 2026/9/23 6:33:53

3步吃透A调源码:别再只背八股,这次真能写项目
3步吃透A调源码:别再只背八股,这次真能写项目

3步吃透A调源码:别再只背八股,这次真能写项目 看了一堆教程还是不会写项目?别怪自己笨,是你缺了“源码解析”这一环。 很多新手卡在“听懂了但手不动”,根源在于只看了表层API,没看懂底层数据流。 今天不讲虚的,直接拆解一个真实场景中的… · 2026/9/23 6:33:53

深度学习进阶:CNN、分布式训练与GPU性能调优实战
深度学习进阶:CNN、分布式训练与GPU性能调优实战

1. 从第51集到第111集:这段内容到底在讲什么如果你正在跟《动手学深度学习》这套课程,大概会有个明显的感受:前50集像是在铺路,把张量、自动求导、线性回归、Softmax这些基础砖块一块块码齐;而从第51集开始&#xff0c… · 2026/9/23 6:33:47

诺基亚c7复刻避坑指南:3个步骤搞定API变更与完整示例
诺基亚c7复刻避坑指南:3个步骤搞定API变更与完整示例

诺基亚c7复刻避坑指南:3个步骤搞定API变更与完整示例 版本升级后 API 全变了,这是很多老程序员接手旧项目时的噩梦。我花了一周时间,把经典的诺基亚c7复刻成现代Web应用,踩了无数坑。今天直接甩出 完整示例 ,帮你省下这周时间。… · 2026/9/23 6:33:47

永辉超市供应商系统图解原理:3步搞定面试高频考点
永辉超市供应商系统图解原理:3步搞定面试高频考点

永辉超市供应商系统图解原理:3步搞定面试高频考点 面试被问原理答不上来,是不是经常脑子一片空白?特别是聊到永辉超市供应商系统这种大型零售后端架构,面试官一句“讲讲核心链路”,你只能支支吾吾。别慌,今天用图解原理的方式,把这套系统最核心的库存… · 2026/9/23 6:33:47

Agent技能体系实战:从工具调用到生产级应用
Agent技能体系实战:从工具调用到生产级应用

先说一下背景。我做AI应用层开发有几年了,最近半年几乎全扑在Agent相关的项目上。从最早拿LangChain拼个Demo,到后来在真实业务里落地带工具调用的Agent服务,中间踩过的坑比写过的代码还多。这个过程中我意识到一个核心问题:很多人… · 2026/9/23 6:33:47

3招搞定手机怎么下载微信面试难题实战项目解析
3招搞定手机怎么下载微信面试难题实战项目解析

3招搞定手机怎么下载微信面试难题实战项目解析 面试被问“手机怎么下载微信”背后的原理,90%的人答不上来。别笑,这看似弱智的问题,实则是考察你对移动应用分发机制、安全校验及网络协议理解的试金石。我带过不少校招新人,他们背了八股文,却连一个A… · 2026/9/23 0:00:03

你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型
你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型

你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型 面试被问“高并发下如何保证消息不丢失”,你张口就是“用Redis”,结果面试官追问“如果Redis宕机了怎么办”,你瞬间卡壳。这种场景太常见了,很多新手在背八股文时,只记住了技术名词… · 2026/9/23 0:00:29

Win7无线热点配置工具源码解析:解决API失效的3个实战技巧
Win7无线热点配置工具源码解析:解决API失效的3个实战技巧

Win7无线热点配置工具源码解析:解决API失效的3个实战技巧 Win7无线热点配置工具在Win10/11上跑不动?不是你的问题,是版本升级后 API 全变了。很多老项目里的 netsh wlan… · 2026/9/23 0:00:36

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码